From a6ceb42c4cd360f205bd42f498d698a3b6cc6c57 Mon Sep 17 00:00:00 2001 From: revi cheng Date: Wed, 25 Oct 2023 14:42:32 -0700 Subject: [PATCH] small fixes --- .../connector/internal/streaming/TopicPartitionChannel.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/main/java/com/snowflake/kafka/connector/internal/streaming/TopicPartitionChannel.java b/src/main/java/com/snowflake/kafka/connector/internal/streaming/TopicPartitionChannel.java index cf15a9c85..f13644b9b 100644 --- a/src/main/java/com/snowflake/kafka/connector/internal/streaming/TopicPartitionChannel.java +++ b/src/main/java/com/snowflake/kafka/connector/internal/streaming/TopicPartitionChannel.java @@ -303,7 +303,7 @@ public boolean shouldInsertRecord(SinkRecord kafkaSinkRecord, long currProcessed return true; } - if (kafkaSinkRecord.kafkaOffset() == currProcessedOffset) { + if (kafkaSinkRecord.kafkaOffset() > currProcessedOffset) { LOGGER.debug( "Insert record because kafkaOffset {} > currProcessedOffset {} for channel:{}", kafkaSinkRecord.kafkaOffset(),