Skip to main content

S2 input connector

Experimental feature

S2 support is an experimental feature of Feldera.

Feldera can consume ordered records from an S2 stream with the s2_input connector. The connector checkpoints S2 sequence numbers and replays checkpointed ranges when a pipeline resumes.

Configuration​

PropertyTypeRequiredDescription
basinstringYesS2 basin name.
streamstringYesS2 stream name.
auth_tokenstringYesS2 authentication token. Use a secret reference in production.
endpointstringNoCustom S2 endpoint, primarily for local development. When omitted, the S2 cloud endpoint is used.
start_fromvariantNoInitial read position when no checkpoint exists. Defaults to "Beginning".

start_from accepts the following values:

  • "Beginning": start at sequence number 0.
  • "Tail": consume records appended after the connector starts. The tail is resolved once, right before the first record is read, and that sequence number is then checkpointed, so it does not keep drifting as the stream grows. Before the first record, the anchor is time-dependent; if you need a fully deterministic start, use SeqNum or Timestamp.
  • {"SeqNum": 42}: start at sequence number 42.
  • {"Timestamp": 1773900000000}: start at a Unix timestamp in milliseconds.
  • {"TailOffset": 100}: start 100 records before the current tail.

Once a checkpoint exists, its sequence number takes precedence over start_from. If the checkpoint was taken after the start position was resolved (a record was consumed or a tail was anchored), the connector always resumes from that absolute sequence number — including sequence 0 — rather than recomputing start_from. The connector preserves ordered delivery across reconnects and validates sequence contiguity after the start position has been resolved. Gaps, duplicates, or S2 trimmed records required to replay a checkpoint fail the connector rather than silently skipping or duplicating data.

Retry behavior​

The S2 SDK performs finite retries for transient service and network failures. If those retries are exhausted, Feldera classifies the resulting S2 error. Errors that are safe to retry, such as rate limiting, service unavailability, request or heartbeat timeouts, and common connection resets, are reported as non-fatal and retried indefinitely with bounded backoff from the last resolved S2 sequence number. The pipeline can still flush records that were already buffered while the connector is waiting to reconnect.

Fatal errors, such as malformed access tokens, validation failures, append condition failures, reads from unwritten positions, replay gaps or short reads, and unknown S2 error codes, stop the connector and require operator action. Pause and disconnect requests cancel any in-progress reconnect backoff promptly.

Example​

CREATE TABLE orders (
id BIGINT,
amount DECIMAL(12, 2)
) WITH (
'connectors' = '[{
"transport": {
"name": "s2_input",
"config": {
"basin": "commerce",
"stream": "orders",
"auth_token": "${secret:kubernetes:s2/auth-token}",
"start_from": "Beginning"
}
},
"format": {
"name": "json",
"config": {
"update_format": "raw"
}
}
}]'
);

Each S2 record body is passed to the configured Feldera input format. For JSON streams containing insert/delete envelopes, set update_format to "insert_delete" instead of "raw".