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 6d5446d9c213c..4f9e0936ee061 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 @@ -37,6 +37,7 @@ 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.ErrorReporter; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; import org.apache.kafka.connect.runtime.errors.Stage; import org.apache.kafka.connect.runtime.errors.ToleranceType; @@ -65,6 +66,7 @@ import java.util.concurrent.CountDownLatch; import java.util.concurrent.Executor; import java.util.concurrent.TimeUnit; +import java.util.function.Supplier; import static org.apache.kafka.connect.runtime.WorkerConfig.TOPIC_TRACKING_ENABLE_CONFIG; @@ -195,6 +197,7 @@ protected abstract void producerSendFailed( private final boolean topicTrackingEnabled; private final TopicCreation topicCreation; private final Executor closeExecutor; + private final Supplier> errorReportersSupplier; // Visible for testing List toSend; @@ -224,7 +227,8 @@ protected AbstractWorkerSourceTask(ConnectorTaskId id, Time time, RetryWithToleranceOperator retryWithToleranceOperator, StatusBackingStore statusBackingStore, - Executor closeExecutor) { + Executor closeExecutor, + Supplier> errorReportersSupplier) { super(id, statusListener, initialState, loader, connectMetrics, errorMetrics, retryWithToleranceOperator, time, statusBackingStore); @@ -242,6 +246,7 @@ protected AbstractWorkerSourceTask(ConnectorTaskId id, this.offsetStore = Objects.requireNonNull(offsetStore, "offset store cannot be null for source tasks"); this.closeExecutor = closeExecutor; this.sourceTaskContext = sourceTaskContext; + this.errorReportersSupplier = errorReportersSupplier; this.stopRequestedLatch = new CountDownLatch(1); this.sourceTaskMetricsGroup = new SourceTaskMetricsGroup(id, connectMetrics); @@ -261,6 +266,7 @@ public void initialize(TaskConfig taskConfig) { @Override protected void initializeAndStart() { + retryWithToleranceOperator.reporters(errorReportersSupplier.get()); prepareToInitializeTask(); offsetStore.start(); // If we try to start the task at all by invoking initialize, then count this as 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 f8f0d9f393c6a..5a123f2cd1e54 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 @@ -28,6 +28,7 @@ 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.ErrorReporter; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; import org.apache.kafka.connect.source.SourceRecord; import org.apache.kafka.connect.source.SourceTask; @@ -47,11 +48,13 @@ import org.slf4j.LoggerFactory; import java.util.LinkedHashMap; +import java.util.List; import java.util.Map; import java.util.Objects; import java.util.Optional; import java.util.concurrent.Executor; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Supplier; /** @@ -94,11 +97,12 @@ public ExactlyOnceWorkerSourceTask(ConnectorTaskId id, SourceConnectorConfig sourceConfig, Executor closeExecutor, Runnable preProducerCheck, - Runnable postProducerCheck) { + Runnable postProducerCheck, + Supplier> errorReportersSupplier) { super(id, task, statusListener, initialState, keyConverter, valueConverter, headerConverter, transformationChain, new WorkerSourceTaskContext(offsetReader, id, configState, buildTransactionContext(sourceConfig)), producer, admin, topicGroups, offsetReader, offsetWriter, offsetStore, workerConfig, connectMetrics, errorMetrics, - loader, time, retryWithToleranceOperator, statusBackingStore, closeExecutor); + loader, time, retryWithToleranceOperator, statusBackingStore, closeExecutor, errorReportersSupplier); this.transactionOpen = false; this.committableRecords = new LinkedHashMap<>(); 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 d0e410c386cc5..2f2aa762038b0 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 @@ -1797,7 +1797,6 @@ public WorkerTask doBuild(Task task, TransformationChain transformationChain = new TransformationChain<>(connectorConfig.transformationStages(), retryWithToleranceOperator); log.info("Initializing: {}", transformationChain); SinkConnectorConfig sinkConfig = new SinkConnectorConfig(plugins, connectorConfig.originalsStrings()); - retryWithToleranceOperator.reporters(sinkTaskReporters(id, sinkConfig, errorHandlingMetrics, connectorClass)); WorkerErrantRecordReporter workerErrantRecordReporter = createWorkerErrantRecordReporter(sinkConfig, retryWithToleranceOperator, keyConverter, valueConverter, headerConverter); @@ -1808,7 +1807,8 @@ public WorkerTask doBuild(Task task, return new WorkerSinkTask(id, (SinkTask) task, statusListener, initialState, config, configState, metrics, keyConverter, valueConverter, errorHandlingMetrics, headerConverter, transformationChain, consumer, classLoader, time, - retryWithToleranceOperator, workerErrantRecordReporter, herder.statusBackingStore()); + retryWithToleranceOperator, workerErrantRecordReporter, herder.statusBackingStore(), + () -> sinkTaskReporters(id, sinkConfig, errorHandlingMetrics, connectorClass)); } } @@ -1837,7 +1837,6 @@ public WorkerTask doBuild(Task task, SourceConnectorConfig sourceConfig = new SourceConnectorConfig(plugins, connectorConfig.originalsStrings(), config.topicCreationEnable()); - retryWithToleranceOperator.reporters(sourceTaskReporters(id, sourceConfig, errorHandlingMetrics)); TransformationChain transformationChain = new TransformationChain<>(sourceConfig.transformationStages(), retryWithToleranceOperator); log.info("Initializing: {}", transformationChain); @@ -1869,7 +1868,7 @@ public WorkerTask doBuild(Task task, 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); + retryWithToleranceOperator, herder.statusBackingStore(), executor, () -> sourceTaskReporters(id, sourceConfig, errorHandlingMetrics)); } } @@ -1905,7 +1904,6 @@ public WorkerTask doBuild(Task task, SourceConnectorConfig sourceConfig = new SourceConnectorConfig(plugins, connectorConfig.originalsStrings(), config.topicCreationEnable()); - retryWithToleranceOperator.reporters(sourceTaskReporters(id, sourceConfig, errorHandlingMetrics)); TransformationChain transformationChain = new TransformationChain<>(sourceConfig.transformationStages(), retryWithToleranceOperator); log.info("Initializing: {}", transformationChain); @@ -1935,7 +1933,8 @@ public WorkerTask doBuild(Task task, return new ExactlyOnceWorkerSourceTask(id, (SourceTask) task, statusListener, initialState, keyConverter, valueConverter, headerConverter, transformationChain, producer, topicAdmin, topicCreationGroups, offsetReader, offsetWriter, offsetStore, config, configState, metrics, errorHandlingMetrics, classLoader, time, retryWithToleranceOperator, - herder.statusBackingStore(), sourceConfig, executor, preProducerCheck, postProducerCheck); + herder.statusBackingStore(), sourceConfig, executor, preProducerCheck, postProducerCheck, + () -> sourceTaskReporters(id, sourceConfig, errorHandlingMetrics)); } } 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 59a7d24781067..7b514016c37a3 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 @@ -40,13 +40,14 @@ import org.apache.kafka.connect.header.ConnectHeaders; import org.apache.kafka.connect.header.Headers; 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.ErrorReporter; +import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; import org.apache.kafka.connect.runtime.errors.Stage; import org.apache.kafka.connect.runtime.errors.WorkerErrantRecordReporter; import org.apache.kafka.connect.sink.SinkRecord; import org.apache.kafka.connect.sink.SinkTask; +import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.storage.Converter; import org.apache.kafka.connect.storage.HeaderConverter; import org.apache.kafka.connect.storage.StatusBackingStore; @@ -61,6 +62,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.function.Supplier; import java.util.regex.Pattern; import java.util.stream.Collectors; @@ -98,6 +100,7 @@ class WorkerSinkTask extends WorkerTask { private boolean committing; private boolean taskStopped; private final WorkerErrantRecordReporter workerErrantRecordReporter; + private final Supplier> errorReportersSupplier; public WorkerSinkTask(ConnectorTaskId id, SinkTask task, @@ -116,7 +119,8 @@ public WorkerSinkTask(ConnectorTaskId id, Time time, RetryWithToleranceOperator retryWithToleranceOperator, WorkerErrantRecordReporter workerErrantRecordReporter, - StatusBackingStore statusBackingStore) { + StatusBackingStore statusBackingStore, + Supplier> errorReportersSupplier) { super(id, statusListener, initialState, loader, connectMetrics, errorMetrics, retryWithToleranceOperator, time, statusBackingStore); @@ -145,6 +149,7 @@ public WorkerSinkTask(ConnectorTaskId id, this.isTopicTrackingEnabled = workerConfig.getBoolean(TOPIC_TRACKING_ENABLE_CONFIG); this.taskStopped = false; this.workerErrantRecordReporter = workerErrantRecordReporter; + this.errorReportersSupplier = errorReportersSupplier; } @Override @@ -299,6 +304,7 @@ public int commitFailures() { @Override protected void initializeAndStart() { SinkConnectorConfig.validate(taskConfig); + retryWithToleranceOperator.reporters(errorReportersSupplier.get()); if (SinkConnectorConfig.hasTopicsConfig(taskConfig)) { List topics = SinkConnectorConfig.parseTopicsList(taskConfig); 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 501018cdce54c..e1ea9efa9bb6d 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 @@ -21,6 +21,7 @@ import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.utils.Time; import org.apache.kafka.connect.errors.ConnectException; +import org.apache.kafka.connect.runtime.errors.ErrorReporter; import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; @@ -40,6 +41,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.List; import java.util.Map; import java.util.Optional; import java.util.concurrent.ExecutionException; @@ -48,6 +50,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Supplier; import static org.apache.kafka.connect.runtime.SubmittedRecords.CommittableOffsets; @@ -84,12 +87,13 @@ public WorkerSourceTask(ConnectorTaskId id, Time time, RetryWithToleranceOperator retryWithToleranceOperator, StatusBackingStore statusBackingStore, - Executor closeExecutor) { + Executor closeExecutor, + Supplier> errorReportersSupplier) { super(id, task, statusListener, initialState, keyConverter, valueConverter, headerConverter, transformationChain, new WorkerSourceTaskContext(offsetReader, id, configState, null), producer, admin, topicGroups, offsetReader, offsetWriter, offsetStore, workerConfig, connectMetrics, errorMetrics, loader, - time, retryWithToleranceOperator, statusBackingStore, closeExecutor); + time, retryWithToleranceOperator, statusBackingStore, closeExecutor, errorReportersSupplier); this.committableOffsets = CommittableOffsets.EMPTY; this.submittedRecords = new SubmittedRecords(); 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 9ac7d5cca283b..26ee2fb4e0cb8 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,8 @@ 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.ErrorReporter; +import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperatorTest; import org.apache.kafka.connect.runtime.isolation.Plugins; import org.apache.kafka.connect.runtime.standalone.StandaloneConfig; @@ -72,6 +74,7 @@ import java.util.Optional; import java.util.Set; import java.util.concurrent.TimeoutException; +import java.util.function.Supplier; import java.util.stream.Collectors; import static org.apache.kafka.connect.integration.MonitorableSourceConnector.TOPIC_CONFIG; @@ -95,6 +98,7 @@ import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -348,7 +352,8 @@ public void testHeadersWithCustomConverter() throws Exception { StringConverter stringConverter = new StringConverter(); SampleConverterWithHeaders testConverter = new SampleConverterWithHeaders(); - createWorkerTask(stringConverter, testConverter, stringConverter); + createWorkerTask(stringConverter, testConverter, stringConverter, RetryWithToleranceOperatorTest.NOOP_OPERATOR, + Collections::emptyList); expectSendRecord(null); expectTopicCreation(TOPIC); @@ -689,6 +694,28 @@ public void testSendRecordsRetriableException() { verify(transformationChain, times(2)).apply(eq(record3)); } + @Test + public void testErrorReportersConfigured() { + RetryWithToleranceOperator retryWithToleranceOperator = mock(RetryWithToleranceOperator.class); + List errorReporters = Collections.singletonList(mock(ErrorReporter.class)); + createWorkerTask(keyConverter, valueConverter, headerConverter, retryWithToleranceOperator, () -> errorReporters); + workerTask.initializeAndStart(); + + ArgumentCaptor> errorReportersCapture = ArgumentCaptor.forClass(List.class); + verify(retryWithToleranceOperator).reporters(errorReportersCapture.capture()); + assertEquals(errorReporters, errorReportersCapture.getValue()); + } + + @Test + public void testErrorReporterConfigurationExceptionPropagation() { + createWorkerTask(keyConverter, valueConverter, headerConverter, RetryWithToleranceOperatorTest.NOOP_OPERATOR, + () -> { + throw new ConnectException("Failed to create error reporters"); + } + ); + assertThrows(ConnectException.class, () -> workerTask.initializeAndStart()); + } + private void expectSendRecord(Headers headers) { if (headers != null) expectConvertHeadersAndKeyValue(headers, TOPIC); @@ -799,15 +826,16 @@ private RecordHeaders emptyHeaders() { } private void createWorkerTask() { - createWorkerTask(keyConverter, valueConverter, headerConverter); + createWorkerTask(keyConverter, valueConverter, headerConverter, RetryWithToleranceOperatorTest.NOOP_OPERATOR, Collections::emptyList); } - private void createWorkerTask(Converter keyConverter, Converter valueConverter, HeaderConverter headerConverter) { + private void createWorkerTask(Converter keyConverter, Converter valueConverter, HeaderConverter headerConverter, + RetryWithToleranceOperator retryWithToleranceOperator, Supplier> errorReportersSupplier) { workerTask = new AbstractWorkerSourceTask( taskId, sourceTask, statusListener, TargetState.STARTED, keyConverter, valueConverter, headerConverter, transformationChain, sourceTaskContext, producer, admin, TopicCreationGroup.configuredGroups(sourceConfig), offsetReader, offsetWriter, offsetStore, - config, metrics, errorHandlingMetrics, plugins.delegatingLoader(), Time.SYSTEM, RetryWithToleranceOperatorTest.NOOP_OPERATOR, - statusBackingStore, Runnable::run) { + config, metrics, errorHandlingMetrics, plugins.delegatingLoader(), Time.SYSTEM, retryWithToleranceOperator, + statusBackingStore, Runnable::run, errorReportersSupplier) { @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 239141ec20d9e..3037803f6c592 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 @@ -33,7 +33,6 @@ import org.apache.kafka.connect.errors.RetriableException; import org.apache.kafka.connect.integration.MonitorableSourceConnector; import org.apache.kafka.connect.json.JsonConverter; -import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; import org.apache.kafka.connect.runtime.errors.ErrorReporter; import org.apache.kafka.connect.runtime.errors.LogReporter; @@ -48,6 +47,7 @@ import org.apache.kafka.connect.sink.SinkTask; import org.apache.kafka.connect.source.SourceRecord; import org.apache.kafka.connect.source.SourceTask; +import org.apache.kafka.connect.storage.ClusterConfigState; import org.apache.kafka.connect.storage.ConnectorOffsetBackingStore; import org.apache.kafka.connect.storage.Converter; import org.apache.kafka.connect.storage.HeaderConverter; @@ -75,12 +75,13 @@ import java.io.IOException; import java.time.Duration; +import java.util.Arrays; +import java.util.Collection; +import java.util.Collections; import java.util.HashMap; +import java.util.List; import java.util.Map; -import java.util.Arrays; import java.util.Set; -import java.util.Collections; -import java.util.Collection; import java.util.concurrent.Executor; import static java.util.Collections.emptyMap; @@ -98,16 +99,15 @@ import static org.apache.kafka.connect.runtime.TopicCreationConfig.REPLICATION_FACTOR_CONFIG; import static org.apache.kafka.connect.runtime.WorkerConfig.TOPIC_CREATION_ENABLE_CONFIG; import static org.junit.Assert.assertEquals; - import static org.mockito.ArgumentMatchers.eq; -import static org.mockito.Mockito.mock; -import static org.mockito.Mockito.spy; -import static org.mockito.Mockito.when; import static org.mockito.Mockito.any; +import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.doThrow; -import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; import static org.mockito.Mockito.times; -import static org.mockito.Mockito.doReturn; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; @RunWith(Parameterized.class) @@ -234,9 +234,8 @@ public void testSinkTasksCloseErrorReporters() throws Exception { ErrorReporter reporter = mock(ErrorReporter.class); RetryWithToleranceOperator retryWithToleranceOperator = operator(); - retryWithToleranceOperator.reporters(singletonList(reporter)); - createSinkTask(initialState, retryWithToleranceOperator); + createSinkTask(initialState, retryWithToleranceOperator, singletonList(reporter)); workerSinkTask.initialize(TASK_CONFIG); workerSinkTask.initializeAndStart(); workerSinkTask.close(); @@ -253,11 +252,11 @@ public void testSourceTasksCloseErrorReporters() throws IOException { ErrorReporter reporter = mock(ErrorReporter.class); RetryWithToleranceOperator retryWithToleranceOperator = operator(); - retryWithToleranceOperator.reporters(singletonList(reporter)); - createSourceTask(initialState, retryWithToleranceOperator); + createSourceTask(initialState, retryWithToleranceOperator, singletonList(reporter)); workerSourceTask.initialize(TASK_CONFIG); + workerSourceTask.initializeAndStart(); workerSourceTask.close(); verifyCloseSource(); verify(reporter).close(); @@ -269,15 +268,15 @@ public void testCloseErrorReportersExceptionPropagation() throws IOException { ErrorReporter reporterB = mock(ErrorReporter.class); RetryWithToleranceOperator retryWithToleranceOperator = operator(); - retryWithToleranceOperator.reporters(Arrays.asList(reporterA, reporterB)); - createSourceTask(initialState, retryWithToleranceOperator); + createSourceTask(initialState, retryWithToleranceOperator, Arrays.asList(reporterA, reporterB)); // Even though the reporters throw exceptions, they should both still be closed. doThrow(new RuntimeException()).when(reporterA).close(); doThrow(new RuntimeException()).when(reporterB).close(); workerSourceTask.initialize(TASK_CONFIG); + workerSourceTask.initializeAndStart(); workerSourceTask.close(); verify(reporterA).close(); @@ -293,9 +292,7 @@ public void testErrorHandlingInSinkTasks() throws Exception { LogReporter reporter = new LogReporter(taskId, connConfig(reportProps), errorHandlingMetrics); RetryWithToleranceOperator retryWithToleranceOperator = operator(); - retryWithToleranceOperator.reporters(singletonList(reporter)); - createSinkTask(initialState, retryWithToleranceOperator); - + createSinkTask(initialState, retryWithToleranceOperator, singletonList(reporter)); // valid json ConsumerRecord record1 = new ConsumerRecord<>( @@ -345,8 +342,7 @@ public void testErrorHandlingInSourceTasks() throws Exception { LogReporter reporter = new LogReporter(taskId, connConfig(reportProps), errorHandlingMetrics); RetryWithToleranceOperator retryWithToleranceOperator = operator(); - retryWithToleranceOperator.reporters(singletonList(reporter)); - createSourceTask(initialState, retryWithToleranceOperator); + createSourceTask(initialState, retryWithToleranceOperator, singletonList(reporter)); // valid json Schema valSchema = SchemaBuilder.struct().field("val", Schema.INT32_SCHEMA).build(); @@ -406,8 +402,7 @@ public void testErrorHandlingInSourceTasksWithBadConverter() throws Exception { LogReporter reporter = new LogReporter(taskId, connConfig(reportProps), errorHandlingMetrics); RetryWithToleranceOperator retryWithToleranceOperator = operator(); - retryWithToleranceOperator.reporters(singletonList(reporter)); - createSourceTask(initialState, retryWithToleranceOperator, badConverter()); + createSourceTask(initialState, retryWithToleranceOperator, singletonList(reporter), badConverter()); // valid json Schema valSchema = SchemaBuilder.struct().field("val", Schema.INT32_SCHEMA).build(); @@ -494,7 +489,8 @@ private void expectTopicCreation(String topic) { } } - private void createSinkTask(TargetState initialState, RetryWithToleranceOperator retryWithToleranceOperator) { + private void createSinkTask(TargetState initialState, RetryWithToleranceOperator retryWithToleranceOperator, + List errorReporters) { JsonConverter converter = new JsonConverter(); Map oo = workerConfig.originalsWithPrefix("value.converter."); oo.put("converter.type", "value"); @@ -509,17 +505,17 @@ private void createSinkTask(TargetState initialState, RetryWithToleranceOperator ClusterConfigState.EMPTY, metrics, converter, converter, errorHandlingMetrics, headerConverter, sinkTransforms, consumer, pluginLoader, time, retryWithToleranceOperator, workerErrantRecordReporter, - statusBackingStore); + statusBackingStore, () -> errorReporters); } - private void createSourceTask(TargetState initialState, RetryWithToleranceOperator retryWithToleranceOperator) { + private void createSourceTask(TargetState initialState, RetryWithToleranceOperator retryWithToleranceOperator, List errorReporters) { JsonConverter converter = new JsonConverter(); Map oo = workerConfig.originalsWithPrefix("value.converter."); oo.put("converter.type", "value"); oo.put("schemas.enable", "false"); converter.configure(oo); - createSourceTask(initialState, retryWithToleranceOperator, converter); + createSourceTask(initialState, retryWithToleranceOperator, errorReporters, converter); } private Converter badConverter() { @@ -531,8 +527,10 @@ private Converter badConverter() { return converter; } - private void createSourceTask(TargetState initialState, RetryWithToleranceOperator retryWithToleranceOperator, Converter converter) { - TransformationChain sourceTransforms = new TransformationChain<>(singletonList(new TransformationStage<>(new FaultyPassthrough())), retryWithToleranceOperator); + private void createSourceTask(TargetState initialState, RetryWithToleranceOperator retryWithToleranceOperator, + List errorReporters, Converter converter) { + TransformationChain sourceTransforms = new TransformationChain<>(singletonList( + new TransformationStage<>(new FaultyPassthrough())), retryWithToleranceOperator); workerSourceTask = spy(new WorkerSourceTask( taskId, sourceTask, statusListener, initialState, converter, @@ -542,7 +540,7 @@ private void createSourceTask(TargetState initialState, RetryWithToleranceOperat offsetReader, offsetWriter, offsetStore, workerConfig, ClusterConfigState.EMPTY, metrics, pluginLoader, time, retryWithToleranceOperator, - statusBackingStore, (Executor) Runnable::run)); + statusBackingStore, (Executor) Runnable::run, () -> errorReporters)); } 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 9939becc9c1fa..6aa5844dd2c6f 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 @@ -291,7 +291,7 @@ private void createWorkerTask(TargetState initialState, Converter keyConverter, workerTask = new ExactlyOnceWorkerSourceTask(taskId, sourceTask, statusListener, initialState, keyConverter, valueConverter, headerConverter, transformationChain, producer, admin, TopicCreationGroup.configuredGroups(sourceConfig), offsetReader, offsetWriter, offsetStore, config, clusterConfigState, metrics, errorHandlingMetrics, plugins.delegatingLoader(), time, RetryWithToleranceOperatorTest.NOOP_OPERATOR, statusBackingStore, - sourceConfig, Runnable::run, preProducerCheck, postProducerCheck); + sourceConfig, Runnable::run, preProducerCheck, postProducerCheck, Collections::emptyList); } @Test 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 5ab27fe78f0af..c943e22948b56 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 @@ -39,6 +39,8 @@ import org.apache.kafka.connect.runtime.ConnectMetrics.MetricGroup; import org.apache.kafka.connect.runtime.WorkerSinkTask.SinkTaskMetricsGroup; import org.apache.kafka.connect.runtime.errors.ErrorHandlingMetrics; +import org.apache.kafka.connect.runtime.errors.ErrorReporter; +import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator; import org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperatorTest; import org.apache.kafka.connect.runtime.isolation.PluginClassLoader; import org.apache.kafka.connect.runtime.standalone.StandaloneConfig; @@ -85,16 +87,19 @@ import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Function; +import java.util.function.Supplier; import java.util.regex.Pattern; import java.util.stream.Collectors; import static java.util.Arrays.asList; import static java.util.Collections.singleton; +import static org.easymock.EasyMock.createMock; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotEquals; import static org.junit.Assert.assertNull; import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertThrows; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; @@ -183,11 +188,16 @@ private void createTask(TargetState initialState) { } private void createTask(TargetState initialState, Converter keyConverter, Converter valueConverter, HeaderConverter headerConverter) { + createTask(initialState, keyConverter, valueConverter, headerConverter, RetryWithToleranceOperatorTest.NOOP_OPERATOR, Collections::emptyList); + } + + private void createTask(TargetState initialState, Converter keyConverter, Converter valueConverter, HeaderConverter headerConverter, + RetryWithToleranceOperator retryWithToleranceOperator, Supplier> errorReportersSupplier) { workerTask = new WorkerSinkTask( - taskId, sinkTask, statusListener, initialState, workerConfig, ClusterConfigState.EMPTY, metrics, - keyConverter, valueConverter, errorHandlingMetrics, headerConverter, - transformationChain, consumer, pluginLoader, time, - RetryWithToleranceOperatorTest.NOOP_OPERATOR, null, statusBackingStore); + taskId, sinkTask, statusListener, initialState, workerConfig, ClusterConfigState.EMPTY, metrics, + keyConverter, valueConverter, errorHandlingMetrics, headerConverter, + transformationChain, consumer, pluginLoader, time, + retryWithToleranceOperator, null, statusBackingStore, errorReportersSupplier); } @After @@ -1951,6 +1961,44 @@ public void testOriginalTopicWithTopicMutatingTransformations() { PowerMock.verifyAll(); } + @Test + public void testErrorReportersConfigured() { + RetryWithToleranceOperator retryWithToleranceOperator = createMock(RetryWithToleranceOperator.class); + List errorReporters = Collections.singletonList(createMock(ErrorReporter.class)); + createTask(initialState, keyConverter, valueConverter, headerConverter, retryWithToleranceOperator, + () -> errorReporters); + + expectInitializeTask(); + Capture> errorReportersCapture = EasyMock.newCapture(); + retryWithToleranceOperator.reporters(EasyMock.capture(errorReportersCapture)); + PowerMock.expectLastCall(); + + PowerMock.replayAll(retryWithToleranceOperator); + + workerTask.initialize(TASK_CONFIG); + workerTask.initializeAndStart(); + + assertEquals(errorReporters, errorReportersCapture.getValue()); + + PowerMock.verifyAll(); + } + + @Test + public void testErrorReporterConfigurationExceptionPropagation() { + createTask(initialState, keyConverter, valueConverter, headerConverter, RetryWithToleranceOperatorTest.NOOP_OPERATOR, + () -> { + throw new ConnectException("Failed to create error reporters"); + } + ); + + PowerMock.replayAll(); + + workerTask.initialize(TASK_CONFIG); + assertThrows(ConnectException.class, () -> workerTask.initializeAndStart()); + + PowerMock.verifyAll(); + } + private void expectInitializeTask() { consumer.subscribe(EasyMock.eq(asList(TOPIC)), EasyMock.capture(rebalanceListener)); PowerMock.expectLastCall(); 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 8e5c209aa0767..db90897a23c6b 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 @@ -145,7 +145,8 @@ public void setup() { taskId, sinkTask, statusListener, initialState, workerConfig, ClusterConfigState.EMPTY, metrics, keyConverter, valueConverter, errorHandlingMetrics, headerConverter, new TransformationChain<>(Collections.emptyList(), RetryWithToleranceOperatorTest.NOOP_OPERATOR), - consumer, pluginLoader, time, RetryWithToleranceOperatorTest.NOOP_OPERATOR, null, statusBackingStore); + consumer, pluginLoader, time, RetryWithToleranceOperatorTest.NOOP_OPERATOR, null, statusBackingStore, + Collections::emptyList); recordsReturned = 0; } 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 ec57aea306718..719dc10d5b1bc 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 @@ -264,7 +264,7 @@ private void createWorkerTask(TargetState initialState, Converter keyConverter, 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); + retryWithToleranceOperator, statusBackingStore, Runnable::run, Collections::emptyList); } @Test