diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractWorkerSourceTask.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractWorkerSourceTask.java index 693ef510f1a4d..671e9aefa7fef 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractWorkerSourceTask.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractWorkerSourceTask.java @@ -36,6 +36,7 @@ import org.apache.kafka.connect.errors.RetriableException; import org.apache.kafka.connect.header.Header; import org.apache.kafka.connect.header.Headers; +import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; import org.apache.kafka.connect.runtime.errors.Stage; import org.apache.kafka.connect.runtime.errors.ToleranceType; @@ -217,13 +218,14 @@ protected AbstractWorkerSourceTask(ConnectorTaskId id, ConnectorOffsetBackingStore offsetStore, WorkerConfig workerConfig, ConnectMetrics connectMetrics, + ErrorHandlingMetrics errorMetrics, ClassLoader loader, Time time, RetryWithToleranceOperator retryWithToleranceOperator, StatusBackingStore statusBackingStore, Executor closeExecutor) { - super(id, statusListener, initialState, loader, connectMetrics, + super(id, statusListener, initialState, loader, connectMetrics, errorMetrics, retryWithToleranceOperator, time, statusBackingStore); this.workerConfig = workerConfig; diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ExactlyOnceWorkerSourceTask.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ExactlyOnceWorkerSourceTask.java index 931917b9e15ce..428ad2fb07b71 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ExactlyOnceWorkerSourceTask.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ExactlyOnceWorkerSourceTask.java @@ -27,6 +27,7 @@ import org.apache.kafka.common.utils.Time; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.connect.errors.ConnectException; +import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; import org.apache.kafka.connect.source.SourceRecord; import org.apache.kafka.connect.source.SourceTask; @@ -85,6 +86,7 @@ public ExactlyOnceWorkerSourceTask(ConnectorTaskId id, WorkerConfig workerConfig, ClusterConfigState configState, ConnectMetrics connectMetrics, + ErrorHandlingMetrics errorMetrics, ClassLoader loader, Time time, RetryWithToleranceOperator retryWithToleranceOperator, @@ -95,7 +97,7 @@ public ExactlyOnceWorkerSourceTask(ConnectorTaskId id, Runnable postProducerCheck) { super(id, task, statusListener, initialState, keyConverter, valueConverter, headerConverter, transformationChain, new WorkerSourceTaskContext(offsetReader, id, configState, buildTransactionContext(sourceConfig)), - producer, admin, topicGroups, offsetReader, offsetWriter, offsetStore, workerConfig, connectMetrics, + producer, admin, topicGroups, offsetReader, offsetWriter, offsetStore, workerConfig, connectMetrics, errorMetrics, loader, time, retryWithToleranceOperator, statusBackingStore, closeExecutor); this.transactionOpen = false; diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java index 293eb0f79d23d..da5dff2aa17d2 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java @@ -1307,7 +1307,7 @@ public WorkerTask doBuild(Task task, KafkaConsumer consumer = new KafkaConsumer<>(consumerProps); return new WorkerSinkTask(id, (SinkTask) task, statusListener, initialState, config, configState, metrics, keyConverter, - valueConverter, headerConverter, transformationChain, consumer, classLoader, time, + valueConverter, errorHandlingMetrics, headerConverter, transformationChain, consumer, classLoader, time, retryWithToleranceOperator, workerErrantRecordReporter, herder.statusBackingStore()); } } @@ -1366,7 +1366,7 @@ public WorkerTask doBuild(Task task, OffsetStorageWriter offsetWriter = new OffsetStorageWriter(offsetStore, id.connector(), internalKeyConverter, internalValueConverter); // Note we pass the configState as it performs dynamic transformations under the covers - return new WorkerSourceTask(id, (SourceTask) task, statusListener, initialState, keyConverter, valueConverter, + return new WorkerSourceTask(id, (SourceTask) task, statusListener, initialState, keyConverter, valueConverter, errorHandlingMetrics, headerConverter, transformationChain, producer, topicAdmin, topicCreationGroups, offsetReader, offsetWriter, offsetStore, config, configState, metrics, classLoader, time, retryWithToleranceOperator, herder.statusBackingStore(), executor); @@ -1434,7 +1434,7 @@ public WorkerTask doBuild(Task task, // Note we pass the configState as it performs dynamic transformations under the covers return new ExactlyOnceWorkerSourceTask(id, (SourceTask) task, statusListener, initialState, keyConverter, valueConverter, headerConverter, transformationChain, producer, topicAdmin, topicCreationGroups, - offsetReader, offsetWriter, offsetStore, config, configState, metrics, classLoader, time, retryWithToleranceOperator, + offsetReader, offsetWriter, offsetStore, config, configState, metrics, errorHandlingMetrics, classLoader, time, retryWithToleranceOperator, herder.statusBackingStore(), sourceConfig, executor, preProducerCheck, postProducerCheck); } } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java index 83faf5287795c..76dae3e7f3de7 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java @@ -42,6 +42,7 @@ import org.apache.kafka.connect.runtime.ConnectMetrics.MetricGroup; import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; +import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; import org.apache.kafka.connect.runtime.errors.Stage; import org.apache.kafka.connect.runtime.errors.WorkerErrantRecordReporter; import org.apache.kafka.connect.sink.SinkRecord; @@ -107,6 +108,7 @@ public WorkerSinkTask(ConnectorTaskId id, ConnectMetrics connectMetrics, Converter keyConverter, Converter valueConverter, + ErrorHandlingMetrics errorMetrics, HeaderConverter headerConverter, TransformationChain transformationChain, KafkaConsumer consumer, @@ -115,7 +117,7 @@ public WorkerSinkTask(ConnectorTaskId id, RetryWithToleranceOperator retryWithToleranceOperator, WorkerErrantRecordReporter workerErrantRecordReporter, StatusBackingStore statusBackingStore) { - super(id, statusListener, initialState, loader, connectMetrics, + super(id, statusListener, initialState, loader, connectMetrics, errorMetrics, retryWithToleranceOperator, time, statusBackingStore); this.workerConfig = workerConfig; diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSourceTask.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSourceTask.java index 37d93a3fe8685..8a8de1b0553db 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSourceTask.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSourceTask.java @@ -23,6 +23,7 @@ import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; +import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; import org.apache.kafka.connect.runtime.errors.Stage; import org.apache.kafka.connect.runtime.errors.ToleranceType; import org.apache.kafka.connect.source.SourceRecord; @@ -66,6 +67,7 @@ public WorkerSourceTask(ConnectorTaskId id, TargetState initialState, Converter keyConverter, Converter valueConverter, + ErrorHandlingMetrics errorMetrics, HeaderConverter headerConverter, TransformationChain transformationChain, Producer producer, @@ -85,7 +87,7 @@ public WorkerSourceTask(ConnectorTaskId id, super(id, task, statusListener, initialState, keyConverter, valueConverter, headerConverter, transformationChain, new WorkerSourceTaskContext(offsetReader, id, configState, null), producer, - admin, topicGroups, offsetReader, offsetWriter, offsetStore, workerConfig, connectMetrics, loader, + admin, topicGroups, offsetReader, offsetWriter, offsetStore, workerConfig, connectMetrics, errorMetrics, loader, time, retryWithToleranceOperator, statusBackingStore, closeExecutor); this.committableOffsets = CommittableOffsets.EMPTY; diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerTask.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerTask.java index e8086b763221b..0d1749901e361 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerTask.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerTask.java @@ -24,6 +24,8 @@ import org.apache.kafka.common.metrics.stats.Frequencies; import org.apache.kafka.common.metrics.stats.Max; import org.apache.kafka.common.utils.Time; +import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; import org.apache.kafka.connect.runtime.AbstractStatus.State; import org.apache.kafka.connect.runtime.ConnectMetrics.MetricGroup; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; @@ -63,6 +65,7 @@ abstract class WorkerTask implements Runnable { private volatile TargetState targetState; private volatile boolean stopping; // indicates whether the Worker has asked the task to stop private volatile boolean cancelled; // indicates whether the Worker has cancelled the task (e.g. because of slow shutdown) + private final ErrorHandlingMetrics errorMetrics; protected final RetryWithToleranceOperator retryWithToleranceOperator; @@ -71,11 +74,13 @@ public WorkerTask(ConnectorTaskId id, TargetState initialState, ClassLoader loader, ConnectMetrics connectMetrics, + ErrorHandlingMetrics errorMetrics, RetryWithToleranceOperator retryWithToleranceOperator, Time time, StatusBackingStore statusBackingStore) { this.id = id; this.taskMetricsGroup = new TaskMetricsGroup(this.id, connectMetrics, statusListener); + this.errorMetrics = errorMetrics; this.statusListener = taskMetricsGroup; this.loader = loader; this.targetState = initialState; @@ -147,7 +152,9 @@ public boolean awaitStop(long timeoutMs) { * Remove all metrics published by this task. */ public void removeMetrics() { - taskMetricsGroup.close(); + // Close quietly here so that we can be sure to close everything even if one attempt fails + Utils.closeQuietly(taskMetricsGroup::close, "Task metrics group"); + Utils.closeQuietly(errorMetrics, "Error handling metrics"); } protected abstract void initializeAndStart(); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/errors/ErrorHandlingMetrics.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/errors/ErrorHandlingMetrics.java index 419bea97f1200..e96832e3403bc 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/errors/ErrorHandlingMetrics.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/errors/ErrorHandlingMetrics.java @@ -23,16 +23,20 @@ import org.apache.kafka.connect.runtime.ConnectMetrics; import org.apache.kafka.connect.runtime.ConnectMetricsRegistry; import org.apache.kafka.connect.util.ConnectorTaskId; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * Contains various sensors used for monitoring errors. */ -public class ErrorHandlingMetrics { +public class ErrorHandlingMetrics implements AutoCloseable { private final Time time = new SystemTime(); private final ConnectMetrics.MetricGroup metricGroup; + private static final Logger log = LoggerFactory.getLogger(ErrorHandlingMetrics.class); + // metrics private final Sensor recordProcessingFailures; private final Sensor recordProcessingErrors; @@ -138,4 +142,13 @@ public void recordErrorTimestamp() { public ConnectMetrics.MetricGroup metricGroup() { return metricGroup; } + + /** + * Close the task Error metrics group when the task is closed + */ + @Override + public void close() { + log.debug("Removing error handling metrics of group {}", metricGroup.groupId()); + metricGroup.close(); + } } diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/AbstractWorkerSourceTaskTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/AbstractWorkerSourceTaskTest.java index d0833dbffc794..f2f63264e3653 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/AbstractWorkerSourceTaskTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/AbstractWorkerSourceTaskTest.java @@ -37,6 +37,7 @@ import org.apache.kafka.connect.errors.RetriableException; import org.apache.kafka.connect.header.ConnectHeaders; import org.apache.kafka.connect.integration.MonitorableSourceConnector; +import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperatorTest; import org.apache.kafka.connect.runtime.isolation.Plugins; import org.apache.kafka.connect.runtime.standalone.StandaloneConfig; @@ -137,6 +138,7 @@ public class AbstractWorkerSourceTaskTest { private WorkerConfig config; private SourceConnectorConfig sourceConfig; private MockConnectMetrics metrics = new MockConnectMetrics(); + @Mock private ErrorHandlingMetrics errorHandlingMetrics; private Capture producerCallbacks; private AbstractWorkerSourceTask workerTask; @@ -789,7 +791,7 @@ private void createWorkerTask(Converter keyConverter, Converter valueConverter, workerTask = new AbstractWorkerSourceTask( taskId, sourceTask, statusListener, TargetState.STARTED, keyConverter, valueConverter, headerConverter, transformationChain, sourceTaskContext, producer, admin, TopicCreationGroup.configuredGroups(sourceConfig), offsetReader, offsetWriter, offsetStore, - config, metrics, plugins.delegatingLoader(), Time.SYSTEM, RetryWithToleranceOperatorTest.NOOP_OPERATOR, + config, metrics, errorHandlingMetrics, plugins.delegatingLoader(), Time.SYSTEM, RetryWithToleranceOperatorTest.NOOP_OPERATOR, statusBackingStore, Runnable::run) { @Override protected void prepareToInitializeTask() { diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ErrorHandlingTaskTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ErrorHandlingTaskTest.java index f1913d9848ed3..780ca4e790e94 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ErrorHandlingTaskTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ErrorHandlingTaskTest.java @@ -575,7 +575,7 @@ private void createSinkTask(TargetState initialState, RetryWithToleranceOperator workerSinkTask = new WorkerSinkTask( taskId, sinkTask, statusListener, initialState, workerConfig, - ClusterConfigState.EMPTY, metrics, converter, converter, + ClusterConfigState.EMPTY, metrics, converter, converter, errorHandlingMetrics, headerConverter, sinkTransforms, consumer, pluginLoader, time, retryWithToleranceOperator, workerErrantRecordReporter, statusBackingStore); } @@ -604,7 +604,7 @@ private void createSourceTask(TargetState initialState, RetryWithToleranceOperat workerSourceTask = PowerMock.createPartialMock( WorkerSourceTask.class, new String[]{"commitOffsets", "isStopping"}, - taskId, sourceTask, statusListener, initialState, converter, converter, headerConverter, sourceTransforms, + taskId, sourceTask, statusListener, initialState, converter, converter, errorHandlingMetrics, headerConverter, sourceTransforms, producer, admin, TopicCreationGroup.configuredGroups(sourceConfig), offsetReader, offsetWriter, offsetStore, workerConfig, ClusterConfigState.EMPTY, metrics, pluginLoader, time, retryWithToleranceOperator, diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ExactlyOnceWorkerSourceTaskTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ExactlyOnceWorkerSourceTaskTest.java index 8346de4d51c33..8de1bbeecb8bf 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ExactlyOnceWorkerSourceTaskTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ExactlyOnceWorkerSourceTaskTest.java @@ -34,6 +34,7 @@ import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.integration.MonitorableSourceConnector; import org.apache.kafka.connect.runtime.ConnectMetrics.MetricGroup; +import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperatorTest; import org.apache.kafka.connect.runtime.isolation.Plugins; import org.apache.kafka.connect.runtime.standalone.StandaloneConfig; @@ -131,6 +132,7 @@ public class ExactlyOnceWorkerSourceTaskTest { private SourceConnectorConfig sourceConfig; private Plugins plugins; private MockConnectMetrics metrics; + @Mock private ErrorHandlingMetrics errorHandlingMetrics; private Time time; private CountDownLatch pollLatch; @Mock private SourceTask sourceTask; @@ -240,7 +242,7 @@ private void createWorkerTask(TargetState initialState) { private void createWorkerTask(TargetState initialState, Converter keyConverter, Converter valueConverter, HeaderConverter headerConverter) { workerTask = new ExactlyOnceWorkerSourceTask(taskId, sourceTask, statusListener, initialState, keyConverter, valueConverter, headerConverter, transformationChain, producer, admin, TopicCreationGroup.configuredGroups(sourceConfig), offsetReader, offsetWriter, offsetStore, - config, clusterConfigState, metrics, plugins.delegatingLoader(), time, RetryWithToleranceOperatorTest.NOOP_OPERATOR, statusBackingStore, + config, clusterConfigState, metrics, errorHandlingMetrics, plugins.delegatingLoader(), time, RetryWithToleranceOperatorTest.NOOP_OPERATOR, statusBackingStore, sourceConfig, Runnable::run, preProducerCheck, postProducerCheck); } diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskTest.java index f63d7e9485a41..d8c79ff23edb8 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskTest.java @@ -42,6 +42,7 @@ import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.runtime.WorkerSinkTask.SinkTaskMetricsGroup; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperatorTest; +import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; import org.apache.kafka.connect.runtime.isolation.PluginClassLoader; import org.apache.kafka.connect.runtime.standalone.StandaloneConfig; import org.apache.kafka.connect.sink.SinkConnector; @@ -155,6 +156,8 @@ public class WorkerSinkTaskTest { private StatusBackingStore statusBackingStore; @Mock private KafkaConsumer consumer; + @Mock + private ErrorHandlingMetrics errorHandlingMetrics; private Capture rebalanceListener = EasyMock.newCapture(); private Capture topicsRegex = EasyMock.newCapture(); @@ -182,7 +185,7 @@ private void createTask(TargetState initialState) { private void createTask(TargetState initialState, Converter keyConverter, Converter valueConverter, HeaderConverter headerConverter) { workerTask = new WorkerSinkTask( taskId, sinkTask, statusListener, initialState, workerConfig, ClusterConfigState.EMPTY, metrics, - keyConverter, valueConverter, headerConverter, + keyConverter, valueConverter, errorHandlingMetrics, headerConverter, transformationChain, consumer, pluginLoader, time, RetryWithToleranceOperatorTest.NOOP_OPERATOR, null, statusBackingStore); } diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskThreadedTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskThreadedTest.java index 7ca7f44117288..096ced35d0e87 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskThreadedTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskThreadedTest.java @@ -31,6 +31,7 @@ import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperatorTest; +import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; import org.apache.kafka.connect.runtime.isolation.PluginClassLoader; import org.apache.kafka.connect.runtime.standalone.StandaloneConfig; import org.apache.kafka.connect.sink.SinkConnector; @@ -125,6 +126,7 @@ public class WorkerSinkTaskThreadedTest { private Capture rebalanceListener = EasyMock.newCapture(); @Mock private TaskStatus.Listener statusListener; @Mock private StatusBackingStore statusBackingStore; + @Mock private ErrorHandlingMetrics errorHandlingMetrics; private long recordsReturned; @@ -141,7 +143,7 @@ public void setup() { workerConfig = new StandaloneConfig(workerProps); workerTask = new WorkerSinkTask( taskId, sinkTask, statusListener, initialState, workerConfig, ClusterConfigState.EMPTY, metrics, keyConverter, - valueConverter, headerConverter, + valueConverter, errorHandlingMetrics, headerConverter, new TransformationChain<>(Collections.emptyList(), RetryWithToleranceOperatorTest.NOOP_OPERATOR), consumer, pluginLoader, time, RetryWithToleranceOperatorTest.NOOP_OPERATOR, null, statusBackingStore); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSourceTaskTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSourceTaskTest.java index 898afacbc4d8f..f411efbde3b50 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSourceTaskTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSourceTaskTest.java @@ -35,6 +35,7 @@ import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperatorTest; +import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; import org.apache.kafka.connect.runtime.isolation.Plugins; import org.apache.kafka.connect.runtime.standalone.StandaloneConfig; import org.apache.kafka.connect.source.SourceRecord; @@ -144,6 +145,7 @@ public class WorkerSourceTaskTest { @Mock private Future sendFuture; @MockStrict private TaskStatus.Listener statusListener; @Mock private StatusBackingStore statusBackingStore; + @Mock private ErrorHandlingMetrics errorHandlingMetrics; private Capture producerCallbacks; @@ -228,7 +230,7 @@ private void createWorkerTask(TargetState initialState, RetryWithToleranceOperat private void createWorkerTask(TargetState initialState, Converter keyConverter, Converter valueConverter, HeaderConverter headerConverter, RetryWithToleranceOperator retryWithToleranceOperator) { - workerTask = new WorkerSourceTask(taskId, sourceTask, statusListener, initialState, keyConverter, valueConverter, headerConverter, + workerTask = new WorkerSourceTask(taskId, sourceTask, statusListener, initialState, keyConverter, valueConverter, errorHandlingMetrics, headerConverter, transformationChain, producer, admin, TopicCreationGroup.configuredGroups(sourceConfig), offsetReader, offsetWriter, offsetStore, config, clusterConfigState, metrics, plugins.delegatingLoader(), Time.SYSTEM, retryWithToleranceOperator, statusBackingStore, Runnable::run); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTaskTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTaskTest.java index ae9c75d525344..b0a23b636bb63 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTaskTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTaskTest.java @@ -22,6 +22,7 @@ import org.apache.kafka.connect.runtime.WorkerTask.TaskMetricsGroup; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperatorTest; +import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; import org.apache.kafka.connect.sink.SinkTask; import org.apache.kafka.connect.storage.StatusBackingStore; import org.apache.kafka.connect.util.ConnectorTaskId; @@ -57,6 +58,7 @@ public class WorkerTaskTest { @Mock private StatusBackingStore statusBackingStore; private ConnectMetrics metrics; private RetryWithToleranceOperator retryWithToleranceOperator; + @Mock private ErrorHandlingMetrics errorHandlingMetrics; @Before public void setup() { @@ -73,9 +75,8 @@ public void tearDown() { public void standardStartup() { ConnectorTaskId taskId = new ConnectorTaskId("foo", 0); - WorkerTask workerTask = new TestWorkerTask(taskId, statusListener, TargetState.STARTED, loader, metrics, + WorkerTask workerTask = new TestWorkerTask(taskId, statusListener, TargetState.STARTED, loader, metrics, errorHandlingMetrics, retryWithToleranceOperator, Time.SYSTEM, statusBackingStore); - workerTask.initialize(TASK_CONFIG); workerTask.run(); workerTask.stop(); @@ -89,7 +90,7 @@ public void standardStartup() { public void stopBeforeStarting() { ConnectorTaskId taskId = new ConnectorTaskId("foo", 0); - WorkerTask workerTask = new TestWorkerTask(taskId, statusListener, TargetState.STARTED, loader, metrics, + WorkerTask workerTask = new TestWorkerTask(taskId, statusListener, TargetState.STARTED, loader, metrics, errorHandlingMetrics, retryWithToleranceOperator, Time.SYSTEM, statusBackingStore) { @Override @@ -116,7 +117,7 @@ public void cancelBeforeStopping() throws Exception { ConnectorTaskId taskId = new ConnectorTaskId("foo", 0); final CountDownLatch stopped = new CountDownLatch(1); - WorkerTask workerTask = new TestWorkerTask(taskId, statusListener, TargetState.STARTED, loader, metrics, + WorkerTask workerTask = new TestWorkerTask(taskId, statusListener, TargetState.STARTED, loader, metrics, errorHandlingMetrics, retryWithToleranceOperator, Time.SYSTEM, statusBackingStore) { @Override @@ -225,9 +226,9 @@ private static abstract class TestSinkTask extends SinkTask { private static class TestWorkerTask extends WorkerTask { public TestWorkerTask(ConnectorTaskId id, Listener statusListener, TargetState initialState, ClassLoader loader, - ConnectMetrics connectMetrics, RetryWithToleranceOperator retryWithToleranceOperator, Time time, + ConnectMetrics connectMetrics, ErrorHandlingMetrics errorHandlingMetrics, RetryWithToleranceOperator retryWithToleranceOperator, Time time, StatusBackingStore statusBackingStore) { - super(id, statusListener, initialState, loader, connectMetrics, retryWithToleranceOperator, time, statusBackingStore); + super(id, statusListener, initialState, loader, connectMetrics, errorHandlingMetrics, retryWithToleranceOperator, time, statusBackingStore); } @Override