Snowflake Ingest
Move data from Snowflake into the data lake. Assumes a Glue SNOWFLAKE connection exists. If not, delegate to connecting-to-data-source.
Contents
Prerequisites
- Glue connection of type
SNOWFLAKE(not JDBC) - Source database, schema, table, and optional query
- Target table in data lake
- Warehouse sized for the read workload (larger warehouse = faster read, more cost)
Read Pattern
The Glue Snowflake connector reads via Snowflake's COPY INTO mechanism under the hood -- efficient for large extracts.
snowflake_df = glueContext.create_dynamic_frame.from_options(
connection_type="snowflake",
connection_options={
"connectionName": args['connection_name'],
"sfDatabase": args['database'],
"sfSchema": args['schema'],
"dbtable": args['table']
}
).toDF()For custom SQL, use query instead of dbtable:
connection_options={
"connectionName": args['connection_name'],
"query": "SELECT id, name, updated_at FROM SALES.ORDERS WHERE status = 'CLOSED'"
}Incremental Loading
Snowflake has reliable timestamps on most tables. Common watermark columns:
- Application-maintained
updated_at/modified_at - Snowflake-maintained
_FIVETRAN_SYNCEDif sourced via Fivetran INFORMATION_SCHEMA.TABLES.LAST_ALTEREDfor schema-level freshness (not row-level)
For tables without an updated_at, options:
- Query
SNOWFLAKE.ACCOUNT_USAGE.QUERY_HISTORYorTABLE_STORAGE_METRICSto identify changed tables for full refresh scheduling - Use Snowflake Streams to capture CDC (advanced; requires Snowflake-side setup -- see Snowflake Streams docs)
Standard watermark filter in the custom query:
connection_options={
"connectionName": args['connection_name'],
"query": f"SELECT * FROM {source_table} WHERE updated_at > '{last_watermark}'"
}See incremental-loading.md for watermark storage and the broader incremental pattern.
Partition Pruning
Snowflake tables are automatically micro-partitioned. Push down filters via the query option -- do not pull full tables and filter in Spark.
Clustered tables benefit most from filter push-down. Check cluster keys:
SHOW TABLES LIKE '<table>' IN SCHEMA <db>.<schema>;
-- Look at CLUSTER_BY columnIf the source table is clustered on created_date and you filter on created_date >= '2026-01-01', Snowflake prunes micro-partitions and returns only relevant data.
Type Mapping
| Snowflake | Iceberg | Notes |
|---|---|---|
| VARCHAR, STRING, TEXT | STRING | |
| NUMBER(p,s) | DECIMAL(p,s) | |
| NUMBER (no scale) | BIGINT | |
| FLOAT, DOUBLE | DOUBLE | |
| BOOLEAN | BOOLEAN | |
| DATE | DATE | |
| TIME | STRING | Iceberg has no TIME type |
| TIMESTAMP_NTZ | TIMESTAMP | Naive timestamp |
| TIMESTAMP_LTZ, TIMESTAMP_TZ | TIMESTAMPTZ | Timezone-aware |
| VARIANT | STRING | Serialize as JSON |
| OBJECT | STRUCT or STRING | Flatten or serialize |
| ARRAY | ARRAY or STRING | |
| BINARY | BINARY | |
| GEOGRAPHY, GEOMETRY | STRING | GeoJSON or WKT |