From 26fc0543ed13b41b0b7f14b7d38e70b7930dc244 Mon Sep 17 00:00:00 2001 From: Song Huang Date: Wed, 22 Jul 2026 21:41:19 -0400 Subject: [PATCH] kafka: ignore tombstone payloads --- crates/adapters/src/test/kafka.rs | 6 +++ .../adapters/src/transport/kafka/ft/input.rs | 8 ++-- .../adapters/src/transport/kafka/ft/test.rs | 39 +++++++++++++++++++ 3 files changed, 50 insertions(+), 3 deletions(-) diff --git a/crates/adapters/src/test/kafka.rs b/crates/adapters/src/test/kafka.rs index a7b269019de..81e2656f480 100644 --- a/crates/adapters/src/test/kafka.rs +++ b/crates/adapters/src/test/kafka.rs @@ -242,6 +242,12 @@ impl TestProducer { self.producer.flush(Timeout::Never).unwrap(); } + pub fn send_tombstone(&self, key: &[u8], topic: &str) { + let record = BaseRecord::<[u8], (), ()>::to(topic).key(key); + self.producer.send(record).unwrap(); + self.producer.flush(Timeout::Never).unwrap(); + } + pub fn send_to_topic_partition(&self, data: &[Vec], topic: &str, partition: i32) { for batch in data { let mut writer = CsvWriterBuilder::new() diff --git a/crates/adapters/src/transport/kafka/ft/input.rs b/crates/adapters/src/transport/kafka/ft/input.rs index 0da96a08fc1..f21144c00bd 100644 --- a/crates/adapters/src/transport/kafka/ft/input.rs +++ b/crates/adapters/src/transport/kafka/ft/input.rs @@ -1255,9 +1255,11 @@ impl PartitionReceiver { } let timestamp = message.timestamp().to_millis().unwrap_or(i64::MIN); - let payload = message.payload().unwrap_or(&[]); - let metadata = self.create_metadata(&message); - let (buffer, errors) = parser.parse(payload, metadata); + let (buffer, errors) = match message.payload() { + Some(payload) => parser.parse(payload, self.create_metadata(&message)), + // Skip tombstones; keep offsets. + None => (None, Vec::new()), + }; self.n_bytes .fetch_add(buffer.len().bytes, Ordering::Relaxed); self.messages diff --git a/crates/adapters/src/transport/kafka/ft/test.rs b/crates/adapters/src/transport/kafka/ft/test.rs index c4cffda5424..340d618489b 100644 --- a/crates/adapters/src/transport/kafka/ft/test.rs +++ b/crates/adapters/src/transport/kafka/ft/test.rs @@ -509,6 +509,45 @@ fn test_input(topic: &str, batch_sizes: &[u32]) { } } +#[test] +fn test_input_tombstone() { + init_test_logger(); + + let topic = "kafka_tombstone_ft"; + let _kafka_resources = KafkaResources::create_topics(&[(topic, 1)]); + + let (endpoint, receiver, reader) = create_reader(topic, None, false); + reader.extend(); + + let producer = TestProducer::new(); + producer.send_message(b"m0", topic, None); + producer.send_tombstone(b"deleted-key", topic); + producer.send_message(b"", topic, None); + producer.send_message(b"m3", topic, None); + + receiver.expect_buffering(3); + reader.queue(false); + + let metadata = Metadata { + offsets: vec![0..4], + }; + receiver.expect(vec![ConsumerCall::Extended { + num_records: 3, + metadata: serde_json::to_value(&metadata).unwrap(), + }]); + assert_eq!(take_flushed(&receiver), vec!["m0", "", "m3"]); + + drop(endpoint); + drop(receiver); + drop(reader); + + let (_endpoint, receiver, reader) = create_reader(topic, None, false); + receiver.inner.drop_buffered.store(true, Ordering::Release); + reader.replay(serde_json::to_value(&metadata).unwrap(), RmpValue::Nil); + receiver.expect(vec![ConsumerCall::Replayed { num_records: 3 }]); + assert_eq!(take_flushed(&receiver), vec!["m0", "", "m3"]); +} + /// A configured header filter drops non-matching messages before parsing, and /// the drop decision is deterministic across replay: a fresh reader replaying /// the checkpointed offset range reproduces exactly the admitted records.