Skip to main content
A streaming connector keeps a session open across many commits. Its most important job is to explain when a portion of the input is complete and when the corresponding output is durable. The runner uses those two promises to decide where a later attempt may safely resume. Start with the streaming lifecycle for the end-to-end picture. This page focuses on the connector’s part in that lifecycle.

Declare what the pair can support

Sources describe whether they emit rows, changes, or messages, which operations they produce, and what ordering and replay guarantees they offer. Continuous planning currently requires replayable, at-least-once input. A best-effort feed cannot promise to recover the messages lost during a worker interruption. Sinks declare continuous write policies separately from bounded policies. Supporting bounded upserts does not automatically provide streaming upserts. The planner checks operations, ordering, durability, and required primary keys before writing. Replace is unavailable for continuous execution because it needs a finite snapshot and a completion boundary. The current planner accepts an order-dependent sink policy only when the source promises strict global order. Partition or transaction ordering alone does not satisfy that requirement. Declare the guarantee the connector can actually maintain through extraction and writing.

Keep the source session alive

Implement StreamSource to open a reusable StreamSession. Opening receives the selected resources, certified resume positions, resource membership, and the admitted worker identity. Use those positions to resume the existing feed; do not silently replace an unavailable position with the latest available one. The runner alternates reads and acknowledgements on the session. A read writes typed rows through StreamRecordSink and returns the source positions covered by that read. Those positions are candidates until the runner has finished the pipeline and committed the epoch. Backpressure can block a row flush for much longer than a network read. Keep protocol heartbeats and any required delivery renewal alive while waiting. Honor cancellation, and join an outstanding read or acknowledgement before closing the session. Closing concurrently with a read is not part of the contract.

Show that the source boundary is complete

Consider a source transaction that updates an order and inserts its line items. Reaching a row limit after the order update is not a safe place to save progress if resuming there would skip the line items. Finish the transaction and report its end before returning coverage. Sources express these boundaries with ordered controls alongside the rows. Transaction begin and end markers describe transaction boundaries; a progress boundary can describe a safe position outside a transaction. Sending a control flushes the affected builders and waits for preceding sink writes and integrity checks. It does not commit the epoch or acknowledge the provider. Before an epoch can be sealed, every covered domain needs an accepted safe boundary after its last row, and no transaction may remain open. Returning a position from Read without the corresponding pipeline evidence is not enough. The source is also responsible for proving it did not skip relevant input on the way to that position. The pipeline cannot discover a missing broker message from the rows it received. Treat record, byte, idle, and age limits as targets for finding that boundary. Cancellation is an interruption, never evidence that a transaction completed. The current runner accepts domain-position coverage; the inbox-claim types in the contracts do not provide an implemented inbox execution path.

Give bookmarks a stable meaning

An ordering domain identifies the part of the source whose positions can be compared, together with its incarnation. A partition offset is useful only if it still refers to the same partition history. Detecting a recreated source prevents an old offset from being mistaken for progress in a new feed. Implement a position codec that validates, compares, and canonicalizes those bookmarks. Codec lookup must work without connector configuration or network access. Register the codecs for runtime persistence as well as exposing them through the source. The datastore needs the same interpretation when it checks for regression or compares a retried certificate. Where the source has a stable event identity, preserve it across retries. It combines the domain, position, and an ordinal for events sharing a position. Do not build that identity from the run or worker attempt: doing so would make a replay look like a new event to a deduplicating sink.

Acknowledge only certified work

Acknowledge receives coverage after destination commit and durable epoch certification. A source must not acknowledge speculative progress during a read. Check current ownership through the supplied authority callback before provider acknowledgements, including redeliveries that you suppress because they are already certified. This check matters when an old worker wakes up after losing its lease. Knowing the correct bookmark does not give it permission to change the provider’s consumer state.

Implement the sink’s epoch lifecycle

A streaming sink opens once, then repeats three steps: begin an epoch, apply its batches, and commit it. StreamingSink supplies that lifecycle through BeginEpoch, CommitEpoch, and AbortEpoch; CloseSession ends the reusable session. The runner does not adapt bounded Commit and Abort methods into streaming behavior. Bind each batch to the active epoch and reject mismatched identities. Keep the same Arrow and encoded integrity obligations as a bounded sink. When commit returns, the writes must have reached the durability promised by the selected policy, and the per-resource receipts must describe the completed work. For PostgreSQL append, this means starting a transaction, copying the batches into it, and returning receipts after transaction commit. For NATS append, messages are already acknowledged by the destination during batch application; the epoch commit summarizes that accepted work. These implementations have different visibility and rollback behavior even though they use the same runner lifecycle. Aborting discards only tentative effects. It cannot undo messages already published or earlier committed epochs. Document that boundary so users can reason about retries and duplicates.

Make shutdown an honest promise

A successful CloseSession means no later destination effects remain possible from that session. Releasing a local client or reaching a timeout does not establish that promise if remote requests could still finish. Return an error when shutdown is uncertain. The runtime records that uncertainty and blocks automatic takeover. Similarly, matching an epoch ID in memory does not prove destination-side owner fencing: that stronger capability requires the destination to check ownership atomically with its effects.

Test interruptions at the boundaries

Successful delivery is only part of streaming support. Exercise a slow sink, cancellation during a read, a partial source transaction, and loss of ownership before acknowledgement. Verify that a source cannot certify a position whose rows have not finished writing. For a sink, test failed writes, uncertain commit responses, repeated lifecycle calls, and shutdown with outstanding work. An interruption after destination commit but before certificate persistence is especially useful: recovery may replay that epoch, and the result must match the connector’s documented duplicate-handling behavior. The public contracts live in stream.go and epoch.go. The coordinator is in runner/continuous.go, and pipeline completion is enforced in pipeline/epoch.go. The Kafka and NATS sources and PostgreSQL and NATS sinks provide concrete implementations to read alongside those contracts.