From 3f700a94f8a5acfc3dfae8f618e0e0f52287333c Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Wed, 16 Aug 2023 13:00:05 -0400 Subject: [PATCH] MINOR: Permit tombstone offsets with any partition for file source connector --- .../connect/file/FileStreamSourceConnector.java | 13 ++++++++----- .../file/FileStreamSourceConnectorTest.java | 15 +++++++++++++++ 2 files changed, 23 insertions(+), 5 deletions(-) diff --git a/connect/file/src/main/java/org/apache/kafka/connect/file/FileStreamSourceConnector.java b/connect/file/src/main/java/org/apache/kafka/connect/file/FileStreamSourceConnector.java index 13193f8f50124..35a1c07aabe1f 100644 --- a/connect/file/src/main/java/org/apache/kafka/connect/file/FileStreamSourceConnector.java +++ b/connect/file/src/main/java/org/apache/kafka/connect/file/FileStreamSourceConnector.java @@ -117,6 +117,14 @@ public boolean alterOffsets(Map connectorConfig, Map, Map> partitionOffset : offsets.entrySet()) { + Map offset = partitionOffset.getValue(); + // null offsets are allowed and represent a deletion of offsets for a partition + // allow tombstones for anything; if there's garbage in the offsets for the connector, we don't + // want to prevent users from being able to clean it up using the REST API + if (offset == null) { + continue; + } + Map partition = partitionOffset.getKey(); if (partition == null) { throw new ConnectException("Partition objects cannot be null"); @@ -126,11 +134,6 @@ public boolean alterOffsets(Map connectorConfig, Map offset = partitionOffset.getValue(); - // null offsets are allowed and represent a deletion of offsets for a partition - if (offset == null) { - continue; - } if (!offset.containsKey(POSITION_FIELD)) { throw new ConnectException("Offset objects should either be null or contain the key '" + POSITION_FIELD + "'"); diff --git a/connect/file/src/test/java/org/apache/kafka/connect/file/FileStreamSourceConnectorTest.java b/connect/file/src/test/java/org/apache/kafka/connect/file/FileStreamSourceConnectorTest.java index 185faa80eb34a..be5019f9c5d8a 100644 --- a/connect/file/src/test/java/org/apache/kafka/connect/file/FileStreamSourceConnectorTest.java +++ b/connect/file/src/test/java/org/apache/kafka/connect/file/FileStreamSourceConnectorTest.java @@ -213,6 +213,21 @@ public void testAlterOffsetsOffsetPositionValues() { assertThrows(ConnectException.class, () -> alterOffsets.apply(-10L)); assertTrue(() -> alterOffsets.apply(10L)); } + @Test + public void testAlterOffsetsOffsetTombstones() { + Function, Boolean> alterOffsets = partition -> + connector.alterOffsets(sourceProperties, Collections.singletonMap(partition, null)); + + assertTrue(alterOffsets.apply(null)); + assertTrue(alterOffsets.apply(Collections.emptyMap())); + Map partition = new HashMap<>(); + partition.put("unused_partition_key", "unused_partition_value"); + assertTrue(alterOffsets.apply(partition)); + partition.put(FILENAME_FIELD, FILENAME); + assertTrue(alterOffsets.apply(partition)); + partition.put("", ""); + assertTrue(alterOffsets.apply(partition)); + } @Test public void testSuccessfulAlterOffsets() {