A Stream durably appends ordered JSON records. A Pipeline reads that Stream,
applies a stateless SQL transform, maintains its own cursor, and retries a
deterministic batch until every Sink confirms it.
Send records
Declare a fixed Stream binding:
Then call send from any Worker handler:
send requires an array of JSON-serializable values, rejects requests over 5
MiB, and resolves only after the Stream returns a durable 2xx acknowledgement.
A single Stream can feed multiple Pipelines with independent cursors.
A Pipeline configuration selects one Stream, maps SQL destination names to
Sinks, and defines explicit row, byte, and time batch limits. The supported SQL
surface is intentionally bounded:
Supported operations include projection, aliases, nested field access,
filters, arithmetic, comparisons, WITH, one UNNEST per SELECT, and scalar
functions such as UPPER, LOWER, TRIM, ROUND, COALESCE, NULLIF, and
CONCAT.
Pipeline SQL is not general analytical SQL. Joins, aggregates, windows,
GROUP BY, ORDER BY, DDL, UPDATE, and DELETE are rejected. Use Pipeline
for record-by-record transformation and a Query read model for maintained
aggregates.
Fan out
Multiple statements can read the same source and target different Sinks:
Each delivered batch has a deterministic identity derived from the Pipeline,
SQL digest, source sequence range, and Sink. The Pipeline advances its cursor
only after all destinations confirm that identity. A retry cannot create a
second logical delivery in an idempotent built-in Sink.
Structured Streams
Streams may enforce an immutable schema with string, integer and float
types, bool, timestamp, json, binary, list, and nested struct.
Invalid positions stay visible in validation metrics but are skipped during
Pipeline processing without renumbering later sequence positions.
Backfill a Stream from a scheduled Worker
A Worker can use a past start_date, bounded concurrency, and its logical
scheduledTime to ingest historical partitions into a Stream while current
cron deadlines continue to run. See
Scheduled Workers and backfills for a daily
market-data example.