🎉 75% of content is free forever — Unlock Premium from $10/mo →
CW
NEWSLIVESearch All Content
đŸ’ŧ Servicesâ„šī¸ Aboutâœ‰ī¸ ContactView Pricing Plansfrom $10

Snowflake Provider Integration with Airflow

đŸŸĸ Free Lesson

Advertisement

Snowflake Provider Integration with Airflow

Snowflake Integration ArchitectureAirflowSchedulerSnowflakeHookConnection + SQLSnowflakeData warehouseExternal StageS3 / GCS / AzureSnowpipeAuto-ingestSnowflakeOperatorExecute SQLGCSToSnowflakeLoad from GCSCopyFromStageCOPY INTOTransferOperatorCross-region copyUse key pair auth for production; Snowpipe for continuous data loading

Architecture Diagram

Formal Definitions

Detailed Explanation

Snowflake Integration Overview

The Snowflake provider enables Airflow to execute SQL, load data, and manage Snowflake resources. It supports COPY INTO operations, Snowpipe management, and data transfer between Snowflake and cloud storage.

Key Insight: Snowpipe auto-ingestion reduces Airflow's role to trigger and monitor. Let Snowflake handle continuous loading for better efficiency.

Operator Selection Guide

OperatorUse CaseKey Parameters
SnowflakeOperatorExecute SQLsql, snowflake_conn_id
SnowflakeToGCSOperatorExport to GCSgcs_bucket, sql
GCSToSnowflakeOperatorLoad from GCSbucket, prefix, table
SnowflakeCopyFromExternalStageCOPY INTOtable, stage, file_format
SnowflakeSensorWait for conditionsql, poke_interval

COPY INTO Options

OptionDescriptionCommon Values
FILE_FORMATFile type specificationTYPE = 'PARQUET', TYPE = 'CSV'
PATTERNRegex for file matchingPATTERN = '.*\\.parquet'
ON_ERRORError handling strategySKIP_FILE, CONTINUE, ABORT_STATEMENT
FORCEReload existing filesFORCE = TRUE
PURGEDelete files after loadPURGE = TRUE

Connection Setup

Basic SQL Execution

GCS to Snowflake Loading

Pre/Post Operator Pattern

Multi-Warehouse Pattern

Key Concepts Table

OperatorPurposeKey Parameters
SnowflakeOperatorExecute SQLsql, snowflake_conn_id
SnowflakeToGCSOperatorExport to GCSsnowflake_conn_id, gcs_bucket, sql
GCSToSnowflakeOperatorLoad from GCSbucket, prefix, table, stage
SnowflakeCopyFromExternalStageCOPY INTOtable, stage, file_format
SnowflakePrePostOperatorPre/post hookssql, pre_query, post_query
SnowflakeSensorWait for conditionsql, poke_interval

COPY INTO Options

OptionDescriptionExample
FILE_FORMATFile type specificationTYPE = 'PARQUET'
PATTERNRegex for file matchingPATTERN = '.*\\.parquet'
ON_ERRORError handling strategySKIP_FILE, CONTINUE, ABORT_STATEMENT
FORCEReload existing filesFORCE = TRUE
SIZE_LIMITMax bytes per loadSIZE_LIMIT = 1073741824
PURGEDelete files after loadPURGE = TRUE
RETURN_FAILED_ONLYReturn only failed filesRETURN_FAILED_ONLY = TRUE

Best Practices

Data Loading

  1. Use named stages for reusable external connections — avoid embedded credentials.
  2. Set ON_ERROR appropriately: SKIP_FILE for bad files, CONTINUE for partial loads.
  3. Partition large COPY INTO operations with SIZE_LIMIT to avoid single large transactions.
  4. Use PURGE = TRUE to clean up files after successful load.

Monitoring and Cost

  1. Set QUERY_TAG for monitoring and cost attribution per Airflow DAG.
  2. Monitor Snowpipe status with SHOW PIPES and SYSTEM$PIPE_STATUS.

Security and Performance

  1. Use key-pair authentication for production — avoid password-based connections.
  2. Leverage clustering keys for frequently joined/filtered columns.

Cost Optimization

StrategySavingsImplementation
Suspend idle warehouses60-80%ALTER WAREHOUSE ... SUSPEND
Use scaling policy30-50%Set AUTO_SUSPEND and AUTO_RESUME
Query taggingAttributionSET QUERY_TAG = 'airflow_etl'
Clustering keys2-10xCLUSTER BY (column)

See Also

—
☆☆☆☆☆
0 ratings

Rate & Feedback

Need Expert Airflow Help?

Get personalized tutoring, project support, or professional consulting.

Advertisement