From 912db22417f2d0d8b5fdc369a67a6026f682bb4a Mon Sep 17 00:00:00 2001 From: Weny Xu Date: Wed, 15 Jul 2026 11:04:11 +0800 Subject: [PATCH] fix(log-store): deduplicate Kafka WAL multipart records (#8514) * fix(log-store): deduplicate Kafka WAL multipart records Signed-off-by: WenyXu * fix(log-store): preserve delayed Kafka WAL entry offsets Signed-off-by: WenyXu * fix(log-store): reject conflicting Kafka WAL last records Signed-off-by: WenyXu * fix(log-store): emit Kafka WAL entries on last record Signed-off-by: WenyXu * fix(log-store): discard duplicate Kafka WAL first records Signed-off-by: WenyXu --------- Signed-off-by: WenyXu --- src/log-store/src/kafka/util/record.rs | 166 ++++++++++++++++++++----- 1 file changed, 137 insertions(+), 29 deletions(-) diff --git a/src/log-store/src/kafka/util/record.rs b/src/log-store/src/kafka/util/record.rs index 720f989139..aefb098e0d 100644 --- a/src/log-store/src/kafka/util/record.rs +++ b/src/log-store/src/kafka/util/record.rs @@ -227,13 +227,17 @@ pub fn remaining_entries( /// - Emits a [RecordType::Full] type record immediately. /// /// For type of [Entry::MultiplePart] Entry: -/// - Emits a complete or incomplete [Entry] while the next same [RegionId] record arrives. +/// - Emits a complete [Entry] immediately when a [RecordType::Last] record arrives with +/// buffered parts. +/// - Emits an incomplete [Entry] when a new [RecordType::First] record replaces +/// buffered incomplete records for the same [RegionId]. /// -/// **Incomplete Entry:** -/// If the records arrive in the following order, it emits **the incomplete [Entry]** when the next record arrives. -/// - **[RecordType::First], [RecordType::Middle]**, [RecordType::First] -/// - **[RecordType::Middle]**, [RecordType::First] -/// - **[RecordType::Last]** +/// **Discarded records:** +/// - A standalone [RecordType::First] record is discarded when another [RecordType::First] +/// record with the same payload for the same [RegionId] arrives. +/// - A [RecordType::Last] record without buffered parts is discarded. +/// +/// A trailing incomplete entry is emitted by [`remaining_entries`]. pub(crate) fn maybe_emit_entry( provider: &Arc, record: Record, @@ -244,7 +248,16 @@ pub(crate) fn maybe_emit_entry( RecordType::Full => entry = Some(convert_to_naive_entry(provider.clone(), record)), RecordType::First => { let region_id = record.meta.ns.region_id.into(); + let duplicate = buffered_records.get(®ion_id).is_some_and(|records| { + records.len() == 1 + && records[0].meta.tp == RecordType::First + && records[0].data == record.data + }); if let Some(records) = buffered_records.insert(region_id, vec![record]) { + if duplicate { + // A duplicate standalone First cannot form a complete entry. + return Ok(None); + } // Incomplete entry entry = Some(convert_to_multiple_entry( provider.clone(), @@ -261,6 +274,11 @@ pub(crate) fn maybe_emit_entry( if !records.is_empty() { // Safety: the records are guaranteed not empty if the key exists. let last_record = records.last().unwrap(); + if matches!(last_record.meta.tp, RecordType::Middle(last_seq) if last_seq == seq) + && last_record.data == record.data + { + return Ok(None); + } let legal = match last_record.meta.tp { // Legal if this record follows a First record. RecordType::First => seq == 1, @@ -290,14 +308,10 @@ pub(crate) fn maybe_emit_entry( provider.clone(), region_id, records, - )) + )); } else { - // Incomplete entry - entry = Some(convert_to_multiple_entry( - provider.clone(), - region_id, - vec![record], - )) + // Intentionally discard a Last record without buffered parts. It is either a + // duplicate Last record or the tail of an incomplete multipart entry. } } } @@ -359,7 +373,16 @@ mod tests { .unwrap() .is_none() ); - let record = new_test_record(RecordType::First, 2, region_id.as_u64(), vec![2; 100]); + let record = new_test_record(RecordType::First, 2, region_id.as_u64(), vec![1; 100]); + // An identical First is a duplicate and does not emit an incomplete entry. + assert!( + maybe_emit_entry(&provider, record, &mut buffer) + .unwrap() + .is_none() + ); + assert_eq!(buffer[®ion_id][0].meta.entry_id, 2); + + let record = new_test_record(RecordType::First, 3, region_id.as_u64(), vec![2; 100]); let incomplete_entry = maybe_emit_entry(&provider, record, &mut buffer) .unwrap() .unwrap(); @@ -369,29 +392,21 @@ mod tests { Entry::MultiplePart(MultiplePartEntry { provider: Provider::Kafka(provider.clone()), region_id, - entry_id: 1, + entry_id: 2, headers: vec![MultiplePartHeader::First], parts: vec![vec![1; 100]], }) ); - // `Last` overwrite `None` + // `Last` after `None` let mut buffer = HashMap::new(); let record = new_test_record(RecordType::Last, 1, region_id.as_u64(), vec![1; 100]); - let incomplete_entry = maybe_emit_entry(&provider, record, &mut buffer) - .unwrap() - .unwrap(); - - assert_eq!( - incomplete_entry, - Entry::MultiplePart(MultiplePartEntry { - provider: Provider::Kafka(provider.clone()), - region_id, - entry_id: 1, - headers: vec![MultiplePartHeader::Last], - parts: vec![vec![1; 100]], - }) + assert!( + maybe_emit_entry(&provider, record, &mut buffer) + .unwrap() + .is_none() ); + assert!(buffer.is_empty()); // `First` overwrite `Middle(0)` let mut buffer = HashMap::new(); @@ -451,6 +466,99 @@ mod tests { assert_matches!(err, error::Error::IllegalSequence { .. }); } + #[test] + fn test_maybe_emit_entry_deduplicates_multipart_records() { + let provider = Arc::new(KafkaProvider::new("my_topic".to_string())); + let region_id = RegionId::new(1, 1); + let mut buffer = HashMap::new(); + + for record in [ + new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]), + new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![2; 100]), + new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![2; 100]), + ] { + assert!( + maybe_emit_entry(&provider, record, &mut buffer) + .unwrap() + .is_none() + ); + } + + let last = new_test_record(RecordType::Last, 1, region_id.as_u64(), vec![3; 100]); + let entry = maybe_emit_entry(&provider, last, &mut buffer) + .unwrap() + .unwrap(); + + assert_eq!( + entry, + Entry::MultiplePart(MultiplePartEntry { + provider: Provider::Kafka(provider.clone()), + region_id, + entry_id: 1, + headers: vec![ + MultiplePartHeader::First, + MultiplePartHeader::Middle(1), + MultiplePartHeader::Last, + ], + parts: vec![vec![1; 100], vec![2; 100], vec![3; 100]], + }) + ); + let duplicate = new_test_record(RecordType::Last, 1, region_id.as_u64(), vec![3; 100]); + assert!( + maybe_emit_entry(&provider, duplicate, &mut buffer) + .unwrap() + .is_none() + ); + } + + #[test] + fn test_maybe_emit_entry_rejects_conflicting_duplicate_middle_record() { + let provider = Arc::new(KafkaProvider::new("my_topic".to_string())); + let region_id = RegionId::new(1, 1); + let mut buffer = HashMap::new(); + + for record in [ + new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]), + new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![2; 100]), + ] { + assert!( + maybe_emit_entry(&provider, record, &mut buffer) + .unwrap() + .is_none() + ); + } + + let duplicate = new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![3; 100]); + let err = maybe_emit_entry(&provider, duplicate, &mut buffer).unwrap_err(); + + assert_matches!(err, error::Error::IllegalSequence { .. }); + } + + #[test] + fn test_maybe_emit_entry_preserves_multipart_entry_order_before_full_record() { + let provider = Arc::new(KafkaProvider::new("my_topic".to_string())); + let region_id = RegionId::new(1, 1); + let mut buffer = HashMap::new(); + + let first = new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]); + assert!( + maybe_emit_entry(&provider, first, &mut buffer) + .unwrap() + .is_none() + ); + let last = new_test_record(RecordType::Last, 2, region_id.as_u64(), vec![2; 100]); + let multipart_entry = maybe_emit_entry(&provider, last, &mut buffer) + .unwrap() + .unwrap(); + let full = new_test_record(RecordType::Full, 3, region_id.as_u64(), vec![3; 100]); + let full_entry = maybe_emit_entry(&provider, full, &mut buffer) + .unwrap() + .unwrap(); + + assert_eq!(multipart_entry.entry_id(), 2); + assert_eq!(full_entry.entry_id(), 3); + } + #[test] fn test_meta_size() { let meta = RecordMeta {