Snowpipe și Snowpipe Streaming

Automatizarea pipeline-urilor de date în Snowflake

Emily Melhuish

Technical Curriculum Developer, Snowflake

Snowpipe

Caz de utilizare:

  • Evenimentele de livrare sosesc continuu

Soluție: Snowpipe

  • Snowpipe încarcă date din fișiere imediat ce sunt disponibile într-un stage.

Captură de ecran 2026-05-11 la 12.24.14 pm.png

1 * Material de învățare Snowflake
Automatizarea pipeline-urilor de date în Snowflake

Problema încărcării batch

  • Fișierele sosesc în S3 la câteva minute interval pe parcursul zilei
  • COPY INTO planificat rulează la miezul nopții — decalaj de 24 de ore
  • Livrările întârziate nu vor apărea decât a doua zi.
  • Snowpipe elimină acest decalaj
-- Nightly batch: runs at 00:00, data arrives all day
COPY INTO logistics.delivery_events
FROM @harbr_s3_stage/events/
FILE_FORMAT = (FORMAT_NAME = 'harbr_json_format');
-- A 9am exception won't appear until tomorrow
Automatizarea pipeline-urilor de date în Snowflake

Ce este Snowpipe?

  • Încapsulează o instrucțiune COPY INTO — aceeași sintaxă, aceleași formate de fișiere
  • Se declanșează automat la sosirea de noi fișiere într-un stage
  • Încarcă în micro-batch-uri, de obicei în câteva minute
  • Serverless — nu necesită provizionarea unui warehouse
CREATE PIPE harbr_events_pipe AS
  COPY INTO logistics.delivery_events
  FROM @harbr_s3_stage/events/
  FILE_FORMAT = (FORMAT_NAME = 'harbr_json_format');
Automatizarea pipeline-urilor de date în Snowflake

Cum funcționează Snowpipe

Fluxul de lucru Snowpipe

  • AUTO_INGEST — bazat pe evenimente; stocarea în cloud publică o notificare
  • Amazon S3 | Azure Event Grid | GCP Pub/Sub
  • Declanșator REST API — apelați direct endpoint-urile insertFiles sau insertReport din codul de orchestrare
Automatizarea pipeline-urilor de date în Snowflake

Facturare Snowpipe

Captură de ecran 2026-05-11 la 12.24.14 pm.png

  • Facturat în funcție de un număr fix de credite per GB consumat
  • Fișiere text: taxare în funcție de dimensiunea necomprimată
  • Fișiere binare: taxare în funcție de dimensiunea observată
Automatizarea pipeline-urilor de date în Snowflake

Snowpipe Streaming

Snowpipe Snowpipe Streaming
Declanșator Fișier ajuns în stage Rând scris de aplicație
Latență Minute Secunde
Caz de utilizare Fluxuri de fișiere GPS, IoT, date în timp real

 

Elimină complet granița de fișier

  • Rândurile sunt scrise direct din aplicație prin Streaming Ingest SDK
  • Fără fișiere, fără stage-uri — latență în secunde
# Snowpipe Streaming: application writes rows directly
channel = client.openChannel('GPS_CHANNEL', 'LOGISTICS', 'GPS_EVENTS')
channel.insertRows(rows=[
    {'vehicle_id': 'V001', 'lat': 51.5, 'lng': -0.12, 'ts': now()}
])
Automatizarea pipeline-urilor de date în Snowflake

Alegerea metodei de ingestie potrivite

Metode de ingestie

Metodă Când se utilizează
COPY INTO Încărcări batch planificate - fișiere nocturne, exporturi săptămânale, latență de ore acceptabilă
Snowpipe Sosiri continue de fișiere, încărcare în câteva minute de la sosire
Snowpipe Streaming Date generate de aplicații - GPS, IoT, piețe financiare - date disponibile în secunde
1 * Resursă de învățare Snowflake
Automatizarea pipeline-urilor de date în Snowflake

Să exersăm!

Automatizarea pipeline-urilor de date în Snowflake

Preparing Video For Download...