Streaming
Run a trained session's data preparation continuously on Kafka, Kinesis or MQTT records, with throughput, lag and dead-letter monitoring. Enterprise plan.
What streaming does (Enterprise)
A stream is inference that never stops. It reads records from a Kafka topic, a Kinesis stream or an MQTT topic, prepares them in micro-batches with the fitted transforms of a completed session, and writes the prepared rows to a sink. Each prepared row is what batch inference would have written to features.csv for the same input.
Streaming is part of the Enterprise plan and must also be switched on for your deployment. Other accounts see an upgrade notice on the Streaming tab, and the API answers 403. While streaming is switched off on the server, every streams route answers 404.
Before you start
- A completed session whose stored transforms are still available. Streams replay it; they never train.
- A saved connection for the source and one for the sink, created on the Connections tab. Supported for both: Kafka, Kinesis and MQTT.
- A stream signs in as the saved connection says: SASL for Kafka, a custom CA certificate and a client certificate for Kafka and MQTT (pasted as PEM text in
ca_cert_pem,client_cert_pemandclient_key_pem), and temporary keys, an assumed role or a custom endpoint for Kinesis. Saving a change to a connection restarts the streams that use it.
Creating a stream
On the Streaming tab, choose New stream and fill in:
| Setting | What it does |
|---|---|
| Name | A label for the stream. |
| Trained session | The completed session to replay. |
| Source | Connection type, saved connection, and topic or stream name. For MQTT, topic wildcards such as sensors/+/readings are accepted. |
| Start position | Latest reads only new records. Earliest replays the records the source still retains. |
| Input format | JSON: one record per message. NDJSON: many records per message, one per line. |
| Sink | Connection type, saved connection, and topic or stream name. It cannot be the same topic as the source. |
| Rows per message | How many prepared rows to pack into each outgoing message. |
| Output format | See below. |
| Partition key field | Kafka and Kinesis: the output column whose value keys each message. Leave it blank to spread messages evenly. |
| QoS | MQTT only: 1 (at least once) or 0 (at most once), for subscribing and for publishing. |
| Dead letters | Optional. Records that fail to parse or transform are written to a separate dead-letter target. Without one, they are counted and the most recent are kept for inspection. |
| Apply the session's anomaly handling | On by default. Replays the anomaly decisions the session made at training time. Turn it off to pass rows straight to the fitted transforms. |
Use Test source connection and Test sink connection to check the connections before saving. A new stream starts as soon as it is created.
Advanced settings
| Setting | Range | Default |
|---|---|---|
| Max batch rows | 1 to 100,000 | 5,000 |
| Max batch time (ms) | 50 to 60,000 | 1,000 |
| Parallelism (workers) | 1 to 256 | 1 |
Drain timeout (drain_timeout_s, API) | 5 to 600 seconds | 30 |
Rows per message (sink_config.rows_per_message) | 1 to 10,000 | 1 |
Max message size (sink_config.max_message_bytes, API) | 1,024 to 10,485,760 bytes | Kafka 900,000; Kinesis 1,000,000; MQTT 128,000 |
The drain timeout is how long a stopping worker may take to finish its current batch. Messages larger than the maximum size are split before sending. A stream stops with status error after 10 consecutive failed batches, or after 3 when the source or sink rejects the connection's credentials; the reason is shown on the stream in plain words. A batch is flushed when it reaches either limit, whichever comes first. More parallelism only helps up to the number of source partitions or shards.
Output formats
| Format | Each row is | Use it when |
|---|---|---|
json | An object with column names as keys | Consumers want self-describing rows and the session is narrow. |
values | An array of values in a fixed column order | Consumers already know the column order. |
sparse | {"i": [positions], "v": [values]} of its non-zero cells; absent positions are 0 | The session has many one-hot or text columns. Sparse is often 100x smaller and several times faster. |
The fixed column order for values and sparse is the stream's schema: GET /api/v1/streams/{id}/schema. Positions never shift between batches; a column missing from a batch is null at its position.
Delivery guarantees
Delivery is at least once. Source offsets are committed only after the sink has flushed, so after a restart some records may be written twice. Downstream consumers must tolerate duplicates.
A record that cannot be processed is isolated from the rest of its batch and sent to the dead-letter target, so one bad record does not hold up the others.
Monitoring
The Streaming tab shows the number of running streams, total rows per second, errors and dead letters. Open a stream for its detail view:
- Rows / s, with the rate expected from the session's measured speed and the stream's parallelism, and a five-minute sparkline.
- Lag: how many records the stream is behind the source.
- Workers live: each worker, the partitions or shards it owns, and its own rate and lag.
- Errors and the most recent dead letters, with the reason and a sample of each payload.
If a stream is slower than expected, the detail view says why: typically a very wide session. The sparse output format and more parallelism are the usual fixes.
Billing
A stream is billed hourly for the rows it scored in that hour, at the inference rate (about a fifth of preparing the same data). There is no price confirmation and no per-run limit for streams. See Account and billing.
Changing a running stream
Streams can be stopped, started, restarted, edited and deleted from the tab. Saving an edit to a running stream drains its in-flight batches and restarts it with the new settings. Deleting a stream stops it; records already written to the sink are not affected.
API
| Method | Path | Purpose |
|---|---|---|
| GET | /api/v1/streams | List your streams |
| POST | /api/v1/streams | Create a stream (starts it unless "start": false) |
| POST | /api/v1/streams/test-connection | Test a source or sink configuration |
| GET / PUT / DELETE | /api/v1/streams/{id} | Read, update or delete a stream |
| POST | /api/v1/streams/{id}/start, /stop, /restart | Change its running state |
| GET | /api/v1/streams/{id}/status | Throughput, lag, workers and recent dead letters |
| GET | /api/v1/streams/{id}/schema | Column order of the values and sparse formats |
A stream body names name, base_session_id, source_type, source_credential_id, source_config, sink_type, sink_credential_id, sink_config, and optionally dead_letter_config, max_batch_rows, max_batch_ms, parallelism, drain_timeout_s, enable_anomaly and start. test-connection takes { connector_type, credential_id, role: "source" | "sink", config }.
LM Readiness sessions
A stream can replay an LM Readiness session: each batch goes through the session's classic pipeline and saved LM transformations, and the stream emits the columns of the session's training partition, without the row id and targets. Vector columns are sent as lists unless the session flattens them; Sparse output needs flattened vectors.
Event sequences and transaction history use each entity's recent events, kept by each worker: the last sequence length events per entity, and for transaction history up to 256 events per entity and counterparty within 30 days. A batch joins that history once delivered. History starts empty when a worker starts; key the source by the entity column so each entity's events reach the same worker.