Skip to main content

Module command_channel

Module command_channel 

Source
Expand description

A channel for sequencing commands between all workers of a Timely cluster.

Compute uses a dataflow to distribute commands between workers. This is to ensure workers retain a consistent dataflow state across reconnects. If each worker would receive its commands directly from the controller, there wouldn’t be any guarantee that after a reconnect all workers have seen the same sequence of commands. This is particularly problematic for CreateDataflow commands, since Timely requires that all workers render the same dataflows in the same order. So the controller instead sends commands only to worker 0, which then broadcasts them to other workers through the Timely fabric, taking care of the correct sequencing.

Commands in the command channel are tagged with a nonce identifying the incarnation of the compute protocol the command belongs to, allowing workers to recognize client reconnects that require a reconciliation.

The channel optionally also carries storage-internal commands, for clusters that host storage objects alongside compute objects. Both command kinds are sequenced through a single lane, so all workers observe one consistent interleaving and therefore construct all dataflows, compute and storage alike, in the same order. Unlike compute commands, storage-internal commands may be injected from any worker (e.g. by health operators triggering a suspend-and-restart), so the channel uses a two-hop structure copied from storage’s command sequencer: producers tag commands with a per-producer index, worker 0 fixes one definitive order and assigns a global index, and receivers restore that order.

Structs§

Receiver
A receiver reading commands from the command channel.
Sender
A sender pushing compute commands onto the command channel.
StorageLaneInput
Per-worker storage-side inputs to the command channel.

Enums§

UnifiedCommand
A command in the unified command lane.

Functions§

render
Render the command channel dataflow.
split_command 🔒
Split the given command into one part per target worker.