RockteMQ-AI commented on code in PR #200:
URL: https://github.com/apache/rocketmq-connect/pull/200#discussion_r3817787787
##########
connectors/rocketmq-connect-metrics-exporter/src/main/java/org/apache/rocket/connect/metrics/export/sink/util/ServiceProvicerUtil.java:
##########
@@ -0,0 +1,23 @@
+package org.apache.rocket.connect.metrics.export.sink.util;
+
+import java.util.ArrayList;
+import java.util.Iterator;
+import java.util.List;
+import java.util.ServiceLoader;
+import org.apache.rocket.connect.metrics.export.sink.connector.MetricsExporter;
+
+/**
+ * @author: ming
+ */
Review Comment:
Typo in class name: ServiceProvicerUtil should be ServiceProviderUtil
(missing d in Provider).
##########
connectors/rocketmq-connect-metrics-exporter/src/main/java/org/apache/rocket/connect/metrics/export/sink/connector/MetricsExportSinkTask.java:
##########
@@ -0,0 +1,37 @@
+package org.apache.rocket.connect.metrics.export.sink.connector;
+
+import io.openmessaging.KeyValue;
+import io.openmessaging.connector.api.component.task.sink.SinkTask;
+import io.openmessaging.connector.api.component.task.sink.SinkTaskContext;
+import io.openmessaging.connector.api.data.ConnectRecord;
+import io.openmessaging.connector.api.errors.ConnectException;
+import java.util.List;
+import org.apache.rocket.connect.metrics.export.sink.util.ServiceProvicerUtil;
+
+
+public class MetricsExportSinkTask extends SinkTask {
+ private List<MetricsExporter> metricsExporters;
+
+ @Override public void put(List<ConnectRecord> sinkRecords) throws
ConnectException {
+ for (MetricsExporter exporter : metricsExporters) {
Review Comment:
metricsExporters is initialized in init() but used in put() without null
check. If put() is somehow called before init(), this will NPE. Consider adding
a null guard or lazy initialization.
##########
connectors/rocketmq-connect-metrics-exporter/src/main/java/org/apache/rocket/connect/metrics/export/sink/connector/MetricsExportSinkConnector.java:
##########
@@ -0,0 +1,41 @@
+package org.apache.rocket.connect.metrics.export.sink.connector;
+
+import io.openmessaging.KeyValue;
+import io.openmessaging.connector.api.component.task.Task;
+import io.openmessaging.connector.api.component.task.sink.SinkConnector;
+import java.util.ArrayList;
+import java.util.List;
+import org.apache.rocket.connect.metrics.export.sink.util.ServiceProvicerUtil;
+
+public class MetricsExportSinkConnector extends SinkConnector {
+ private KeyValue config;
+
+ private List<MetricsExporter> metricsExporters;
+ {
+ metricsExporters = ServiceProvicerUtil.getMetricsExporterServices();
Review Comment:
Instance initializer block runs at construction time, before start().
Consider moving ServiceLoader initialization to start() to avoid potential
classpath issues.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]