Skip to main content

log_store/kafka/util/
record.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::HashMap;
16use std::sync::Arc;
17
18use rskafka::record::Record as KafkaRecord;
19use serde::{Deserialize, Serialize};
20use snafu::{OptionExt, ResultExt, ensure};
21use store_api::logstore::entry::{Entry, MultiplePartEntry, MultiplePartHeader, NaiveEntry};
22use store_api::logstore::provider::{KafkaProvider, Provider};
23use store_api::storage::RegionId;
24
25use crate::error::{
26    DecodeJsonSnafu, EncodeJsonSnafu, IllegalSequenceSnafu, MetaLengthExceededLimitSnafu,
27    MissingKeySnafu, MissingValueSnafu, Result,
28};
29use crate::kafka::{EntryId, NamespaceImpl};
30
31/// The current version of Record.
32pub(crate) const VERSION: u32 = 0;
33
34/// The estimated size in bytes of a serialized RecordMeta.
35/// A record is guaranteed to have sizeof(meta) + sizeof(data) <= max_batch_byte - ESTIMATED_META_SIZE.
36pub(crate) const ESTIMATED_META_SIZE: usize = 256;
37
38/// The type of a record.
39///
40/// - If the entry is able to fit into a Kafka record, it's converted into a Full record.
41///
42/// - If the entry is too large to fit into a Kafka record, it's converted into a collection of records.
43///
44/// Those records must contain exactly one First record and one Last record, and potentially several
45///   Middle records. There may be no Middle record.
46#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq)]
47pub enum RecordType {
48    /// The record is self-contained, i.e. an entry's data is fully stored into this record.
49    Full,
50    /// The record contains the first part of an entry's data.
51    First,
52    /// The record contains one of the middle parts of an entry's data.
53    /// The sequence of the record is identified by the inner field.
54    Middle(usize),
55    /// The record contains the last part of an entry's data.
56    Last,
57}
58
59/// The metadata of a record.
60#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
61pub struct RecordMeta {
62    /// The version of the record. Used for backward compatibility.
63    version: u32,
64    /// The type of the record.
65    pub tp: RecordType,
66    /// The id of the entry the record associated with.
67    pub entry_id: EntryId,
68    /// The namespace of the entry the record associated with.
69    pub ns: NamespaceImpl,
70}
71
72/// The minimal storage unit in the Kafka log store.
73///
74/// An entry will be first converted into several Records before producing.
75/// If an entry is able to fit into a KafkaRecord, it converts to a single Record.
76/// If otherwise an entry cannot fit into a KafkaRecord, it will be split into a collection of Records.
77///
78/// A KafkaRecord is the minimal storage unit used by Kafka client and Kafka server.
79/// The Kafka client produces KafkaRecords and consumes KafkaRecords, and Kafka server stores
80/// a collection of KafkaRecords.
81#[derive(Debug, Clone, PartialEq)]
82pub(crate) struct Record {
83    /// The metadata of the record.
84    pub(crate) meta: RecordMeta,
85    /// The payload of the record.
86    data: Vec<u8>,
87}
88
89impl TryFrom<Record> for KafkaRecord {
90    type Error = crate::error::Error;
91
92    fn try_from(record: Record) -> Result<Self> {
93        let key = serde_json::to_vec(&record.meta).context(EncodeJsonSnafu)?;
94        ensure!(
95            key.len() < ESTIMATED_META_SIZE,
96            MetaLengthExceededLimitSnafu {
97                limit: ESTIMATED_META_SIZE,
98                actual: key.len()
99            }
100        );
101        Ok(KafkaRecord {
102            key: Some(key),
103            value: Some(record.data),
104            timestamp: chrono::Utc::now(),
105            headers: Default::default(),
106        })
107    }
108}
109
110// TODO(niebayes): improve the performance of decoding kafka record.
111impl TryFrom<KafkaRecord> for Record {
112    type Error = crate::error::Error;
113
114    fn try_from(kafka_record: KafkaRecord) -> Result<Self> {
115        let key = kafka_record.key.context(MissingKeySnafu)?;
116        let meta = serde_json::from_slice(&key).context(DecodeJsonSnafu)?;
117        let data = kafka_record.value.context(MissingValueSnafu)?;
118        Ok(Self { meta, data })
119    }
120}
121
122pub(crate) fn convert_to_kafka_records(entry: Entry) -> Result<Vec<KafkaRecord>> {
123    match entry {
124        Entry::Naive(entry) => Ok(vec![KafkaRecord::try_from(Record {
125            meta: RecordMeta {
126                version: VERSION,
127                tp: RecordType::Full,
128                // TODO(weny): refactor the record meta.
129                entry_id: 0,
130                ns: NamespaceImpl {
131                    region_id: entry.region_id.as_u64(),
132                    // TODO(weny): refactor the record meta.
133                    topic: String::new(),
134                },
135            },
136            data: entry.data,
137        })?]),
138        Entry::MultiplePart(entry) => {
139            let mut entries = Vec::with_capacity(entry.parts.len());
140
141            for (idx, part) in entry.parts.into_iter().enumerate() {
142                let tp = match entry.headers[idx] {
143                    MultiplePartHeader::First => RecordType::First,
144                    MultiplePartHeader::Middle(i) => RecordType::Middle(i),
145                    MultiplePartHeader::Last => RecordType::Last,
146                };
147                entries.push(KafkaRecord::try_from(Record {
148                    meta: RecordMeta {
149                        version: VERSION,
150                        tp,
151                        // TODO(weny): refactor the record meta.
152                        entry_id: 0,
153                        ns: NamespaceImpl {
154                            region_id: entry.region_id.as_u64(),
155                            topic: String::new(),
156                        },
157                    },
158                    data: part,
159                })?)
160            }
161            Ok(entries)
162        }
163    }
164}
165
166fn convert_to_naive_entry(provider: Arc<KafkaProvider>, record: Record) -> Entry {
167    let region_id = RegionId::from_u64(record.meta.ns.region_id);
168
169    Entry::Naive(NaiveEntry {
170        provider: Provider::Kafka(provider),
171        region_id,
172        entry_id: record.meta.entry_id,
173        data: record.data,
174    })
175}
176
177fn convert_to_multiple_entry(
178    provider: Arc<KafkaProvider>,
179    region_id: RegionId,
180    records: Vec<Record>,
181) -> Entry {
182    let mut headers = Vec::with_capacity(records.len());
183    let mut parts = Vec::with_capacity(records.len());
184    let entry_id = records.last().map(|r| r.meta.entry_id).unwrap_or_default();
185
186    for record in records {
187        let header = match record.meta.tp {
188            RecordType::Full => unreachable!(),
189            RecordType::First => MultiplePartHeader::First,
190            RecordType::Middle(i) => MultiplePartHeader::Middle(i),
191            RecordType::Last => MultiplePartHeader::Last,
192        };
193        headers.push(header);
194        parts.push(record.data);
195    }
196
197    Entry::MultiplePart(MultiplePartEntry {
198        provider: Provider::Kafka(provider),
199        region_id,
200        entry_id,
201        headers,
202        parts,
203    })
204}
205
206/// Constructs entries from `buffered_records`
207pub fn remaining_entries(
208    provider: &Arc<KafkaProvider>,
209    buffered_records: &mut HashMap<RegionId, Vec<Record>>,
210) -> Option<Vec<Entry>> {
211    if buffered_records.is_empty() {
212        None
213    } else {
214        let mut entries = Vec::with_capacity(buffered_records.len());
215        for (region_id, records) in buffered_records.drain() {
216            entries.push(convert_to_multiple_entry(
217                provider.clone(),
218                region_id,
219                records,
220            ));
221        }
222        Some(entries)
223    }
224}
225
226/// For type of [Entry::Naive] Entry:
227/// - Emits a [RecordType::Full] type record immediately.
228///
229/// For type of [Entry::MultiplePart] Entry:
230/// - Emits a complete [Entry] immediately when a [RecordType::Last] record arrives with
231///   buffered parts.
232/// - Emits an incomplete [Entry] when a new [RecordType::First] record replaces
233///   buffered incomplete records for the same [RegionId].
234///
235/// **Discarded records:**
236/// - A standalone [RecordType::First] record is discarded when another [RecordType::First]
237///   record with the same payload for the same [RegionId] arrives.
238/// - A [RecordType::Last] record without buffered parts is discarded.
239///
240/// A trailing incomplete entry is emitted by [`remaining_entries`].
241pub(crate) fn maybe_emit_entry(
242    provider: &Arc<KafkaProvider>,
243    record: Record,
244    buffered_records: &mut HashMap<RegionId, Vec<Record>>,
245) -> Result<Option<Entry>> {
246    let mut entry = None;
247    match record.meta.tp {
248        RecordType::Full => entry = Some(convert_to_naive_entry(provider.clone(), record)),
249        RecordType::First => {
250            let region_id = record.meta.ns.region_id.into();
251            let duplicate = buffered_records.get(&region_id).is_some_and(|records| {
252                records.len() == 1
253                    && records[0].meta.tp == RecordType::First
254                    && records[0].data == record.data
255            });
256            if let Some(records) = buffered_records.insert(region_id, vec![record]) {
257                if duplicate {
258                    // A duplicate standalone First cannot form a complete entry.
259                    return Ok(None);
260                }
261                // Incomplete entry
262                entry = Some(convert_to_multiple_entry(
263                    provider.clone(),
264                    region_id,
265                    records,
266                ))
267            }
268        }
269        RecordType::Middle(seq) => {
270            let region_id = record.meta.ns.region_id.into();
271            let records = buffered_records.entry(region_id).or_default();
272
273            // Only validate complete entries.
274            if !records.is_empty() {
275                // Safety: the records are guaranteed not empty if the key exists.
276                let last_record = records.last().unwrap();
277                if matches!(last_record.meta.tp, RecordType::Middle(last_seq) if last_seq == seq)
278                    && last_record.data == record.data
279                {
280                    return Ok(None);
281                }
282                let legal = match last_record.meta.tp {
283                    // Legal if this record follows a First record.
284                    RecordType::First => seq == 1,
285                    // Legal if this record follows a Middle record just prior to this record.
286                    RecordType::Middle(last_seq) => last_seq + 1 == seq,
287                    // Illegal sequence.
288                    _ => false,
289                };
290                ensure!(
291                    legal,
292                    IllegalSequenceSnafu {
293                        error: format!(
294                            "Illegal sequence of a middle record, last record: {:?}, incoming record: {:?}",
295                            last_record.meta.tp, record.meta.tp
296                        )
297                    }
298                );
299            }
300
301            records.push(record);
302        }
303        RecordType::Last => {
304            let region_id = record.meta.ns.region_id.into();
305            if let Some(mut records) = buffered_records.remove(&region_id) {
306                records.push(record);
307                entry = Some(convert_to_multiple_entry(
308                    provider.clone(),
309                    region_id,
310                    records,
311                ));
312            } else {
313                // Intentionally discard a Last record without buffered parts. It is either a
314                // duplicate Last record or the tail of an incomplete multipart entry.
315            }
316        }
317    }
318    Ok(entry)
319}
320
321#[cfg(test)]
322mod tests {
323    use std::assert_matches;
324    use std::sync::Arc;
325
326    use super::*;
327    use crate::error;
328
329    fn new_test_record(tp: RecordType, entry_id: EntryId, region_id: u64, data: Vec<u8>) -> Record {
330        Record {
331            meta: RecordMeta {
332                version: VERSION,
333                tp,
334                ns: NamespaceImpl {
335                    region_id,
336                    topic: "greptimedb_wal_topic".to_string(),
337                },
338                entry_id,
339            },
340            data,
341        }
342    }
343
344    #[test]
345    fn test_maybe_emit_entry_emit_naive_entry() {
346        let provider = Arc::new(KafkaProvider::new("my_topic".to_string()));
347        let region_id = RegionId::new(1, 1);
348        let mut buffer = HashMap::new();
349        let record = new_test_record(RecordType::Full, 1, region_id.as_u64(), vec![1; 100]);
350        let entry = maybe_emit_entry(&provider, record, &mut buffer)
351            .unwrap()
352            .unwrap();
353        assert_eq!(
354            entry,
355            Entry::Naive(NaiveEntry {
356                provider: Provider::Kafka(provider),
357                region_id,
358                entry_id: 1,
359                data: vec![1; 100]
360            })
361        );
362    }
363
364    #[test]
365    fn test_maybe_emit_entry_emit_incomplete_entry() {
366        let provider = Arc::new(KafkaProvider::new("my_topic".to_string()));
367        let region_id = RegionId::new(1, 1);
368        // `First` overwrite `First`
369        let mut buffer = HashMap::new();
370        let record = new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]);
371        assert!(
372            maybe_emit_entry(&provider, record, &mut buffer)
373                .unwrap()
374                .is_none()
375        );
376        let record = new_test_record(RecordType::First, 2, region_id.as_u64(), vec![1; 100]);
377        // An identical First is a duplicate and does not emit an incomplete entry.
378        assert!(
379            maybe_emit_entry(&provider, record, &mut buffer)
380                .unwrap()
381                .is_none()
382        );
383        assert_eq!(buffer[&region_id][0].meta.entry_id, 2);
384
385        let record = new_test_record(RecordType::First, 3, region_id.as_u64(), vec![2; 100]);
386        let incomplete_entry = maybe_emit_entry(&provider, record, &mut buffer)
387            .unwrap()
388            .unwrap();
389
390        assert_eq!(
391            incomplete_entry,
392            Entry::MultiplePart(MultiplePartEntry {
393                provider: Provider::Kafka(provider.clone()),
394                region_id,
395                entry_id: 2,
396                headers: vec![MultiplePartHeader::First],
397                parts: vec![vec![1; 100]],
398            })
399        );
400
401        // `Last` after `None`
402        let mut buffer = HashMap::new();
403        let record = new_test_record(RecordType::Last, 1, region_id.as_u64(), vec![1; 100]);
404        assert!(
405            maybe_emit_entry(&provider, record, &mut buffer)
406                .unwrap()
407                .is_none()
408        );
409        assert!(buffer.is_empty());
410
411        // `First` overwrite `Middle(0)`
412        let mut buffer = HashMap::new();
413        let record = new_test_record(RecordType::Middle(0), 1, region_id.as_u64(), vec![1; 100]);
414        assert!(
415            maybe_emit_entry(&provider, record, &mut buffer)
416                .unwrap()
417                .is_none()
418        );
419        let record = new_test_record(RecordType::First, 2, region_id.as_u64(), vec![2; 100]);
420        let incomplete_entry = maybe_emit_entry(&provider, record, &mut buffer)
421            .unwrap()
422            .unwrap();
423
424        assert_eq!(
425            incomplete_entry,
426            Entry::MultiplePart(MultiplePartEntry {
427                provider: Provider::Kafka(provider),
428                region_id,
429                entry_id: 1,
430                headers: vec![MultiplePartHeader::Middle(0)],
431                parts: vec![vec![1; 100]],
432            })
433        );
434    }
435
436    #[test]
437    fn test_maybe_emit_entry_illegal_seq() {
438        let provider = Arc::new(KafkaProvider::new("my_topic".to_string()));
439        let region_id = RegionId::new(1, 1);
440        let mut buffer = HashMap::new();
441        let record = new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]);
442        assert!(
443            maybe_emit_entry(&provider, record, &mut buffer)
444                .unwrap()
445                .is_none()
446        );
447        let record = new_test_record(RecordType::Middle(2), 1, region_id.as_u64(), vec![2; 100]);
448        let err = maybe_emit_entry(&provider, record, &mut buffer).unwrap_err();
449        assert_matches!(err, error::Error::IllegalSequence { .. });
450
451        let mut buffer = HashMap::new();
452        let record = new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]);
453        assert!(
454            maybe_emit_entry(&provider, record, &mut buffer)
455                .unwrap()
456                .is_none()
457        );
458        let record = new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![2; 100]);
459        assert!(
460            maybe_emit_entry(&provider, record, &mut buffer)
461                .unwrap()
462                .is_none()
463        );
464        let record = new_test_record(RecordType::Middle(3), 1, region_id.as_u64(), vec![2; 100]);
465        let err = maybe_emit_entry(&provider, record, &mut buffer).unwrap_err();
466        assert_matches!(err, error::Error::IllegalSequence { .. });
467    }
468
469    #[test]
470    fn test_maybe_emit_entry_deduplicates_multipart_records() {
471        let provider = Arc::new(KafkaProvider::new("my_topic".to_string()));
472        let region_id = RegionId::new(1, 1);
473        let mut buffer = HashMap::new();
474
475        for record in [
476            new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]),
477            new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![2; 100]),
478            new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![2; 100]),
479        ] {
480            assert!(
481                maybe_emit_entry(&provider, record, &mut buffer)
482                    .unwrap()
483                    .is_none()
484            );
485        }
486
487        let last = new_test_record(RecordType::Last, 1, region_id.as_u64(), vec![3; 100]);
488        let entry = maybe_emit_entry(&provider, last, &mut buffer)
489            .unwrap()
490            .unwrap();
491
492        assert_eq!(
493            entry,
494            Entry::MultiplePart(MultiplePartEntry {
495                provider: Provider::Kafka(provider.clone()),
496                region_id,
497                entry_id: 1,
498                headers: vec![
499                    MultiplePartHeader::First,
500                    MultiplePartHeader::Middle(1),
501                    MultiplePartHeader::Last,
502                ],
503                parts: vec![vec![1; 100], vec![2; 100], vec![3; 100]],
504            })
505        );
506        let duplicate = new_test_record(RecordType::Last, 1, region_id.as_u64(), vec![3; 100]);
507        assert!(
508            maybe_emit_entry(&provider, duplicate, &mut buffer)
509                .unwrap()
510                .is_none()
511        );
512    }
513
514    #[test]
515    fn test_maybe_emit_entry_rejects_conflicting_duplicate_middle_record() {
516        let provider = Arc::new(KafkaProvider::new("my_topic".to_string()));
517        let region_id = RegionId::new(1, 1);
518        let mut buffer = HashMap::new();
519
520        for record in [
521            new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]),
522            new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![2; 100]),
523        ] {
524            assert!(
525                maybe_emit_entry(&provider, record, &mut buffer)
526                    .unwrap()
527                    .is_none()
528            );
529        }
530
531        let duplicate = new_test_record(RecordType::Middle(1), 1, region_id.as_u64(), vec![3; 100]);
532        let err = maybe_emit_entry(&provider, duplicate, &mut buffer).unwrap_err();
533
534        assert_matches!(err, error::Error::IllegalSequence { .. });
535    }
536
537    #[test]
538    fn test_maybe_emit_entry_preserves_multipart_entry_order_before_full_record() {
539        let provider = Arc::new(KafkaProvider::new("my_topic".to_string()));
540        let region_id = RegionId::new(1, 1);
541        let mut buffer = HashMap::new();
542
543        let first = new_test_record(RecordType::First, 1, region_id.as_u64(), vec![1; 100]);
544        assert!(
545            maybe_emit_entry(&provider, first, &mut buffer)
546                .unwrap()
547                .is_none()
548        );
549        let last = new_test_record(RecordType::Last, 2, region_id.as_u64(), vec![2; 100]);
550        let multipart_entry = maybe_emit_entry(&provider, last, &mut buffer)
551            .unwrap()
552            .unwrap();
553        let full = new_test_record(RecordType::Full, 3, region_id.as_u64(), vec![3; 100]);
554        let full_entry = maybe_emit_entry(&provider, full, &mut buffer)
555            .unwrap()
556            .unwrap();
557
558        assert_eq!(multipart_entry.entry_id(), 2);
559        assert_eq!(full_entry.entry_id(), 3);
560    }
561
562    #[test]
563    fn test_meta_size() {
564        let meta = RecordMeta {
565            version: VERSION,
566            tp: RecordType::Middle(usize::MAX),
567            entry_id: u64::MAX,
568            ns: NamespaceImpl {
569                region_id: RegionId::new(u32::MAX, u32::MAX).as_u64(),
570                topic: format!("greptime_kafka_cluster/1024/2048/{}", uuid::Uuid::new_v4()),
571            },
572        };
573        let serialized = serde_json::to_vec(&meta).unwrap();
574        // The len of serialized data is 202.
575        assert!(serialized.len() < ESTIMATED_META_SIZE);
576    }
577}