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 bda9ecbaa8546..d0e410c386cc5 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 @@ -83,8 +83,10 @@ import org.apache.kafka.connect.storage.HeaderConverter; import org.apache.kafka.connect.storage.KafkaOffsetBackingStore; import org.apache.kafka.connect.storage.OffsetBackingStore; +import org.apache.kafka.connect.storage.OffsetStorageReader; import org.apache.kafka.connect.storage.OffsetStorageReaderImpl; import org.apache.kafka.connect.storage.OffsetStorageWriter; +import org.apache.kafka.connect.storage.OffsetUtils; import org.apache.kafka.connect.util.Callback; import org.apache.kafka.connect.util.ConnectUtils; import org.apache.kafka.connect.util.ConnectorTaskId; @@ -1547,9 +1549,11 @@ void modifySourceConnectorOffsets(String connName, Connector connector, Map, Map> normalizedOffsets = normalizeSourceConnectorOffsets(offsetsToWrite); + boolean alterOffsetsResult; try { - alterOffsetsResult = ((SourceConnector) connector).alterOffsets(connectorConfig, offsetsToWrite); + alterOffsetsResult = ((SourceConnector) connector).alterOffsets(connectorConfig, normalizedOffsets); } catch (UnsupportedOperationException e) { log.error("Failed to modify offsets for connector {} because it doesn't support external modification of offsets", connName, e); @@ -1561,7 +1565,7 @@ void modifySourceConnectorOffsets(String connName, Connector connector, Map offsetWriterCallback = new FutureCallback<>(); offsetWriter.doFlush(offsetWriterCallback); if (config.exactlyOnceSourceEnabled()) { @@ -1608,6 +1612,32 @@ void modifySourceConnectorOffsets(String connName, Connector connector, Map + * Visible for testing. + * + * @param originalOffsets the offsets that are to be normalized + * @return the normalized offsets + */ + @SuppressWarnings("unchecked") + Map, Map> normalizeSourceConnectorOffsets(Map, Map> originalOffsets) { + Map, Map> normalizedOffsets = new HashMap<>(); + for (Map.Entry, Map> entry : originalOffsets.entrySet()) { + OffsetUtils.validateFormat(entry.getKey()); + OffsetUtils.validateFormat(entry.getValue()); + byte[] serializedKey = internalKeyConverter.fromConnectData("", null, entry.getKey()); + byte[] serializedValue = internalValueConverter.fromConnectData("", null, entry.getValue()); + Object deserializedKey = internalKeyConverter.toConnectData("", serializedKey).value(); + Object deserializedValue = internalValueConverter.toConnectData("", serializedValue).value(); + normalizedOffsets.put((Map) deserializedKey, (Map) deserializedValue); + } + + return normalizedOffsets; + } + /** * Update the provided timer, check if it's expired and throw a {@link ConnectException} with the provided error * message if it is. diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java index 646ced752bec2..0de444be2025e 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java @@ -124,6 +124,7 @@ import static org.apache.kafka.clients.consumer.ConsumerConfig.ISOLATION_LEVEL_CONFIG; import static org.apache.kafka.clients.producer.ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG; import static org.apache.kafka.clients.producer.ProducerConfig.TRANSACTIONAL_ID_CONFIG; +import static org.apache.kafka.connect.json.JsonConverterConfig.SCHEMAS_ENABLE_CONFIG; import static org.apache.kafka.connect.runtime.ConnectorConfig.CONNECTOR_CLASS_CONFIG; import static org.apache.kafka.connect.runtime.ConnectorConfig.CONNECTOR_CLIENT_ADMIN_OVERRIDES_PREFIX; import static org.apache.kafka.connect.runtime.ConnectorConfig.CONNECTOR_CLIENT_PRODUCER_OVERRIDES_PREFIX; @@ -1918,6 +1919,7 @@ public void testGetSourceConnectorOffsetsError() { public void testAlterOffsetsConnectorDoesNotSupportOffsetAlteration() { mockKafkaClusterId(); + mockInternalConverters(); worker = new Worker(WORKER_ID, new MockTime(), plugins, config, offsetBackingStore, Executors.newSingleThreadExecutor(), allConnectorClientConfigOverridePolicy, null); worker.start(); @@ -1946,6 +1948,7 @@ public void testAlterOffsetsConnectorDoesNotSupportOffsetAlteration() { @SuppressWarnings("unchecked") public void testAlterOffsetsSourceConnector() throws Exception { mockKafkaClusterId(); + mockInternalConverters(); worker = new Worker(WORKER_ID, new MockTime(), plugins, config, offsetBackingStore, Executors.newSingleThreadExecutor(), allConnectorClientConfigOverridePolicy, null); worker.start(); @@ -1982,6 +1985,7 @@ public void testAlterOffsetsSourceConnector() throws Exception { @SuppressWarnings("unchecked") public void testAlterOffsetsSourceConnectorError() throws Exception { mockKafkaClusterId(); + mockInternalConverters(); worker = new Worker(WORKER_ID, new MockTime(), plugins, config, offsetBackingStore, Executors.newSingleThreadExecutor(), allConnectorClientConfigOverridePolicy, null); worker.start(); @@ -2015,6 +2019,28 @@ public void testAlterOffsetsSourceConnectorError() throws Exception { verifyKafkaClusterId(); } + @Test + public void testNormalizeSourceConnectorOffsets() throws Exception { + Map, Map> offsets = Collections.singletonMap( + Collections.singletonMap("filename", "/path/to/filename"), + Collections.singletonMap("position", 20) + ); + + assertTrue(offsets.values().iterator().next().get("position") instanceof Integer); + + mockInternalConverters(); + + mockKafkaClusterId(); + worker = new Worker(WORKER_ID, new MockTime(), plugins, config, offsetBackingStore, executorService, + noneConnectorClientConfigOverridePolicy, null); + + Map, Map> normalizedOffsets = worker.normalizeSourceConnectorOffsets(offsets); + assertEquals(1, normalizedOffsets.size()); + + // The integer value 20 gets deserialized as a long value by the JsonConverter + assertTrue(normalizedOffsets.values().iterator().next().get("position") instanceof Long); + } + @Test public void testAlterOffsetsSinkConnectorNoDeletes() throws Exception { @SuppressWarnings("unchecked") @@ -2266,6 +2292,7 @@ public void testResetOffsetsSourceConnectorExactlyOnceSupportEnabled() throws Ex workerProps.put("status.storage.topic", "connect-statuses"); config = new DistributedConfig(workerProps); mockKafkaClusterId(); + mockInternalConverters(); worker = new Worker(WORKER_ID, new MockTime(), plugins, config, offsetBackingStore, Executors.newSingleThreadExecutor(), allConnectorClientConfigOverridePolicy, null); worker.start(); @@ -2277,8 +2304,8 @@ public void testResetOffsetsSourceConnectorExactlyOnceSupportEnabled() throws Ex OffsetStorageWriter offsetWriter = mock(OffsetStorageWriter.class); Set> connectorPartitions = new HashSet<>(); - connectorPartitions.add(Collections.singletonMap("partitionKey1", new Object())); - connectorPartitions.add(Collections.singletonMap("partitionKey2", new Object())); + connectorPartitions.add(Collections.singletonMap("partitionKey", "partitionValue1")); + connectorPartitions.add(Collections.singletonMap("partitionKey", "partitionValue2")); when(offsetStore.connectorPartitions(eq(CONNECTOR_ID))).thenReturn(connectorPartitions); when(offsetWriter.doFlush(any())).thenAnswer(invocation -> { invocation.getArgument(0, Callback.class).onCompletion(null, null); @@ -2382,6 +2409,7 @@ public void testResetOffsetsSinkConnectorDeleteConsumerGroupError() throws Excep public void testModifySourceConnectorOffsetsTimeout() throws Exception { mockKafkaClusterId(); Time time = new MockTime(); + mockInternalConverters(); worker = new Worker(WORKER_ID, time, plugins, config, offsetBackingStore, Executors.newSingleThreadExecutor(), allConnectorClientConfigOverridePolicy, null); worker.start(); @@ -2498,14 +2526,14 @@ private void verifyStorage() { } private void mockInternalConverters() { - Converter internalKeyConverter = mock(JsonConverter.class); - Converter internalValueConverter = mock(JsonConverter.class); + JsonConverter jsonConverter = new JsonConverter(); + jsonConverter.configure(Collections.singletonMap(SCHEMAS_ENABLE_CONFIG, false), false); when(plugins.newInternalConverter(eq(true), anyString(), anyMap())) - .thenReturn(internalKeyConverter); + .thenReturn(jsonConverter); when(plugins.newInternalConverter(eq(false), anyString(), anyMap())) - .thenReturn(internalValueConverter); + .thenReturn(jsonConverter); } private void verifyConverters() {