Expand description
Render an operator that persists a source collection.
§Implementation
This module defines the persist_sink operator, that writes
a collection produced by source rendering into a persist shard.
It attempts to use all workers to write data to persist, and uses single-instance workers to coordinate work. The below diagram is an overview how it it shaped. There is more information in the doc comments of the top-level functions of this module.
,------------.
| source |
| collection |
+---+--------+
/ |
/ |
/ |
/ |
/ |
/ |
/ |
/ |
/ ,-+-----------------------.
/ | mint_batch_descriptions |
/ | one arbitrary worker |
| +-,--,--------+----+------+
,----------´.-´ | \
_.-´ | .-´ | \
_.-´ | .-´ | \
.-´ .------+----|-------+---------|--------\-----.
/ / | | | \ \
,--------------. ,-----------------. | ,-----------------.
| write_batches| | write_batches | | | write_batches |
| worker 0 | | worker 1 | | | worker N |
+-----+--------+ +-+---------------+ | +--+--------------+
\ \ | /
`-. `, | /
`-._ `-. | /
`-._ `-. | /
`---------. `-. | /
+`---`---+-------------,
| append_batches |
| one arbitrary worker |
+------+---------------+mint_batch_descriptions also takes the ingestion’s remap upper as a disconnected input, and
broadcasts a committed ceiling to the writers, which append_batches does not see. Both edges
are left out of the diagram to avoid clutter.
§Similarities with mz_compute::sink::persist_sink
This module has many similarities with the compute version of the same concept, and in fact, is entirely derived from it.
Compute requires that its persist_sink is self-correcting;
that is, it corrects what the collection in persist
accumulates to if the collection has values changed at
previous timestamps. It does this by continually comparing
the input stream with the collection as read back from persist.
Source collections, while definite, cannot be reliably by
re-produced once written down, which means compute’s
persist_sink’s self-correction mechanism would need to be
skipped on operator startup, and would cause unnecessary read
load on persist.
Additionally, persisting sources requires we use bounded amounts of memory, even if a single timestamp represents a huge amount of data. This is not (currently) possible to guarantee while also performing self-correction.
Because of this, we have ripped out the self-correction
mechanism, and aggressively simplified the sub-operators.
Some, particularly append_batches could be merged with
the compute version, but that requires some amount of
onerous refactoring that we have chosen to skip for now.
Structs§
- Batch
Builder 🔒AndMetadata - Manages batches and metrics.
- Batch
Metrics 🔒 - Metrics about batches.
- Batch
Set 🔒 - Holds finished batches for
append_batches. - Commitment 🔒
- A bound the minter commits ahead of the frontier, broadcast to the writers.
- Finished
Batch 🔒 - Hollow
Batch 🔒AndMetadata - A batch or data + metrics moved from
write_batchestoappend_batches.
Enums§
- Mint 🔒
- What
mint_batch_descriptionsdoes with the frontier and the data it has seen.
Functions§
- append_
batches 🔒 - Fuses written batches together and appends them to persist using one
compare_and_appendcall. Writing only happens for batch descriptions where we know that no future batches will arrive, that is, for those batch descriptions that are not beyond the frontier of both thebatch_descriptionsandbatchesinputs. - description_
lookahead 🔒 - How far past the remap upper the minter commits its ceiling while a snapshot pins the
frontier, in milliseconds, the unit of
mz_repr::Timestamp.Noneturns committing ahead off. - mint_
batch_ 🔒descriptions - 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.
- next_
mint 🔒 - The next thing the minter emits, or
Nonewhen there is nothing to do. - render 🔒
- Continuously writes the
desired_streaminto persist This is done via a multi-stage operator graph: - stage_
update 🔒 - Adds one update to
builder, keeping the batch metrics in step. - write_
batches 🔒 - Writes
desired_collectionto persist, but only for updates that fall into batch a description that we get viabatch_descriptions. This forwards aHollowBatch(with additional metadata) for any batch of updates that was written.
Type Aliases§
- Source
Batch 🔒Builder - The batch builder the source sink writes with.