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


##########
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`. This caps the connector to exactly one task even when the 
runtime requests more. For a metrics exporter this may be intentional, but it 
should be documented or explicitly guarded — silently ignoring `maxTasks` can 
confuse operators who configure a higher parallelism.



##########
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 {

Review Comment:
   No error handling in `put()`. If any single `MetricsExporter.export()` 
throws an exception, the remaining exporters in the list are skipped and the 
entire batch is lost. Consider catching per-exporter exceptions so one faulty 
exporter does not block the others, and log or route failures appropriately.



##########
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;

Review Comment:
   `metricsExporters` is not initialized at declaration — it is only assigned 
inside `init()`. If the framework ever calls `start()` or `put()` before 
`init()`, this will throw a NullPointerException. In 
`MetricsExportSinkConnector`, the equivalent field is initialized via an 
instance-initializer block. The task should do the same, or at minimum 
initialize at declaration: `private List<MetricsExporter> metricsExporters = 
ServiceProvicerUtil.getMetricsExporterServices();`



##########
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>
+
+    <groupId>org.apache.rocketmq</groupId>
+    <artifactId>rocketmq-connect-metrics-exporter</artifactId>

Review Comment:
   This module defines its own `<groupId>` and `<version>` instead of 
inheriting from the parent POM (`<parent>` block is absent). This means it 
won't be built as part of the reactor, won't inherit dependency management, and 
won't receive the project-wide license/enforcer plugins. It will be an orphaned 
module that must be built and released independently.



##########
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) {
+            exporter.export(sinkRecords);
+        }
+    }
+
+    @Override public void start(KeyValue config) {

Review Comment:
   `start()` does not call `super.start(config)`. While the current `SinkTask` 
base may not require it, omitting the super call can break lifecycle contracts 
in future framework versions. Same applies to `stop()`.



##########
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 contains a typo: `ServiceProvicerUtil` should be 
`ServiceProviderUtil`. This typo propagates to every call site (connector and 
task classes) and will be a permanent API blemish if not fixed before merge.



##########
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:
   The `metricsExporters` field is populated in an instance-initializer block, 
meaning ServiceLoader runs at construction time — before `start()` or 
`validate()`. If the connector is constructed but never started (e.g., 
validation fails), the SPI scan is wasted work. Consider lazy initialization in 
`start()` for consistency with the connector lifecycle.



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