Snowflake Integration Architecture Airflow Scheduler SnowflakeHook Connection + SQL Snowflake Data warehouse External Stage S3 / GCS / Azure Snowpipe Auto-ingest SnowflakeOperator Execute SQL GCSToSnowflake Load from GCS CopyFromStage COPY INTO TransferOperator Cross-region copy Use 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
Operator Use Case Key Parameters SnowflakeOperator Execute SQL sql, snowflake_conn_idSnowflakeToGCSOperator Export to GCS gcs_bucket, sqlGCSToSnowflakeOperator Load from GCS bucket, prefix, tableSnowflakeCopyFromExternalStage COPY INTO table, stage, file_formatSnowflakeSensor Wait for condition sql, poke_interval
COPY INTO Options
Option Description Common Values FILE_FORMAT File type specification TYPE = 'PARQUET', TYPE = 'CSV'PATTERN Regex for file matching PATTERN = '.*\\.parquet'ON_ERROR Error handling strategy SKIP_FILE, CONTINUE, ABORT_STATEMENTFORCE Reload existing files FORCE = TRUEPURGE Delete files after load PURGE = TRUE
Connection Setup
# Airflow connection for Snowflake
# Connection ID: snowflake_default
# Connection Type: Snowflake
# Host: <account>.snowflakecomputing.com
# Schema: PUBLIC
# Database: ANALYTICS
# Warehouse: COMPUTE_WH
# Account: <account_identifier>
# Login: <username>
# Password: <password> (or use private_key for key-pair auth)
# Extra: {
# "region": "us-east-1",
# "warehouse": "COMPUTE_WH",
# "database": "ANALYTICS",
# "role": "SYSADMIN"
# }
Basic SQL Execution
from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
with DAG(
dag_id='snowflake_etl',
start_date=datetime(2024, 1, 1),
schedule_interval='@daily',
catchup=False,
tags=['snowflake', 'etl'],
) as dag:
create_schema = SnowflakeOperator(
task_id='create_schema',
snowflake_conn_id='snowflake_default',
sql="""
CREATE SCHEMA IF NOT EXISTS raw;
CREATE SCHEMA IF NOT EXISTS staging;
CREATE SCHEMA IF NOT EXISTS analytics;
""",
)
create_table = SnowflakeOperator(
task_id='create_orders_table',
snowflake_conn_id='snowflake_default',
sql="""
CREATE TABLE IF NOT EXISTS staging.orders (
order_id INTEGER,
customer_id INTEGER,
amount DECIMAL(10,2),
order_date TIMESTAMP_NTZ,
status VARCHAR(20)
)
CLUSTER BY (order_date);
""",
)
load_data = SnowflakeOperator(
task_id='load_orders',
snowflake_conn_id='snowflake_default',
sql="""
COPY INTO staging.orders
FROM @raw_stage/orders/
FILE_FORMAT = (
TYPE = 'PARQUET'
COMPRESSION = 'SNAPPY'
)
PATTERN = '.*\\.parquet'
ON_ERROR = 'SKIP_FILE';
""",
)
transform = SnowflakeOperator(
task_id='transform_orders',
snowflake_conn_id='snowflake_default',
sql="""
MERGE INTO analytics.fct_orders AS target
USING (
SELECT
order_id,
customer_id,
SUM(amount) as total_amount,
COUNT(*) as order_count,
MAX(order_date) as last_order_date
FROM staging.orders
WHERE order_date >= DATEADD(day, -1, CURRENT_DATE())
GROUP BY order_id, customer_id
) AS source
ON target.order_id = source.order_id
WHEN MATCHED THEN UPDATE SET
total_amount = source.total_amount,
order_count = source.order_count,
last_order_date = source.last_order_date
WHEN NOT MATCHED THEN INSERT
(order_id, customer_id, total_amount, order_count, last_order_date)
VALUES (source.order_id, source.customer_id, source.total_amount,
source.order_count, source.last_order_date);
""",
)
create_schema >> create_table >> load_data >> transform
GCS to Snowflake Loading
from airflow.providers.snowflake.operators.snowflake import (
SnowflakeOperator,
)
# Step 1: Create a named stage pointing to GCS
create_stage = SnowflakeOperator(
task_id='create_gcs_stage',
snowflake_conn_id='snowflake_default',
sql="""
CREATE OR REPLACE STAGE raw.raw_stage
URL = 'gcs://data-lake/raw/orders/'
STORAGE_INTEGRATION = gcs_integration;
""",
)
# Step 2: List files in stage
list_files = SnowflakeOperator(
task_id='list_stage_files',
snowflake_conn_id='snowflake_default',
sql="LIST @raw.raw_stage;",
)
# Step 3: Load with COPY INTO
load_from_stage = SnowflakeOperator(
task_id='load_from_stage',
snowflake_conn_id='snowflake_default',
sql="""
COPY INTO staging.orders
FROM @raw.raw_stage
FILE_FORMAT = (
TYPE = 'CSV'
FIELD_DELIMITER = '|'
SKIP_HEADER = 1
NULL_IF = ('NULL', 'null', '')
)
PATTERN = 'orders_.*\\.csv'
ON_ERROR = 'CONTINUE';
""",
)
create_stage >> list_files >> load_from_stage
Pre/Post Operator Pattern
from airflow.providers.snowflake.operators.snowflake import (
SnowflakeOperator,
)
# Pre-task: Set session parameters
pre_task = SnowflakeOperator(
task_id='set_session_params',
snowflake_conn_id='snowflake_default',
sql="""
ALTER SESSION SET TIMEZONE = 'UTC';
ALTER SESSION SET STATEMENT_TIMEOUT_IN_SECONDS = 3600;
ALTER SESSION SET QUERY_TAG = 'airflow_etl_{{ ds }}';
""",
)
# Main ETL task
etl_task = SnowflakeOperator(
task_id='run_etl',
snowflake_conn_id='snowflake_default',
sql="CALL run_etl_procedure('{{ ds }}');",
)
# Post-task: Clean up
post_task = SnowflakeOperator(
task_id='cleanup',
snowflake_conn_id='snowflake_default',
sql="""
ALTER SESSION UNSET QUERY_TAG;
TRUNCATE TABLE staging.temp_data;
""",
)
pre_task >> etl_task >> post_task
Multi-Warehouse Pattern
from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
warehouses = {
'heavy_etl': 'COMPUTE_WH_HEAVY',
'light_transform': 'COMPUTE_WH_LIGHT',
'ad_hoc': 'COMPUTE_WH_ADHOC',
}
for task_name, warehouse in warehouses.items():
SnowflakeOperator(
task_id=f'task_{task_name}',
snowflake_conn_id='snowflake_default',
sql=f"""
ALTER WAREHOUSE {warehouse} RESUME IF SUSPENDED;
-- ETL logic here
""",
warehouse=warehouse,
)
Key Concepts Table
Operator Purpose Key Parameters SnowflakeOperator Execute SQL sql, snowflake_conn_idSnowflakeToGCSOperator Export to GCS snowflake_conn_id, gcs_bucket, sqlGCSToSnowflakeOperator Load from GCS bucket, prefix, table, stageSnowflakeCopyFromExternalStage COPY INTO table, stage, file_formatSnowflakePrePostOperator Pre/post hooks sql, pre_query, post_querySnowflakeSensor Wait for condition sql, poke_interval
COPY INTO Options
Option Description Example FILE_FORMAT File type specification TYPE = 'PARQUET'PATTERN Regex for file matching PATTERN = '.*\\.parquet'ON_ERROR Error handling strategy SKIP_FILE, CONTINUE, ABORT_STATEMENTFORCE Reload existing files FORCE = TRUESIZE_LIMIT Max bytes per load SIZE_LIMIT = 1073741824PURGE Delete files after load PURGE = TRUERETURN_FAILED_ONLY Return only failed files RETURN_FAILED_ONLY = TRUE
Best Practices
Data Loading
Use named stages for reusable external connections â avoid embedded credentials.
Set ON_ERROR appropriately : SKIP_FILE for bad files, CONTINUE for partial loads.
Partition large COPY INTO operations with SIZE_LIMIT to avoid single large transactions.
Use PURGE = TRUE to clean up files after successful load.
Monitoring and Cost
Set QUERY_TAG for monitoring and cost attribution per Airflow DAG.
Monitor Snowpipe status with SHOW PIPES and SYSTEM$PIPE_STATUS.
Security and Performance
Use key-pair authentication for production â avoid password-based connections.
Leverage clustering keys for frequently joined/filtered columns.
Cost Optimization
Strategy Savings Implementation Suspend idle warehouses 60-80% ALTER WAREHOUSE ... SUSPENDUse scaling policy 30-50% Set AUTO_SUSPEND and AUTO_RESUME Query tagging Attribution SET QUERY_TAG = 'airflow_etl'Clustering keys 2-10x CLUSTER BY (column)
See Also
BigQuery Provider â Google BigQuery integration patterns
Databricks Provider â Databricks cluster and job management
Operators and Hooks â Operator lifecycle and hook architecture
XCom Communications â Task communication and data passing