Skip to main content

Module pipeline

Module pipeline 

Source
Expand description

Thread-per-core connector pipeline. Streaming connector pipeline. Each source connector runs as a tokio task pushing batches via crossfire mpsc to the StreamingCoordinator, which drives SQL execution cycles, routes results to sinks, and manages checkpoint barriers. See the streaming_coordinator submodule for the runtime topology.

Re-exports§

pub use callback::BarrierOutcome;
pub use callback::CycleError;
pub use callback::CycleOutcome;
pub use callback::PipelineCallback;
pub use callback::SkipReason;
pub use callback::SourceRegistration;
pub use config::CheckpointSchedule;
pub use config::PipelineConfig;
pub use streaming_coordinator::ExitReason;
pub use streaming_coordinator::StreamingCoordinator;
pub use streaming_coordinator::StreamingCoordinatorRuntime;

Modules§

callback
Pipeline callback trait and source registration types.
config
Pipeline configuration.
streaming_coordinator
Single-task pipeline coordinator on the dedicated laminar-compute thread.

Structs§

ControlMsg
Opaque live-DDL message used by the streaming coordinator.