Fluxuri și captura datelor modificate

Automatizarea pipeline-urilor de date în Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Procesarea doar a datelor modificate

CDC: Change Data Capture Două comparații privind modul de procesare a datelor

Automatizarea pipeline-urilor de date în Snowflake

Fluxuri

Ce face un flux

  • Urmărește fiecare INSERT, UPDATE și DELETE pe o tabelă sursă

Screenshot 2026-05-11 at 10.50.31 am.png

  • Menține un jurnal continuu al modificărilor — fără duplicarea datelor
  • După consum, offset-ul avansează; următoarea citire începe de la zero
1 * Snowflake Learning Resource
Automatizarea pipeline-urilor de date în Snowflake

Tipuri de fluxuri

 

Tip flux Capturează Recomandat pentru
Standard Toate tipurile de tabele și view-uri & toate modificările DML - urmărește inserări, actualizări, ștergeri Tabele unde orice rând poate fi modificat (ex. expedieri)
Append-only Toate tipurile de tabele și view-uri, cu excepția tabelelor externe - urmărește doar inserările Tabele cu inserare unică (ex. evenimente de livrare) - mai eficient
Insert-only Tabele Apache Iceberg și tabele externe gestionate extern - urmărește doar inserările Tabele externe

Tabelele de directoare expun metadatele fișierelor dintr-un stage (nume, dimensiune, timestamp ultima modificare)

Automatizarea pipeline-urilor de date în Snowflake

Crearea unui flux

Flux standard pe tabela de expedieri

CREATE STREAM shipments_stream
  ON TABLE logistics.shipments;

Flux append-only pe tabela de evenimente de livrare

CREATE STREAM delivery_events_stream
  ON TABLE logistics.delivery_events
  APPEND_ONLY = TRUE;
Automatizarea pipeline-urilor de date în Snowflake

Coloanele de metadate ale fluxului

SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
  • METADATA$ACTION: INSERT sau DELETE
  • METADATA$ISUPDATE: TRUE când face parte dintr-o pereche de actualizare
  • METADATA$ROW_ID: Identificator fizic unic al rândului
  • Actualizările apar ca o pereche DELETE + INSERT, ambele marcate cu METADATA$ISUPDATE = TRUE

Coloanele de metadate ale fluxului: METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID cu date exemplu

Automatizarea pipeline-urilor de date în Snowflake

Offset-ul fluxului

Diagramă cronologică - Flux creat (offset începe aici) → Modificări în tabela sursă (fluxul acumulează înregistrări) → Flux consumat într-o tranzacție (offset avansează la momentul curent)

Automatizarea pipeline-urilor de date în Snowflake

Prezentare generală: fluxuri într-un pipeline

  • Fluxurile se asociază cu taskuri — obiecte Snowflake care execută SQL după un program
  • Taskul citește doar rândurile modificate; cu 10M rânduri: 2 secunde față de 2 minute

Screenshot 2026-05-11 at 10.48.51 am.png

1 * Snowflake Learning Resource
Automatizarea pipeline-urilor de date în Snowflake

Interogare: fluxuri într-un pipeline

  • Fluxurile se asociază cu taskuri - obiecte Snowflake care execută SQL după un program
  • Taskul citește doar rândurile modificate; cu 10M rânduri: 2 secunde față de 2 minute
CREATE TASK logistics.sync_shipments
  WAREHOUSE = compute_wh
  SCHEDULE = '5 MINUTE'
  WHEN SYSTEM$STREAM_HAS_DATA('logistics.staging_shipments_stream')
AS
  INSERT INTO logistics.shipments
  SELECT shipment_id, region, carrier, delivery_days
  FROM logistics.staging_shipments_stream
  WHERE METADATA$ACTION = 'INSERT';
Automatizarea pipeline-urilor de date în Snowflake

Să exersăm!

Automatizarea pipeline-urilor de date în Snowflake

Preparing Video For Download...