Automatizarea pipeline-urilor de date în Snowflake
Emily Melhuish
Technical Curriculum Developer, Snowflake
CDC: Change Data Capture

Ce face un flux

| 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)
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;
SELECT product, quantity, METADATA$ACTION, METADATA$ISUPDATE, METADATA$ROW_ID
FROM shipments_stream;
METADATA$ACTION: INSERT sau DELETEMETADATA$ISUPDATE: TRUE când face parte dintr-o pereche de actualizareMETADATA$ROW_ID: Identificator fizic unic al rânduluiMETADATA$ISUPDATE = TRUE


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