Skip to main content

maybe_emit_entry

Function maybe_emit_entry 

Source
pub(crate) fn maybe_emit_entry(
    provider: &Arc<KafkaProvider>,
    record: Record,
    buffered_records: &mut HashMap<RegionId, Vec<Record>>,
) -> Result<Option<Entry>>
Expand description

For type of [Entry::Naive] Entry:

For type of [Entry::MultiplePart] Entry:

  • 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].

Discarded records:

A trailing incomplete entry is emitted by remaining_entries.