Skip to main content
Version: 3.0 (next)

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​

FieldWhat you chooseDetails
ParametersConnection, Function, Function Parameters, Timeout OverrideSelect the connection profile, function, configure function parameters with expression support, and optionally override timeout.
SettingsDescription, Timeout (seconds), Retry on Timeout, Retry on Fail, On ErrorNode 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

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 NamePurposeCommon Use Cases
Stream RowsAppend rows to a table through Snowpipe StreamingMachine 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 nodeWrites
{{ $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

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 NamePurposeCommon Use Cases
Channel StatusRead the state of a table's channelsData-quality alarms, ingest monitoring dashboards

Output​

Every Snowflake Streaming node delivers its data under result, and execution facts (success, functionId, durationMs, timestamp) under _metadata:

NodeExpressionDescription
Stream Rows$node["Name"].result.tableThe fully qualified table written, DATABASE.SCHEMA.TABLE
$node["Name"].result.pipeThe pipe the rows went through — <TABLE>-STREAMING unless a custom pipe was set
$node["Name"].result.channelThe streaming channel used
$node["Name"].result.rowsSentHow many rows the call carried
$node["Name"].result.requestsHow many requests they took — a batch over 4 MB is split; 0 when deduplicated
$node["Name"].result.bytesSentCompressed bytes sent to Snowflake
$node["Name"].result.deduplicatedtrue when the rows were already in Snowflake from an earlier attempt of this same call, so nothing was written again
$node["Name"].result.offsetTokenThe offset token of the last request. Absent when the call was deduplicated before sending
$node["Name"].result.committedWhether Snowflake committed the rows before the call returned. Present only with Wait for Commit on
$node["Name"].result.rowsInsertedRows Snowflake committed into the table. Wait for Commit only, and absent when the call was deduplicated
$node["Name"].result.rowsRejectedRows 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.lastErrorSnowflake'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.tableCreatedtrue when this call created the table — Create Table If Not Exists on and the table was missing
$node["Name"].result.warningWhy 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.healthyfalse when any channel must be reopened or had a row rejected within the Error Window
$node["Name"].result.tableThe fully qualified table
$node["Name"].result.pipeThe pipe whose channels are reported
$node["Name"].result.channelsOne 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:

FieldWhen presentDescription
channel, existsAlwaysThe channel name, and whether Snowflake knows it. false for a channel never written, or dropped after 30 idle days
statusCodeexists is trueSnowflake's channel status: SUCCESS (or ACTIVE) is usable, anything else must be reopened
lastCommittedOffsetTokenexists is trueThe offset token of the last committed request
rowsInserted, rowsParsed, rowsErrorCountexists is trueSnowflake's counters for the channel: rows committed, rows read, rows rejected
avgProcessingLatencyMsexists is trueSnowflake's average processing latency for the channel
lastErrorMessage, lastErrorOffsetTokenA row was ever rejectedSnowflake's message for the last rejected row (value and column hidden), and the offset token of its request
lastErrorAtA row was ever rejectedWhen that row was rejected
healthy, reasonexists is true (reason only when unhealthy)Whether this channel is healthy, and why not — e.g. Snowflake rejected rows within the last 1h0m0s
pendingRows, pendingSinceThis 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, deduplicatedBatchesSame as pendingRowsHow 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).

Without Wait for Commit, a successful node does not mean every row landed

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.