RockteMQ-AI commented on code in PR #200:
URL: https://github.com/apache/rocketmq-connect/pull/200#discussion_r3839554906


##########
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:
   In put(), if one exporter throws an exception, all subsequent exporters in 
the loop are skipped and the entire batch is lost. For a multi-exporter setup 
this means one faulty exporter blocks all others. Consider catching 
per-exporter, logging, and continuing — or aggregating exceptions before 
rethrowing.



##########
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:
   If metricsExporters is empty (which it will be given the broken SPI file), 
put() silently discards all incoming ConnectRecords with no log or error. A 
warning or exception when the list is empty during start() or init() would 
prevent silent data loss.



##########
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:
   metricsExporters is loaded in an instance initializer block but never 
closed/stopped in stop(). If exporters hold resources (e.g., HTTP clients, 
thread pools for Prometheus scraping), they will leak because the connector's 
stop() only nulls out config. Consider iterating metricsExporters in stop() and 
calling stop() on each, or documenting that lifecycle is fully managed by the 
task.



##########
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
+ */
+public class ServiceProvicerUtil {
+

Review Comment:
   Class name 'ServiceProvicerUtil' contains a typo — should be 
'ServiceProviderUtil'. This is a public API surface; fixing it later would be a 
breaking change.



##########
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();
+    }
+
+    @Override public List<KeyValue> taskConfigs(int maxTasks) {
+        List<KeyValue> configs = new ArrayList<>();
+        configs.add(config);
+        return configs;
+    }

Review Comment:
   taskConfigs(int maxTasks) always returns a single-element list regardless of 
maxTasks. If the framework distributes work across multiple tasks for 
parallelism, this silently caps it at one. Either honor maxTasks or explicitly 
document that this connector supports only a single task.



##########
connectors/rocketmq-connect-metrics-exporter/pom.xml:
##########
@@ -0,0 +1,47 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<project xmlns="http://maven.apache.org/POM/4.0.0";
+         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance";
+         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 
http://maven.apache.org/xsd/maven-4.0.0.xsd";>
+    <modelVersion>4.0.0</modelVersion>
+

Review Comment:
   pom.xml has no parent POM and uses a standalone hardcoded version 
'1.0-SNAPSHOT' with a standalone groupId. Other connector modules in this repo 
likely inherit from a parent POM for consistent dependency management. This may 
cause version drift and build inconsistency.



##########
connectors/rocketmq-connect-metrics-exporter/pom.xml:
##########
@@ -0,0 +1,47 @@
+<?xml version="1.0" encoding="UTF-8"?>

Review Comment:
   No test files are included in this PR. There is no unit or integration test 
coverage for the connector, task, or SPI loading logic — critical paths like 
put(), start()/stop() lifecycle, and validate() are completely untested.



-- 
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]

Reply via email to