Skip to main content

write_batches

Function write_batches 

Source
fn write_batches<'scope>(
    scope: Scope<'scope, Timestamp>,
    collection_id: GlobalId,
    operator_name: &str,
    target: &CollectionMetadata,
    batch_descriptions: Stream<'scope, Timestamp, Vec<(Antichain<Timestamp>, Antichain<Timestamp>)>>,
    commitments: StreamVec<'scope, Timestamp, Commitment>,
    desired_collection: VecCollection<'scope, Timestamp, Result<Row, DataflowError>, Diff>,
    persist_clients: Arc<PersistClientCache>,
    source_statistics: SourceStatistics,
    busy_signal: Arc<Semaphore>,
) -> (StreamVec<'scope, Timestamp, HollowBatchAndMetadata<Timestamp>>, PressOnDropButton)
Expand description

Writes desired_collection to persist, but only for updates that fall into batch a description that we get via batch_descriptions. This forwards a HollowBatch (with additional metadata) for any batch of updates that was written.

Every update below the ceiling mint_batch_descriptions commits goes into one open builder, whatever its timestamp, so a pinned frontier costs one batch rather than one per timestamp. An update that outruns the ceiling, or arrives while none is committed, has no bound to be grouped under and writes a batch of its own timestamp, which is what the sink does for every update when nothing is ever committed ahead of the frontier.

This operator assumes that the desired_collection comes pre-sharded.

This also and updates various metrics.