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]