Snowflake Streaming Nodes
Load rows into Snowflake continuously with the Snowflake Streaming connector: rows are queryable within seconds and no warehouse runs.
Configuration Quick Reference
| Field | What you choose | Details |
|---|---|---|
| Parameters | Connection, Function, Function Parameters, Timeout Override | Select the connection profile, function, configure function parameters with expression support, and optionally override timeout. |
| Settings | Description, Timeout (seconds), Retry on Timeout, Retry on Fail, On Error | Node description, maximum execution time, retry behavior on timeout or failure, and error handling strategy. All execution settings default to pipeline-level values. |
Snowflake Stream Rows Node
Snowflake Stream Rows Node
Appends one row (a JSON object) or a batch (an array of objects) to a table. Retries and store-and-forward replays never write a row twice.
Supported Function Types:
| Function Name | Purpose | Common Use Cases |
|---|---|---|
| Stream Rows | Append rows to a table through Snowpipe Streaming | Machine telemetry, production counts, quality results, batch records |
Store-and-forward is available on this node: a write that fails during a network outage is buffered and replayed later, and a replay of rows that already reached Snowflake is recognised and not written again.
Feeding rows into the node. Set the function's Rows to ((data)). On the node, under Function Parameters, give data the rows — one object, an array of objects, or an expression that yields them:
data on the node | Writes |
|---|---|
{{ $node["Buffer"].result }} | The Buffer's flushed messages, one row each (turn the Buffer's Wait For Flush on) |
{{ $trigger.result }} | The trigger's message, when it is one row object or an array of them |
{"asset": "LINE1-PRESS", "count": 1} | A fixed row |
Every other ((parameter)) of the function — in Table, Pipe or Schema — is set the same way. The rules for Rows, and what the call refuses, are in the connector guide: Rows.
Snowflake Channel Status Node
Snowflake Channel Status Node
Reports each streaming channel of a table — committed position, rows inserted and rejected, the last error Snowflake recorded — and one healthy flag to branch on. It only reads.
With the function's Channel empty, it reports the connection's own channel plus every channel this MaestroHub instance has written to the table through since it last started. On a table nothing has been streamed to yet, the node fails with ERR_PIPE_DOES_NOT_EXIST_OR_NOT_AUTHORIZED — Snowflake creates the table's default pipe on the first write.
Supported Function Types:
| Function Name | Purpose | Common Use Cases |
|---|---|---|
| Channel Status | Read the state of a table's channels | Data-quality alarms, ingest monitoring dashboards |
Output
Every Snowflake Streaming node delivers its data under result, and execution facts (success, functionId, durationMs, timestamp) under _metadata:
| Node | Expression | Description |
|---|---|---|
| Stream Rows | $node["Name"].result.table | The fully qualified table written, DATABASE.SCHEMA.TABLE |
$node["Name"].result.pipe | The pipe the rows went through — <TABLE>-STREAMING unless a custom pipe was set | |
$node["Name"].result.channel | The streaming channel used | |
$node["Name"].result.rowsSent | How many rows the call carried | |
$node["Name"].result.requests | How many requests they took — a batch over 4 MB is split; 0 when deduplicated | |
$node["Name"].result.bytesSent | Compressed bytes sent to Snowflake | |
$node["Name"].result.deduplicated | true when the rows were already in Snowflake from an earlier attempt of this same call, so nothing was written again | |
$node["Name"].result.offsetToken | The offset token of the last request. Absent when the call was deduplicated before sending | |
$node["Name"].result.committed | Whether Snowflake committed the rows before the call returned. Present only with Wait for Commit on | |
$node["Name"].result.rowsInserted | Rows Snowflake committed into the table. Wait for Commit only, and absent when the call was deduplicated | |
$node["Name"].result.rowsRejected | Rows Snowflake skipped because they did not fit the table — a non-zero count fails the call. Wait for Commit only, and absent when the call was deduplicated | |
$node["Name"].result.lastError | Snowflake's message for the last rejected row. Only when rows were rejected. Snowflake hides the offending value and names no column here; the table's ERROR_TABLE() has both | |
$node["Name"].result.tableCreated | true when this call created the table — Create Table If Not Exists on and the table was missing | |
$node["Name"].result.warning | Why an option could not take effect: today, Allow Schema Evolution is on but the existing table has schema evolution switched off. Reported once per table | |
| Channel Status | $node["Name"].result.healthy | false when any channel must be reopened or had a row rejected within the Error Window |
$node["Name"].result.table | The fully qualified table | |
$node["Name"].result.pipe | The pipe whose channels are reported | |
$node["Name"].result.channels | One entry per channel, in name order — never null, so it can be looped over without a check. E.g. $node["Name"].result.channels[0].rowsErrorCount |
Each entry of channels:
| Field | When present | Description |
|---|---|---|
channel, exists | Always | The channel name, and whether Snowflake knows it. false for a channel never written, or dropped after 30 idle days |
statusCode | exists is true | Snowflake's channel status: SUCCESS (or ACTIVE) is usable, anything else must be reopened |
lastCommittedOffsetToken | exists is true | The offset token of the last committed request |
rowsInserted, rowsParsed, rowsErrorCount | exists is true | Snowflake's counters for the channel: rows committed, rows read, rows rejected |
avgProcessingLatencyMs | exists is true | Snowflake's average processing latency for the channel |
lastErrorMessage, lastErrorOffsetToken | A row was ever rejected | Snowflake's message for the last rejected row (value and column hidden), and the offset token of its request |
lastErrorAt | A row was ever rejected | When that row was rejected |
healthy, reason | exists is true (reason only when unhealthy) | Whether this channel is healthy, and why not — e.g. Snowflake rejected rows within the last 1h0m0s |
pendingRows, pendingSince | This MaestroHub instance has written through the channel since it started (pendingSince only while rows are pending) | Rows acknowledged but not yet committed, and since when |
reopens, deduplicatedBatches | Same as pendingRows | How often this instance reopened the channel, and how many replayed calls it recognised and did not write again |
A channel with exists: false has no healthy field and does not make the overall healthy false.
Both nodes also deliver the connector's own facts under _metadata: $node["Name"]._metadata.method (the operation that ran), $node["Name"]._metadata.connectionId and $node["Name"]._metadata.protocol (snowflake_streaming).
Snowflake answers a write as soon as it has safely received the rows, and only afterwards skips rows that do not fit the table. With Wait for Commit off the node cannot know about those rows yet: the connector logs them when it sees them, and Channel Status reports them. Turn Wait for Commit on where a rejected row must not go unnoticed: the call then fails, and with the default Delivery mode it is set aside in Failed messages with Snowflake's message while later writes keep flowing.