S2 input connector
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
| Property | Type | Required | Description |
|---|---|---|---|
basin | string | Yes | S2 basin name. |
stream | string | Yes | S2 stream name. |
auth_token | string | Yes | S2 authentication token. Use a secret reference in production. |
endpoint | string | No | Custom S2 endpoint, primarily for local development. When omitted, the S2 cloud endpoint is used. |
start_from | variant | No | Initial 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, useSeqNumorTimestamp.{"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".