fix(log-store): deduplicate Kafka WAL multipart records (#8514)

* fix(log-store): deduplicate Kafka WAL multipart records

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fix(log-store): preserve delayed Kafka WAL entry offsets

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fix(log-store): reject conflicting Kafka WAL last records

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fix(log-store): emit Kafka WAL entries on last record

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fix(log-store): discard duplicate Kafka WAL first records

Signed-off-by: WenyXu <wenymedia@gmail.com>

---------

Signed-off-by: WenyXu <wenymedia@gmail.com>
This commit is contained in:
Weny Xu
2026-07-15 11:04:11 +08:00
committed by GitHub
parent 6e8bd8a5f9
commit 912db22417
+137 -29
View File
@@ -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<KafkaProvider>,
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(&region_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[&region_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 {