What you will be able to do
- Describe how a stream on a stage's directory table finds newly landed document files for a processing pipeline
- Explain what a stream stores (an offset) and what it returns (change records with METADATA$ columns)
- Predict when a stream's offset advances, and use that to make sure each document is processed only once
Key concept
Stream offset — A stream does not copy any data. It keeps a pointer, called an offset, into the change history of its source object. The offset moves forward only when a DML statement consumes the changes, and that is what lets a pipeline process each newly landed document exactly once.
1.The shape of an automated document pipeline
An automated document pipeline has to answer two questions every time it wakes up: *which files are new since the last run?* and *what should be done with them?* Snowflake splits these between two objects. A stream tracks what changed. A task runs SQL against those changes. The documentation sums up the pairing like this: "Combining tasks with table streams is a convenient and powerful way to continuously process new or changed data." On each run, a task can either consume the change data or skip the run if there is nothing new.
For documents, the thing being tracked is a set of files on a stage, not rows in a table. Snowflake's directory-table pipeline example shows the pattern. It detects PDF files added to a stage, extracts data from them, and inserts the results into a table. It "uses a stream to detect changes to a directory table on the stage", and a task calls a UDF to do the extraction. First you create the stage with a directory table enabled. The example also sets server-side encryption, which enables unstructured data access on the stage.
CREATE OR REPLACE STAGE my_pdf_stage
ENCRYPTION = ( TYPE = 'SNOWFLAKE_SSE')
DIRECTORY = ( ENABLE = TRUE);CREATE STREAM my_pdf_stream ON STAGE my_pdf_stage;Checkpoint 1 of 3· Put it in order
Put the steps of the directory-table document pipeline example in order.
- 1.Create an internal stage with a directory table enabled
- 2.Create a task that uses the stream to process new files
- 3.Create a stream on the stage to track changes to its directory table
- 4.Create a UDF that extracts data from the PDF files
The stream depends on the stage's directory table, and the task depends on both the stream and the extraction logic. The example builds the task in step 5, from the stream created earlier.
“In step 5 of this example, we use this stream to construct a task.”Source: docs.snowflake.com
2.What a stream actually holds: an offset and change records
A stream records DML changes (inserts, updates and deletes) to a source object, together with metadata about each change. This is change data capture (CDC). The list of objects a stream can track includes standard tables, views, dynamic tables, external tables and directory tables. The directory table is the one that matters for documents, because its rows describe the files on a stage.
When you create a stream, it takes a logical snapshot of the source by setting an offset at the object's current transactional version. From then on it reports the changes made after that point. The stream holds no data of its own. The documentation suggests you "think of a stream as a bookmark" in the pages of a book. You can drop a bookmark and put new ones in at other places. When you query a stream, it returns rows shaped like the source object, plus three extra metadata columns.
| Column | What it tells you |
|---|---|
| METADATA$ACTION | The DML operation recorded: INSERT or DELETE |
| METADATA$ISUPDATE | TRUE when the INSERT/DELETE pair came from an UPDATE statement |
| METADATA$ROW_ID | A unique, immutable row ID for tracking changes over time |
Updates don't show up as a single row. They appear as a DELETE record and an INSERT record, both with METADATA$ISUPDATE set to TRUE. A stream records the difference between two offsets. So if a row is added and then updated within the same offset window, the stream shows one new row with METADATA$ISUPDATE = FALSE. For a document pipeline, METADATA$ACTION is the column that distinguishes newly arrived files from removed ones.
Checkpoint 2 of 3· Match them up
Match each stream metadata column to what it records.
Tap a term, then the definition that fits it.
Streams return the source columns plus these three. An update shows up as a DELETE/INSERT pair flagged with METADATA$ISUPDATE = TRUE.
“Updates to rows in the source object are represented as a pair of DELETE and INSERT records in the stream”Source: docs.snowflake.com
Sources3
3.When the offset moves, and why that prevents reprocessing
This rule is what stops a pipeline from processing the same document twice. "A stream advances the offset only when it is used in a DML transaction". That includes CTAS and COPY INTO location, in both explicit and autocommit transactions. Until a DML statement consumes the stream, any number of queries can read the same change data without moving the offset. A task that inserts extraction results FROM my_pdf_stream therefore moves the bookmark forward as part of its own commit.
Sometimes you need to skip a backlog without processing it. You can move the offset to the current version without consuming anything in two ways. One is to recreate the stream with CREATE OR REPLACE STREAM. The other is to insert from the stream into a temporary table with a WHERE clause that matches nothing, such as WHERE 0 = 1.
When several statements need to see the same batch of changes, for example one that writes parsed text and another that writes a log row, wrap them in an explicit transaction. Doing so locks the stream. Streams use repeatable read isolation, so every statement in the transaction sees the same set of records. The position moves to the transaction start time if the commit succeeds. If the transaction fails, the position stays put and the same files are offered again on the next run.
Checkpoint 3 of 3· Check yourself
A task body inserts parsed results from a stream and then writes an audit row, also selected from the stream, as two separate autocommit statements. What risk does this create?
Each autocommit DML statement that selects from the stream consumes it. To give both statements the same records, put them inside BEGIN .. COMMIT.
“To ensure multiple statements access the same change records in the stream, surround them with an explicit transaction statement (BEGIN .. COMMIT).”Source: docs.snowflake.com
Sources3
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.Selecting from a stream marks those changes as processed, so previewing new files with SELECT will make the pipeline skip them.Why is that wrong?
Only a DML transaction that uses the stream advances the offset. Plain queries, even inside an explicit transaction, leave it unchanged.
Covered in When the offset moves, and why that prevents reprocessing
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“Combining tasks with table streams is a convenient and powerful way to continuously process new or changed data.”
↩︎ The shape of an automated document pipeline“Each time a task runs, it can either consume the change data or skip the current run if no change data exists.”
↩︎ The shape of an automated document pipeline - 2.
“uses a stream to detect changes to a directory table on the stage”
↩︎ The shape of an automated document pipeline“The example statement sets the ENCRYPTION type to SNOWFLAKE_SSE to enable unstructured data access on the stage.”
↩︎ The shape of an automated document pipeline“In step 5 of this example, we use this stream to construct a task.”
↩︎ Checkpoint - 3.
“It might be useful to think of a stream as a bookmark”
↩︎ What a stream actually holds: an offset and change records“Specifies a unique, immutable row ID for tracking changes over time.”
↩︎ What a stream actually holds: an offset and change records“Multiple queries can independently consume the same change data from a stream without changing the offset.”
↩︎ When the offset moves, and why that prevents reprocessing“Recreate the stream (using the CREATE OR REPLACE STREAM syntax).”
↩︎ When the offset moves, and why that prevents reprocessing“The stream position advances to the transaction start time if the transaction commits; otherwise it stays at the same position.”
↩︎ When the offset moves, and why that prevents reprocessing“a stream itself does not contain any table data. A stream only stores an offset for the source object”
↩︎ Key concept“A stream advances the offset only when it is used in a DML transaction.”
↩︎ Exam trap 1“Updates to rows in the source object are represented as a pair of DELETE and INSERT records in the stream”
↩︎ Checkpoint“Querying a stream alone does not advance its offset, even within an explicit transaction; the stream contents must be consumed in a DML statement.”
↩︎ Prediction“To ensure multiple statements access the same change records in the stream, surround them with an explicit transaction statement (BEGIN .. COMMIT).”
↩︎ Checkpoint