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.