Data reaches Snowflake as files: exports from an application database, nightly drops from a partner, logs written to cloud storage. The standard way to load them is a stage (where files wait), a file format (how to read them) and COPY INTO (load them into a table). In this part we set up all three, load the TPC-H history into RAW.TPCH, and learn how to check exactly what was loaded.
Code for this part: snowflake/02_load_raw.sql, airflow/dags/sql/load_raw_orders.sql. Full project: tpch_analytics on GitHub.
Run this part's SQL as the LOADER role on LOADING_WH, created in part 2. The full script is snowflake/02_load_raw.sql.
Step 1: a file format and a stage
The file format describes the files once, so every COPY statement can reuse it. The stage is a named location for files. An internal stage lives inside Snowflake; an external stage points at an S3 bucket, Azure container or GCS bucket. COPY works the same way with both, so this tutorial uses an internal stage to avoid cloud setup.
SQL · SnowflakeUSE ROLE LOADER;
USE WAREHOUSE LOADING_WH;
USE SCHEMA RAW.TPCH;
-- 1. How to read the files
CREATE FILE FORMAT IF NOT EXISTS CSV_GZ
TYPE = CSV
COMPRESSION = GZIP
FIELD_OPTIONALLY_ENCLOSED_BY = '"'
SKIP_HEADER = 1
NULL_IF = ('')
EMPTY_FIELD_AS_NULL = TRUE;
-- 2. Where the files land (an internal stage; an S3/ADLS/GCS external stage works the same)
CREATE STAGE IF NOT EXISTS LANDING
FILE_FORMAT = CSV_GZ
COMMENT = 'Files waiting to be loaded into RAW.TPCH';Step 2: the raw tables
The five small reference tables change rarely, so we copy them straight from the sample data. The two large tables, ORDERS and LINEITEM, get the same columns as the source plus two audit columns: _LOADED_AT (when the row arrived) and _SOURCE_FILE (which file it came from). WHERE FALSE creates the table structure without copying any rows.
SQL · Snowflake-- 3. Small reference tables: copy straight from the sample data
CREATE OR REPLACE TABLE CUSTOMER AS SELECT * FROM SNOWFLAKE_SAMPLE_DATA.TPCH_SF1.CUSTOMER;
CREATE OR REPLACE TABLE NATION AS SELECT * FROM SNOWFLAKE_SAMPLE_DATA.TPCH_SF1.NATION;
CREATE OR REPLACE TABLE REGION AS SELECT * FROM SNOWFLAKE_SAMPLE_DATA.TPCH_SF1.REGION;
CREATE OR REPLACE TABLE PART AS SELECT * FROM SNOWFLAKE_SAMPLE_DATA.TPCH_SF1.PART;
CREATE OR REPLACE TABLE SUPPLIER AS SELECT * FROM SNOWFLAKE_SAMPLE_DATA.TPCH_SF1.SUPPLIER;
-- 4. Fact tables: same columns as the source plus two audit columns
CREATE TABLE IF NOT EXISTS ORDERS AS
SELECT *, CURRENT_TIMESTAMP()::TIMESTAMP_LTZ AS _LOADED_AT, ''::VARCHAR AS _SOURCE_FILE
FROM SNOWFLAKE_SAMPLE_DATA.TPCH_SF1.ORDERS WHERE FALSE;
CREATE TABLE IF NOT EXISTS LINEITEM AS
SELECT *, CURRENT_TIMESTAMP()::TIMESTAMP_LTZ AS _LOADED_AT, ''::VARCHAR AS _SOURCE_FILE
FROM SNOWFLAKE_SAMPLE_DATA.TPCH_SF1.LINEITEM WHERE FALSE;The audit columns pay off later: _LOADED_AT drives dbt's source freshness checks (part 6) and incremental models (part 7), and _SOURCE_FILE lets you trace any bad row back to the file it came from.
Step 3: play the source system
In a real project another system writes the files. Here we produce them ourselves with COPY INTO a stage (the reverse direction, called unloading). We write everything before July 1998 now and keep the last weeks back, so part 7 can deliver them as "tomorrow's" data.
SQL · SnowflakeCOPY INTO @LANDING/orders/batch_1998_06/
FROM (SELECT * FROM SNOWFLAKE_SAMPLE_DATA.TPCH_SF1.ORDERS WHERE O_ORDERDATE < '1998-07-01')
HEADER = TRUE;
COPY INTO @LANDING/lineitem/batch_1998_06/
FROM (SELECT * FROM SNOWFLAKE_SAMPLE_DATA.TPCH_SF1.LINEITEM
WHERE L_ORDERKEY IN (SELECT O_ORDERKEY FROM SNOWFLAKE_SAMPLE_DATA.TPCH_SF1.ORDERS
WHERE O_ORDERDATE < '1998-07-01'))
HEADER = TRUE;
LIST @LANDING;Snowflake splits large unloads into several compressed files automatically; LIST shows each file with its size.
Step 4: look before you load
You can query staged files directly with positional columns ($1, $2…). It is the quickest way to confirm the delimiter, header and column order are what you expect:
SQL · SnowflakeSELECT $1 AS orderkey, $2 AS custkey, $5 AS orderdate, METADATA$FILENAME
FROM @LANDING/orders/ (FILE_FORMAT => 'CSV_GZ')
LIMIT 5;Step 5: load with COPY INTO
This is the statement Airflow will run every day in part 8, saved as airflow/dags/sql/load_raw_orders.sql. It selects the file columns in order and adds the two audit values. METADATA$FILENAME is a column Snowflake provides for every staged file.
SQL · Snowflake-- Load any new files from the stage. COPY INTO remembers which files it has
-- already loaded (for 64 days), so re-running this never duplicates data.
COPY INTO RAW.TPCH.ORDERS
FROM (
SELECT $1, $2, $3, $4, $5, $6, $7, $8, $9,
CURRENT_TIMESTAMP(), METADATA$FILENAME
FROM @RAW.TPCH.LANDING/orders/
)
FILE_FORMAT = (FORMAT_NAME = 'RAW.TPCH.CSV_GZ')
ON_ERROR = 'ABORT_STATEMENT';
COPY INTO RAW.TPCH.LINEITEM
FROM (
SELECT $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16,
CURRENT_TIMESTAMP(), METADATA$FILENAME
FROM @RAW.TPCH.LANDING/lineitem/
)
FILE_FORMAT = (FORMAT_NAME = 'RAW.TPCH.CSV_GZ')
ON_ERROR = 'ABORT_STATEMENT';Why running it twice is safe
Snowflake keeps load metadata for each table: the name and checksum of every file loaded in the last 64 days. Run the same COPY again and it reports that no new files were found, loading nothing. That makes the daily job idempotent: Airflow can retry it freely. The exceptions: FORCE = TRUE reloads everything (and duplicates rows), and files older than 64 days need LOAD_UNCERTAIN_FILES = TRUE if you must reload them, so keep stage paths tidy and archive processed files.
Choosing ON_ERROR
| Option | What happens on a bad row | Use when |
|---|---|---|
| ABORT_STATEMENT (default for bulk loads) | The whole COPY fails; nothing is loaded | Scheduled pipelines: fail loudly and fix the file |
| SKIP_FILE | The file with the error is skipped; other files load | Many independent files where one bad file should not block the rest |
| CONTINUE | Bad rows are skipped, good rows load | Exploration only; silent data loss in production |
Step 6: check what happened
COPY_HISTORY reports each file: whether it loaded, how many rows were parsed and loaded, and the first error if any.
SQL · SnowflakeSELECT FILE_NAME, STATUS, ROW_COUNT, ROW_PARSED, FIRST_ERROR_MESSAGE, LAST_LOAD_TIME
FROM TABLE(INFORMATION_SCHEMA.COPY_HISTORY(
TABLE_NAME => 'RAW.TPCH.ORDERS',
START_TIME => DATEADD(HOUR, -24, CURRENT_TIMESTAMP())));A quick sanity check that RAW matches the source for the loaded period:
SQL · SnowflakeSELECT COUNT(*), MIN(O_ORDERDATE), MAX(O_ORDERDATE), MAX(_LOADED_AT)
FROM RAW.TPCH.ORDERS;The latest order date should be just before 1 July 1998, and the row count should match the same filter on SNOWFLAKE_SAMPLE_DATA.TPCH_SF1.ORDERS.
Scheduled COPY or Snowpipe?
Snowpipe runs COPY automatically whenever a file lands in an external stage, using cloud event notifications. It suits continuous, low-latency feeds. A scheduled COPY from Airflow, as in this series, suits daily batches: it is simpler to set up, runs on your own warehouse, and sits in the same DAG as the dbt build that depends on it.
| Scheduled COPY (Airflow) | Snowpipe | |
|---|---|---|
| Latency | As often as you schedule it | Typically within minutes of a file arriving |
| Compute | Your warehouse, billed per second | Snowflake-managed, billed separately |
| Setup | Stage + COPY statement | External stage + pipe + cloud notifications |
| Best for | Daily or hourly batches | Continuous feeds, many small files |
Next
RAW now holds history up to June 1998. In part 4 we install dbt, connect it to Snowflake with the key pair from part 2, declare these tables as dbt sources and run our first model.
The series
- The stack and what we will build
- Setting up Snowflake for dbt
- Loading raw data into Snowflake (this post)
- Your first dbt project on Snowflake
- Modeling with dbt: staging to marts
- Testing and documenting with dbt
- Incremental models and snapshots
- Orchestrating dbt with Airflow
- CI/CD for dbt with GitHub Actions
- Running the pipeline in production
All parts: Data Pipeline Series. Every file shown in this series is in the companion project tpch_analytics on GitHub; clone it to follow along, or run it locally on DuckDB without a Snowflake account.
Comments
Post a Comment