From a95604cc657336a8aa55c5404f161a08ab87c6b7 Mon Sep 17 00:00:00 2001 From: Diego Erdody Date: Wed, 3 Nov 2021 11:14:25 -0700 Subject: [PATCH 1/2] Connect Sink connectors: Add support for topic-mutating SMTs to async connectors --- .../apache/kafka/connect/sink/SinkRecord.java | 30 ++++++- .../apache/kafka/connect/sink/SinkTask.java | 14 ++-- .../kafka/connect/sink/SinkRecordTest.java | 36 +++++++- .../connect/runtime/InternalSinkRecord.java | 11 +-- .../kafka/connect/runtime/WorkerSinkTask.java | 7 +- .../connect/tools/VerifiableSinkTask.java | 3 + .../ErrantRecordSinkConnector.java | 16 +--- .../integration/MonitorableSinkConnector.java | 32 ++++--- .../TransformationIntegrationTest.java | 84 ++++++++++++++++++- .../connect/runtime/WorkerSinkTaskTest.java | 67 +++++++++++++++ .../runtime/WorkerSinkTaskThreadedTest.java | 2 +- 11 files changed, 257 insertions(+), 45 deletions(-) diff --git a/connect/api/src/main/java/org/apache/kafka/connect/sink/SinkRecord.java b/connect/api/src/main/java/org/apache/kafka/connect/sink/SinkRecord.java index 12c7ee119abbb..123ad8a0dc368 100644 --- a/connect/api/src/main/java/org/apache/kafka/connect/sink/SinkRecord.java +++ b/connect/api/src/main/java/org/apache/kafka/connect/sink/SinkRecord.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.connect.sink; +import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.record.TimestampType; import org.apache.kafka.connect.connector.ConnectRecord; import org.apache.kafka.connect.data.Schema; @@ -32,6 +33,7 @@ public class SinkRecord extends ConnectRecord { private final long kafkaOffset; private final TimestampType timestampType; + private final TopicPartition originalTopicPartition; public SinkRecord(String topic, int partition, Schema keySchema, Object key, Schema valueSchema, Object value, long kafkaOffset) { this(topic, partition, keySchema, key, valueSchema, value, kafkaOffset, null, TimestampType.NO_TIMESTAMP_TYPE); @@ -44,9 +46,15 @@ public SinkRecord(String topic, int partition, Schema keySchema, Object key, Sch public SinkRecord(String topic, int partition, Schema keySchema, Object key, Schema valueSchema, Object value, long kafkaOffset, Long timestamp, TimestampType timestampType, Iterable
headers) { + this(topic, partition, keySchema, key, valueSchema, value, kafkaOffset, timestamp, timestampType, headers, null); + } + + public SinkRecord(String topic, int partition, Schema keySchema, Object key, Schema valueSchema, Object value, long kafkaOffset, + Long timestamp, TimestampType timestampType, Iterable
headers, TopicPartition originalTopicPartition) { super(topic, partition, keySchema, key, valueSchema, value, timestamp, headers); this.kafkaOffset = kafkaOffset; this.timestampType = timestampType; + this.originalTopicPartition = originalTopicPartition; } public long kafkaOffset() { @@ -57,6 +65,14 @@ public TimestampType timestampType() { return timestampType; } + /** + * @return topic and partition corresponding to the kafka record before transformations were applied. This is + * necessary for internal offset tracking, to be compatible with SMTs that mutate the topic name. + */ + public TopicPartition originalTopicPartition() { + return originalTopicPartition; + } + @Override public SinkRecord newRecord(String topic, Integer kafkaPartition, Schema keySchema, Object key, Schema valueSchema, Object value, Long timestamp) { return newRecord(topic, kafkaPartition, keySchema, key, valueSchema, value, timestamp, headers().duplicate()); @@ -65,7 +81,17 @@ public SinkRecord newRecord(String topic, Integer kafkaPartition, Schema keySche @Override public SinkRecord newRecord(String topic, Integer kafkaPartition, Schema keySchema, Object key, Schema valueSchema, Object value, Long timestamp, Iterable
headers) { - return new SinkRecord(topic, kafkaPartition, keySchema, key, valueSchema, value, kafkaOffset(), timestamp, timestampType, headers); + return new SinkRecord(topic, + kafkaPartition, + keySchema, + key, + valueSchema, + value, + kafkaOffset(), + timestamp, + timestampType, + headers, + originalTopicPartition); } @Override @@ -98,6 +124,8 @@ public String toString() { return "SinkRecord{" + "kafkaOffset=" + kafkaOffset + ", timestampType=" + timestampType + + ", originalTopicPartition=" + originalTopicPartition + "} " + super.toString(); } + } diff --git a/connect/api/src/main/java/org/apache/kafka/connect/sink/SinkTask.java b/connect/api/src/main/java/org/apache/kafka/connect/sink/SinkTask.java index 5d308c439625a..4647a59dfe788 100644 --- a/connect/api/src/main/java/org/apache/kafka/connect/sink/SinkTask.java +++ b/connect/api/src/main/java/org/apache/kafka/connect/sink/SinkTask.java @@ -29,7 +29,7 @@ * from those partitions. As records are fetched from Kafka, they will be passed to the sink task using the * {@link #put(Collection)} API, which should either write them to the downstream system or batch them for * later writing. Periodically, Connect will call {@link #flush(Map)} to ensure that batched records are - * actually pushed to the downstream system.. + * actually pushed to the downstream system. * * Below we describe the lifecycle of a SinkTask. * @@ -96,6 +96,9 @@ public void initialize(SinkTaskContext context) { * be stopped immediately. {@link SinkTaskContext#timeout(long)} can be used to set the maximum time before the * batch will be retried. * + * In case the connector does his own offset tracking (e.g. used by preCommit), the {@link SinkRecord#originalTopicPartition()} + * will have to be used, otherwise, the behavior will not be compatible with SMTs that mutate the topic name. + * * @param records the set of records to send */ public abstract void put(Collection records); @@ -104,8 +107,8 @@ public void initialize(SinkTaskContext context) { * Flush all records that have been {@link #put(Collection)} for the specified topic-partitions. * * @param currentOffsets the current offset state as of the last call to {@link #put(Collection)}}, - * provided for convenience but could also be determined by tracking all offsets included in the {@link SinkRecord}s - * passed to {@link #put}. + * provided for convenience but could also be determined by tracking all offsets included in the + * {@link SinkRecord#originalTopicPartition()}s passed to {@link #put}. */ public void flush(Map currentOffsets) { } @@ -116,8 +119,8 @@ public void flush(Map currentOffsets) { * The default implementation simply invokes {@link #flush(Map)} and is thus able to assume all {@code currentOffsets} are safe to commit. * * @param currentOffsets the current offset state as of the last call to {@link #put(Collection)}}, - * provided for convenience but could also be determined by tracking all offsets included in the {@link SinkRecord}s - * passed to {@link #put}. + * provided for convenience but could also be determined by tracking all offsets included in the + * {@link SinkRecord#originalTopicPartition()}s passed to {@link #put}. * * @return an empty map if Connect-managed offset commit is not desired, otherwise a map of offsets by topic-partition that are safe to commit. */ @@ -171,4 +174,5 @@ public void onPartitionsRevoked(Collection partitions) { */ @Override public abstract void stop(); + } diff --git a/connect/api/src/test/java/org/apache/kafka/connect/sink/SinkRecordTest.java b/connect/api/src/test/java/org/apache/kafka/connect/sink/SinkRecordTest.java index 02a6b2a78f28f..546a8d5ddc59a 100644 --- a/connect/api/src/test/java/org/apache/kafka/connect/sink/SinkRecordTest.java +++ b/connect/api/src/test/java/org/apache/kafka/connect/sink/SinkRecordTest.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.connect.sink; +import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.record.TimestampType; import org.apache.kafka.connect.data.Schema; import org.apache.kafka.connect.data.Values; @@ -46,7 +47,7 @@ public class SinkRecordTest { @BeforeEach public void beforeEach() { record = new SinkRecord(TOPIC_NAME, PARTITION_NUMBER, Schema.STRING_SCHEMA, "key", Schema.BOOLEAN_SCHEMA, false, KAFKA_OFFSET, - KAFKA_TIMESTAMP, TS_TYPE, null); + KAFKA_TIMESTAMP, TS_TYPE, null, new TopicPartition(TOPIC_NAME, PARTITION_NUMBER)); } @Test @@ -125,4 +126,37 @@ public void shouldModifyRecordHeader() { Header header = record.headers().lastWithName("intHeader"); assertEquals(100, (int) Values.convertToInteger(header.schema(), header.value())); } + + @Test + public void shouldCreateSinkRecordWithOriginalTopicPartition() { + TopicPartition tp = new TopicPartition("originalTopic", 100); + record = new SinkRecord(TOPIC_NAME, PARTITION_NUMBER, Schema.STRING_SCHEMA, "key", Schema.BOOLEAN_SCHEMA, false, KAFKA_OFFSET, + KAFKA_TIMESTAMP, TS_TYPE, null, tp); + assertSame(tp, record.originalTopicPartition()); + } + + @Test + public void shouldKeepOriginalTopicPartition() { + TopicPartition tp = new TopicPartition("originalTopic", 100); + record = new SinkRecord(TOPIC_NAME, PARTITION_NUMBER, Schema.STRING_SCHEMA, "key", Schema.BOOLEAN_SCHEMA, false, KAFKA_OFFSET, + KAFKA_TIMESTAMP, TS_TYPE, null, tp); + record = record.newRecord("otherTopic", 200, Schema.STRING_SCHEMA, "key", Schema.BOOLEAN_SCHEMA, false, KAFKA_OFFSET); + + assertSame(tp, record.originalTopicPartition()); + } + + @Test + public void shouldNotConsiderOriginalTopicPartitionForEquality() { + // For backwards compatibility + TopicPartition tp1 = new TopicPartition("originalTopic1", 1); + TopicPartition tp2 = new TopicPartition("originalTopic2", 1); + SinkRecord record1 = new SinkRecord(TOPIC_NAME, PARTITION_NUMBER, Schema.STRING_SCHEMA, "key", Schema.BOOLEAN_SCHEMA, false, KAFKA_OFFSET, + KAFKA_TIMESTAMP, TS_TYPE, null, tp1); + + SinkRecord record2 = new SinkRecord(TOPIC_NAME, PARTITION_NUMBER, Schema.STRING_SCHEMA, "key", Schema.BOOLEAN_SCHEMA, false, KAFKA_OFFSET, + KAFKA_TIMESTAMP, TS_TYPE, null, tp2); + + assertTrue(record1.equals(record2)); + } + } \ No newline at end of file diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/InternalSinkRecord.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/InternalSinkRecord.java index 69554ffb306b5..4ad4d948dbafd 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/InternalSinkRecord.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/InternalSinkRecord.java @@ -18,6 +18,7 @@ package org.apache.kafka.connect.runtime; import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.record.TimestampType; import org.apache.kafka.connect.data.Schema; import org.apache.kafka.connect.header.Header; @@ -32,18 +33,18 @@ public class InternalSinkRecord extends SinkRecord { private final ConsumerRecord originalRecord; - public InternalSinkRecord(ConsumerRecord originalRecord, SinkRecord record) { + public InternalSinkRecord(ConsumerRecord originalRecord, TopicPartition originalTopicPartition, SinkRecord record) { super(record.topic(), record.kafkaPartition(), record.keySchema(), record.key(), record.valueSchema(), record.value(), record.kafkaOffset(), record.timestamp(), - record.timestampType(), record.headers()); + record.timestampType(), record.headers(), originalTopicPartition); this.originalRecord = originalRecord; } - protected InternalSinkRecord(ConsumerRecord originalRecord, String topic, + protected InternalSinkRecord(ConsumerRecord originalRecord, TopicPartition originalTopicPartition, String topic, int partition, Schema keySchema, Object key, Schema valueSchema, Object value, long kafkaOffset, Long timestamp, TimestampType timestampType, Iterable
headers) { - super(topic, partition, keySchema, key, valueSchema, value, kafkaOffset, timestamp, timestampType, headers); + super(topic, partition, keySchema, key, valueSchema, value, kafkaOffset, timestamp, timestampType, headers, originalTopicPartition); this.originalRecord = originalRecord; } @@ -51,7 +52,7 @@ protected InternalSinkRecord(ConsumerRecord originalRecord, Stri public SinkRecord newRecord(String topic, Integer kafkaPartition, Schema keySchema, Object key, Schema valueSchema, Object value, Long timestamp, Iterable
headers) { - return new InternalSinkRecord(originalRecord, topic, kafkaPartition, keySchema, key, + return new InternalSinkRecord(originalRecord, originalTopicPartition(), topic, kafkaPartition, keySchema, key, valueSchema, value, kafkaOffset(), timestamp, timestampType(), headers()); } 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 ed7ad73b2a722..d598042f80c34 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 @@ -503,6 +503,7 @@ private SinkRecord convertAndTransformRecord(final ConsumerRecord record) { diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/tools/VerifiableSinkTask.java b/connect/runtime/src/main/java/org/apache/kafka/connect/tools/VerifiableSinkTask.java index ee58213cb5e9d..4ac82c61d370e 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/tools/VerifiableSinkTask.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/tools/VerifiableSinkTask.java @@ -72,7 +72,10 @@ public void put(Collection records) { data.put("topic", record.topic()); data.put("time_ms", nowMs); data.put("seqno", record.value()); + data.put("partition", record.kafkaPartition()); data.put("offset", record.kafkaOffset()); + data.put("originalTopic", record.originalTopicPartition().topic()); + data.put("originalPartition", record.originalTopicPartition().partition()); String dataJson; try { dataJson = JSON_SERDE.writeValueAsString(data); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ErrantRecordSinkConnector.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ErrantRecordSinkConnector.java index 0fe2f88083883..1d59e32f1d22c 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ErrantRecordSinkConnector.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ErrantRecordSinkConnector.java @@ -17,13 +17,10 @@ package org.apache.kafka.connect.integration; -import org.apache.kafka.common.TopicPartition; import org.apache.kafka.connect.connector.Task; import org.apache.kafka.connect.sink.ErrantRecordReporter; import org.apache.kafka.connect.sink.SinkRecord; -import java.util.Collection; -import java.util.HashMap; import java.util.Map; public class ErrantRecordSinkConnector extends MonitorableSinkConnector { @@ -47,15 +44,10 @@ public void start(Map props) { } @Override - public void put(Collection records) { - for (SinkRecord rec : records) { - taskHandle.record(); - TopicPartition tp = cachedTopicPartitions - .computeIfAbsent(rec.topic(), v -> new HashMap<>()) - .computeIfAbsent(rec.kafkaPartition(), v -> new TopicPartition(rec.topic(), rec.kafkaPartition())); - committedOffsets.put(tp, committedOffsets.getOrDefault(tp, 0L) + 1); - reporter.report(rec, new Throwable()); - } + protected void processRecord(SinkRecord rec) { + super.processRecord(rec); + reporter.report(rec, new Throwable()); } + } } diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/MonitorableSinkConnector.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/MonitorableSinkConnector.java index 7b9afa4419008..bf33abeb8b075 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/MonitorableSinkConnector.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/MonitorableSinkConnector.java @@ -88,17 +88,14 @@ public ConfigDef config() { public static class MonitorableSinkTask extends SinkTask { - private String connectorName; + private final Set assignments; + private final Map committedOffsets; private String taskId; - TaskHandle taskHandle; - Set assignments; - Map committedOffsets; - Map> cachedTopicPartitions; + private TaskHandle taskHandle; public MonitorableSinkTask() { this.assignments = new HashSet<>(); this.committedOffsets = new HashMap<>(); - this.cachedTopicPartitions = new HashMap<>(); } @Override @@ -109,7 +106,7 @@ public String version() { @Override public void start(Map props) { taskId = props.get("task.id"); - connectorName = props.get("connector.name"); + String connectorName = props.get("connector.name"); taskHandle = RuntimeHandles.get().connectorHandle(connectorName).taskHandle(taskId); log.debug("Starting task {}", taskId); taskHandle.recordTaskStart(); @@ -125,28 +122,29 @@ public void open(Collection partitions) { @Override public void put(Collection records) { for (SinkRecord rec : records) { - taskHandle.record(rec); - TopicPartition tp = cachedTopicPartitions - .computeIfAbsent(rec.topic(), v -> new HashMap<>()) - .computeIfAbsent(rec.kafkaPartition(), v -> new TopicPartition(rec.topic(), rec.kafkaPartition())); - committedOffsets.put(tp, committedOffsets.getOrDefault(tp, 0L) + 1); - log.trace("Task {} obtained record (key='{}' value='{}')", taskId, rec.key(), rec.value()); + processRecord(rec); } } + protected void processRecord(SinkRecord rec) { + taskHandle.record(rec); + committedOffsets.put(rec.originalTopicPartition(), rec.kafkaOffset() + 1); + log.trace("Task {} obtained record (key='{}' value='{}')", taskId, rec.key(), rec.value()); + } + @Override - public Map preCommit(Map offsets) { + public Map preCommit(Map currentOffsets) { + Map offsets = new HashMap<>(); for (TopicPartition tp : assignments) { Long recordsSinceLastCommit = committedOffsets.get(tp); if (recordsSinceLastCommit == null) { log.warn("preCommit was called with topic-partition {} that is not included " + "in the assignments of this task {}", tp, assignments); } else { - taskHandle.commit(recordsSinceLastCommit.intValue()); log.error("Forwarding to framework request to commit additional {} for {}", recordsSinceLastCommit, tp); - taskHandle.commit((int) (long) recordsSinceLastCommit); - committedOffsets.put(tp, 0L); + taskHandle.commit(recordsSinceLastCommit.intValue()); + offsets.put(tp, new OffsetAndMetadata(committedOffsets.remove(tp))); } } return offsets; diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java index 7b71f2f5f70a5..d8f662793a6ab 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java @@ -17,11 +17,15 @@ package org.apache.kafka.connect.integration; import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.OffsetAndMetadata; +import org.apache.kafka.common.TopicPartition; import org.apache.kafka.connect.storage.StringConverter; import org.apache.kafka.connect.transforms.Filter; +import org.apache.kafka.connect.transforms.RegexRouter; import org.apache.kafka.connect.transforms.predicates.HasHeaderKey; import org.apache.kafka.connect.transforms.predicates.RecordIsTombstone; import org.apache.kafka.connect.transforms.predicates.TopicNameMatches; +import org.apache.kafka.connect.util.SinkUtils; import org.apache.kafka.connect.util.clusters.EmbeddedConnectCluster; import org.apache.kafka.test.IntegrationTest; import org.junit.After; @@ -236,7 +240,7 @@ public void testFilterOnTombstonesWithSinkConnector() throws Exception { props.put(PREDICATES_CONFIG + ".barPredicate.type", RecordIsTombstone.class.getSimpleName()); // expect only half the records to be consumed by the connector - connectorHandle.expectedCommits(numRecords); + connectorHandle.expectedCommits(numRecords / 2); connectorHandle.expectedRecords(numRecords / 2); // start a sink connector @@ -324,4 +328,82 @@ public void testFilterOnHasHeaderKeyWithSourceConnectorAndTopicCreation() throws // delete connector connect.deleteConnector(CONNECTOR_NAME); } + + /** + * Verify that connectors overriding preCommit are compatible with topic-mutating SMTs + */ + @Test + public void testTopicMutationWithPreCommit() throws Exception { + assertConnectReady(); + + Map observedRecords = observeRecords(); + + // create test topics + int partitions = 1; + String fooTopic = "foo-topic"; + String barTopic = "bar-topic"; + int numFooRecords = 1024; + int numBarRecords = 666; + connect.kafka().createTopic(fooTopic, partitions); + connect.kafka().createTopic(barTopic, partitions); + + // setup up props for the sink connector + Map props = new HashMap<>(); + props.put("name", CONNECTOR_NAME); + props.put(CONNECTOR_CLASS_CONFIG, SINK_CONNECTOR_CLASS_NAME); + props.put(TASKS_MAX_CONFIG, String.valueOf(NUM_TASKS)); + props.put(TOPICS_CONFIG, String.join(",", fooTopic, barTopic)); + props.put(KEY_CONVERTER_CLASS_CONFIG, StringConverter.class.getName()); + props.put(VALUE_CONVERTER_CLASS_CONFIG, StringConverter.class.getName()); + props.put(TRANSFORMS_CONFIG, "example"); + props.put(TRANSFORMS_CONFIG + ".example.type", RegexRouter.class.getName()); + props.put(TRANSFORMS_CONFIG + ".example.regex", "(.*)"); + props.put(TRANSFORMS_CONFIG + ".example.replacement", "static-topic"); + + // expect all records to be consumed by the connector + connectorHandle.expectedRecords(numFooRecords + numBarRecords); + + // expect all records to be consumed by the connector + connectorHandle.expectedCommits(numFooRecords + numBarRecords); + + // start a sink connector + connect.configureConnector(CONNECTOR_NAME, props); + assertConnectorRunning(); + + // produce some messages into source topic partitions + for (int i = 0; i < numBarRecords; i++) { + connect.kafka().produce(barTopic, i % partitions, "key", "simple-message-value-" + i); + } + for (int i = 0; i < numFooRecords; i++) { + connect.kafka().produce(fooTopic, i % partitions, "key", "simple-message-value-" + i); + } + + // consume all records from the source topic or fail, to ensure that they were correctly produced. + assertEquals("Unexpected number of records consumed", numFooRecords, + connect.kafka().consume(numFooRecords, RECORD_TRANSFER_DURATION_MS, fooTopic).count()); + assertEquals("Unexpected number of records consumed", numBarRecords, + connect.kafka().consume(numBarRecords, RECORD_TRANSFER_DURATION_MS, barTopic).count()); + + // wait for the connector tasks to consume all records. + connectorHandle.awaitRecords(RECORD_TRANSFER_DURATION_MS); + + // wait for the connector tasks to commit all records. + connectorHandle.awaitCommits(RECORD_TRANSFER_DURATION_MS); + + // Assert all records went to the static topic + Map expectedRecordCounts = singletonMap("static-topic", (long) numFooRecords + numBarRecords); + assertObservedRecords(observedRecords, expectedRecordCounts); + + // Verify all offsets were properly committed + Map offsetMap = connect.kafka().createAdminClient() + .listConsumerGroupOffsets(SinkUtils.consumerGroupId(CONNECTOR_NAME)) + .partitionsToOffsetAndMetadata() + .get(); + + assertEquals(numFooRecords, offsetMap.get(new TopicPartition(fooTopic, 0)).offset()); + assertEquals(numBarRecords, offsetMap.get(new TopicPartition(barTopic, 0)).offset()); + + // delete connector + connect.deleteConnector(CONNECTOR_NAME); + } } 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 7a2a6e4ffa193..4ac7a80bc4df1 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 @@ -18,6 +18,8 @@ import java.util.Arrays; import java.util.Iterator; + +import com.google.common.collect.ImmutableList; import org.apache.kafka.clients.consumer.ConsumerRebalanceListener; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; @@ -87,6 +89,7 @@ import static java.util.Arrays.asList; import static java.util.Collections.singleton; +import static java.util.stream.Collectors.toList; import static org.junit.Assert.assertSame; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; @@ -796,6 +799,70 @@ public void testPreCommit() throws Exception { PowerMock.verifyAll(); } + @Test + public void testPreCommitWithNewTopicName() throws Exception { + String testPrefix = "my-prefix-"; + + createTask(initialState); + + expectInitializeTask(); + expectTaskGetTopic(true); + + // iter 1 + expectPollInitialAssignment(); + + // iter 2 + expectConsumerPoll(2); + expectConversionAndTransformation(2, testPrefix); + + Capture> recordsCapture = EasyMock.newCapture(); + sinkTask.put(EasyMock.capture(recordsCapture)); + EasyMock.expectLastCall(); + + final Map workerCurrentOffsets = new HashMap<>(); + workerCurrentOffsets.put(TOPIC_PARTITION, new OffsetAndMetadata(FIRST_OFFSET + 2)); + workerCurrentOffsets.put(TOPIC_PARTITION2, new OffsetAndMetadata(FIRST_OFFSET)); + + final Map taskOffsets = new HashMap<>(); + taskOffsets.put(TOPIC_PARTITION, new OffsetAndMetadata(FIRST_OFFSET + 1)); // act like FIRST_OFFSET+2 has not yet been flushed by the task + + final Map committableOffsets = new HashMap<>(); + committableOffsets.put(TOPIC_PARTITION, new OffsetAndMetadata(FIRST_OFFSET + 1)); + committableOffsets.put(TOPIC_PARTITION2, new OffsetAndMetadata(FIRST_OFFSET)); + + EasyMock.expect(sinkTask.preCommit(workerCurrentOffsets)).andReturn(taskOffsets); + + final Capture callback = EasyMock.newCapture(); + consumer.commitAsync(EasyMock.eq(committableOffsets), EasyMock.capture(callback)); + EasyMock.expectLastCall().andAnswer(() -> { + callback.getValue().onComplete(committableOffsets, null); + return null; + }); + expectConsumerPoll(0); + sinkTask.put(EasyMock.anyObject()); + EasyMock.expectLastCall(); + + PowerMock.replayAll(); + + workerTask.initialize(TASK_CONFIG); + workerTask.initializeAndStart(); + workerTask.iteration(); // iter 1 -- initial assignment + workerTask.iteration(); // iter 2 -- deliver 2 records + + assertEquals(ImmutableList.of(TOPIC, TOPIC), recordsCapture.getValue().stream() + .map(sr -> sr.originalTopicPartition().topic()) + .collect(toList())); + + assertEquals(ImmutableList.of(testPrefix + TOPIC, testPrefix + TOPIC), recordsCapture.getValue().stream() + .map(SinkRecord::topic) + .collect(toList())); + + sinkTaskContext.getValue().requestCommit(); + workerTask.iteration(); // iter 3 -- commit + + PowerMock.verifyAll(); + } + @Test public void testIgnoredCommit() throws Exception { createTask(initialState); 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 5918747fd1955..5fdd22daf28e2 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 @@ -184,7 +184,7 @@ public void testPollsInBackground() throws Exception { SinkRecord referenceSinkRecord = new SinkRecord(TOPIC, PARTITION, KEY_SCHEMA, KEY, VALUE_SCHEMA, VALUE, FIRST_OFFSET + offset, TIMESTAMP, TIMESTAMP_TYPE); InternalSinkRecord referenceInternalSinkRecord = - new InternalSinkRecord(null, referenceSinkRecord); + new InternalSinkRecord(null, new TopicPartition(TOPIC, PARTITION), referenceSinkRecord); assertEquals(referenceInternalSinkRecord, rec); offset++; } From 83ecba91df8744f58281b2e72eacbe74f471cae1 Mon Sep 17 00:00:00 2001 From: Diego Erdody Date: Wed, 3 Nov 2021 13:54:49 -0700 Subject: [PATCH 2/2] Fix missing dependency --- .../org/apache/kafka/connect/runtime/WorkerSinkTaskTest.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) 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 4ac7a80bc4df1..402532c698e87 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 @@ -19,7 +19,6 @@ import java.util.Arrays; import java.util.Iterator; -import com.google.common.collect.ImmutableList; import org.apache.kafka.clients.consumer.ConsumerRebalanceListener; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; @@ -849,11 +848,11 @@ public void testPreCommitWithNewTopicName() throws Exception { workerTask.iteration(); // iter 1 -- initial assignment workerTask.iteration(); // iter 2 -- deliver 2 records - assertEquals(ImmutableList.of(TOPIC, TOPIC), recordsCapture.getValue().stream() + assertEquals(Arrays.asList(TOPIC, TOPIC), recordsCapture.getValue().stream() .map(sr -> sr.originalTopicPartition().topic()) .collect(toList())); - assertEquals(ImmutableList.of(testPrefix + TOPIC, testPrefix + TOPIC), recordsCapture.getValue().stream() + assertEquals(Arrays.asList(testPrefix + TOPIC, testPrefix + TOPIC), recordsCapture.getValue().stream() .map(SinkRecord::topic) .collect(toList()));