Skip to main content

mint_batch_descriptions

Function mint_batch_descriptions 

Source
fn mint_batch_descriptions<'scope>(
    scope: Scope<'scope, Timestamp>,
    collection_id: GlobalId,
    operator_name: &str,
    target: &CollectionMetadata,
    desired_collection: VecCollection<'scope, Timestamp, Result<Row, DataflowError>, Diff>,
    remap_upper: StreamVec<'scope, Timestamp, ()>,
    persist_clients: Arc<PersistClientCache>,
    lookahead: Option<u64>,
    snapshot_time: Option<Antichain<Timestamp>>,
) -> (StreamVec<'scope, Timestamp, (Antichain<Timestamp>, Antichain<Timestamp>)>, StreamVec<'scope, Timestamp, Commitment>, StreamVec<'scope, Timestamp, (Result<Row, DataflowError>, Timestamp, Diff)>, PressOnDropButton)
Expand description

Whenever the frontier advances, this mints a new batch description (lower and upper) that writers should use for writing the next set of batches to persist.

With a lookahead, and while the frontier sits at snapshot_time, it also commits to a ceiling that far past the data and broadcasts it on the second output. A ceiling is not a description: it gives the writers a bound to group updates under before the frontier certifies anything, and it binds this operator, which mints nothing below an outstanding ceiling. So the whole snapshot and the catch-up behind it become one description, emitted when the frontier reaches the ceiling. See next_mint.

Only one of the workers does this, meaning there will only be one description in the stream, even in case of multiple timely workers. Use broadcast() to, ahem, broadcast, the one description to all downstream write operators/workers.

remap_upper carries the ingestion’s remap upper as its frontier, which paces the commitments. Reclocking stamps every update below that upper before the update exists, so a ceiling ahead of it leads every row that can still arrive, however long this export’s own data has been quiet.