September 11, 2025

Unit Testing PySpark Bordereaux Transformations with pytest

In the SQL Server world, testing ETL often meant "run it in UAT and compare row counts". tSQLt existed, but few teams used it. With Spark, the same habit often shows up in notebooks: logic that's only ever tested by running the whole notebook against real bordereaux and eyeballing the output.

That's risky in delegated authority. A small mistake in how commission is netted off, how cancellations are signed, or which resubmission wins can misstate written premium for a whole binding authority, and nobody notices until the reconciliation with the broker fails. PySpark transformations are just Python functions that take DataFrames and return DataFrames, so they can be unit tested with pytest, locally, in seconds, without a cluster. Spark 3.5 also added built-in testing helpers that remove most of the boilerplate. Here's a setup that works well.

Step 1: Get the logic out of the notebook

You can't easily unit test a notebook cell that reads a table, transforms it and writes it back. Separate the three:

 1: # src/transforms/premium.py
 2: from pyspark.sql import DataFrame, functions as F, Window
 3: 
 4: 
 5: def latest_submission(df: DataFrame) -> DataFrame:
 6:     """Keep the latest submitted version of each premium transaction."""
 7:     keys = ["umr", "section_no", "certificate_ref", "transaction_seq"]
 8:     w = Window.partitionBy(*keys).orderBy(F.col("submitted_at").desc(),
 9:                                           F.col("bdx_version").desc())
10:     return (df.withColumn("_rn", F.row_number().over(w))
11:               .filter("_rn = 1")
12:               .drop("_rn"))
13: 
14: 
15: def derive_premium_amounts(df: DataFrame) -> DataFrame:
16:     """Sign cancellations negative and derive commission and net premium."""
17:     sign = F.when(F.col("transaction_type") == "CANCELLATION", F.lit(-1)).otherwise(F.lit(1))
18:     gross = F.abs(F.col("gross_premium")) * sign
19:     commission = F.round(gross * F.coalesce(F.col("commission_pct"), F.lit(0)) / 100, 2)
20:     return (df.withColumn("gross_premium_signed", gross)
21:               .withColumn("commission_amount", commission)
22:               .withColumn("net_premium", gross - commission))

The notebook (or job) becomes a thin wrapper:

1: from transforms.premium import latest_submission, derive_premium_amounts
2: 
3: raw   = spark.read.table("bronze.premium_bdx")
4: clean = derive_premium_amounts(latest_submission(raw))
5: clean.write.mode("overwrite").saveAsTable("silver.premium_transaction")

In Fabric you can package src/ as a wheel and attach it to an Environment. In Glue, use --extra-py-files. On Databricks, a workspace file or wheel. The tests stay the same everywhere.

Step 2: A SparkSession fixture

 1: # tests/conftest.py
 2: import pytest
 3: from pyspark.sql import SparkSession
 4: 
 5: 
 6: @pytest.fixture(scope="session")
 7: def spark():
 8:     spark = (SparkSession.builder
 9:              .master("local[1]")
10:              .appName("unit-tests")
11:              .config("spark.sql.shuffle.partitions", "1")
12:              .config("spark.default.parallelism", "1")
13:              .config("spark.ui.enabled", "false")
14:              .config("spark.sql.session.timeZone", "UTC")
15:              .getOrCreate())
16:     yield spark
17:     spark.stop()

Three settings do most of the work:

  • scope="session": start Spark once for the whole test run. Starting a session takes a few seconds, and doing it per test makes the suite painfully slow.
  • shuffle.partitions = 1: the default of 200 partitions is designed for clusters. On tiny test data it just creates 200 empty tasks.
  • A fixed session time zone: otherwise submission timestamp tests pass on your laptop and fail on the CI runner.

Step 3: Tests with assertDataFrameEqual

Spark 3.5 added pyspark.testing with assertDataFrameEqual and assertSchemaEqual. They give readable diffs when something doesn't match, which is a big step up from comparing collect() output by hand.

 1: # tests/test_premium.py
 2: from datetime import datetime
 3: from pyspark.testing import assertDataFrameEqual
 4: from transforms.premium import latest_submission, derive_premium_amounts
 5: 
 6: BDX_SCHEMA = ("umr STRING, section_no STRING, certificate_ref STRING, "
 7:               "transaction_seq INT, gross_premium DOUBLE, submitted_at TIMESTAMP, "
 8:               "bdx_version INT")
 9: 
10: 
11: def test_resubmission_replaces_original(spark):
12:     src = spark.createDataFrame(
13:         [
14:             ("B0999CH001", "S1", "CERT-1001", 1, 1200.0, datetime(2025, 9, 5, 10, 0), 1),
15:             ("B0999CH001", "S1", "CERT-1001", 1, 1250.0, datetime(2025, 9, 12, 9, 0), 2),  # corrected
16:             ("B0999CH001", "S1", "CERT-1002", 1,  800.0, datetime(2025, 9, 5, 10, 0), 1),
17:         ],
18:         BDX_SCHEMA,
19:     )
20: 
21:     expected = spark.createDataFrame(
22:         [
23:             ("B0999CH001", "S1", "CERT-1001", 1, 1250.0, datetime(2025, 9, 12, 9, 0), 2),
24:             ("B0999CH001", "S1", "CERT-1002", 1,  800.0, datetime(2025, 9, 5, 10, 0), 1),
25:         ],
26:         BDX_SCHEMA,
27:     )
28: 
29:     assertDataFrameEqual(latest_submission(src), expected)
30: 
31: 
32: def test_same_timestamp_uses_highest_version(spark):
33:     ts = datetime(2025, 9, 5, 10, 0)
34:     src = spark.createDataFrame(
35:         [
36:             ("B0999CH001", "S1", "CERT-1001", 1, 1200.0, ts, 1),
37:             ("B0999CH001", "S1", "CERT-1001", 1, 1300.0, ts, 2),
38:         ],
39:         BDX_SCHEMA,
40:     )
41:     result = latest_submission(src).select("gross_premium")
42:     assertDataFrameEqual(result, [(1300.0,)])
43: 
44: 
45: def test_cancellation_is_negative_and_commission_nets_off(spark):
46:     src = spark.createDataFrame(
47:         [
48:             ("CERT-1", "NEW",           1000.0, 25.0),
49:             ("CERT-2", "CANCELLATION",   400.0, 25.0),   # some coverholders send positives
50:             ("CERT-3", "CANCELLATION",  -400.0, 25.0),   # others send negatives
51:             ("CERT-4", "MTA",            100.0, None),   # missing commission
52:         ],
53:         "certificate_ref STRING, transaction_type STRING, gross_premium DOUBLE, commission_pct DOUBLE",
54:     )
55: 
56:     result = derive_premium_amounts(src).select(
57:         "certificate_ref", "gross_premium_signed", "commission_amount", "net_premium")
58: 
59:     expected = spark.createDataFrame(
60:         [
61:             ("CERT-1", 1000.0,  250.0,  750.0),
62:             ("CERT-2", -400.0, -100.0, -300.0),
63:             ("CERT-3", -400.0, -100.0, -300.0),
64:             ("CERT-4",  100.0,    0.0,  100.0),
65:         ],
66:         "certificate_ref STRING, gross_premium_signed DOUBLE, "
67:         "commission_amount DOUBLE, net_premium DOUBLE",
68:     )
69: 
70:     assertDataFrameEqual(result, expected)

Some useful behaviour:

  • Row order is ignored by default. Pass checkRowOrder=True when order matters, for example after a sort.
  • Floating-point values are compared with a tolerance (rtol/atol), so 0.1 + 0.2 doesn't fail your test. For real money amounts, use DECIMAL, which the helpers compare exactly.
  • The expected value can be a list of rows instead of a DataFrame, which is handy for small cases like the version tie-break test.

What to test

Don't test Spark itself. You don't need a test proving filter works. Focus on your business rules, especially:

  • Sign conventions: cancellations and return premiums. Coverholders are inconsistent about whether they send them as positive or negative amounts, as the test above shows.
  • Nulls: missing commission percentages, missing currencies, missing inception dates. Every CASE WHEN needs a null test.
  • Resubmissions and ties: two versions with the same timestamp. Without a tie-breaker, the result isn't deterministic, which is why bdx_version is in the window ordering.
  • Boundaries: risks incepting on the first or last day of the contract section's period. Off-by-one on < vs <= is a classic bug that assigns a risk to the wrong year of account.
  • Schema contracts: use assertSchemaEqual on outputs that the premium fact or Power BI depends on, so a renamed column fails in CI rather than in an underwriter's report.
  • Empty input: a coverholder with no business this month should produce an empty DataFrame with the correct schema, not an error.

Running it in CI

If you already run build and deployment workflows in GitHub Actions, adding a test step is a small change:

1: # requirements-dev.txt
2: pyspark==3.5.*
3: pytest
 1: # .github/workflows/tests.yml
 2: name: tests
 3: on: [push, pull_request]
 4: jobs:
 5:   test:
 6:     runs-on: ubuntu-latest
 7:     steps:
 8:       - uses: actions/checkout@v4
 9:       - uses: actions/setup-java@v4
10:         with: { distribution: temurin, java-version: '17' }
11:       - uses: actions/setup-python@v5
12:         with: { python-version: '3.11' }
13:       - run: pip install -r requirements-dev.txt -e .
14:       - run: pytest -q

Match the PySpark version to your target runtime, whether that's the Fabric runtime, the Glue version or the Databricks Runtime, so behaviour is the same in tests and production. A suite of 50–100 tests like these runs in well under a minute.

Integration tests are still needed

Unit tests catch logic bugs. They won't catch a permissions problem, a missing table, or a coverholder who suddenly changes their bordereau layout. You still want an end-to-end run in a test workspace with a sample month of bordereaux before release. With unit tests in place, that run should fail far less often, and when it does, the cause is more likely to be environmental than a business rule. That's the kind of failure you want to be left with.

June 26, 2025

Mirroring Snowflake into Microsoft Fabric OneLake

Many insurers and MGAs now run both Snowflake and Microsoft Fabric. Bordereaux processing and data engineering live in Snowflake, while Power BI, and increasingly the business-facing analytics for underwriters and DA oversight, lives in Fabric. The question is how to get the data across without building and running yet another set of pipelines.

Mirroring is Fabric's answer. You point a mirrored database at a Snowflake database, choose the tables, and Fabric keeps a near-real-time copy in OneLake as Delta tables. Today I bring Snowflake data into OneLake with Fabric pipelines and shortcuts. Mirroring is the lower-code option, and I've been exploring it with a sample delegated authority model. Here's how to set it up and what to watch for.

What mirroring does

  • Takes an initial snapshot of each selected table, then replicates changes continuously.
  • Writes the data into OneLake in Delta format, inside a Mirrored database item.
  • Gives you a SQL analytics endpoint on the mirrored data and a default semantic model, so Power BI can use Direct Lake.
  • The mirrored tables are read-only in Fabric. Snowflake stays the system of record.

The replication compute and the OneLake storage for mirrored data are largely covered by your Fabric capacity. There's a free mirroring storage allowance that scales with capacity size. But the Snowflake side still costs money, which is the part people forget (more on that below).

What to mirror

Mirror the gold star schema, not the bordereaux staging layers:

  • Facts: fact_risk_transaction, fact_premium_transaction, fact_claim_transaction
  • Dimensions: dim_coverholder, dim_broker, dim_contract (binding authority), dim_contract_section, dim_class_of_business, dim_date

Step 1: Prepare Snowflake

Create a dedicated user and role for Fabric. Mirroring uses Snowflake table streams to pick up changes, so the role needs to be able to create streams and read the tables.

 1: USE ROLE SECURITYADMIN;
 2: 
 3: CREATE ROLE IF NOT EXISTS fabric_mirror_role;
 4: CREATE USER IF NOT EXISTS fabric_mirror_user
 5:   DEFAULT_ROLE = fabric_mirror_role
 6:   DEFAULT_WAREHOUSE = mirror_wh
 7:   TYPE = SERVICE;      -- authentication set up per your org's standards
 8: 
 9: GRANT ROLE fabric_mirror_role TO USER fabric_mirror_user;
10: 
11: USE ROLE SYSADMIN;
12: CREATE WAREHOUSE IF NOT EXISTS mirror_wh
13:   WAREHOUSE_SIZE = XSMALL AUTO_SUSPEND = 60 AUTO_RESUME = TRUE;
14: 
15: GRANT USAGE ON WAREHOUSE mirror_wh                  TO ROLE fabric_mirror_role;
16: GRANT USAGE ON DATABASE  da_analytics               TO ROLE fabric_mirror_role;
17: GRANT USAGE ON SCHEMA    da_analytics.gold          TO ROLE fabric_mirror_role;
18: GRANT SELECT ON ALL TABLES IN SCHEMA da_analytics.gold TO ROLE fabric_mirror_role;
19: GRANT CREATE STREAM ON SCHEMA da_analytics.gold     TO ROLE fabric_mirror_role;

Check the current Microsoft documentation for the exact privilege list before you set this up. It's short, but it has changed as the feature has matured. Using a dedicated warehouse lets you see mirroring's Snowflake cost on its own line.

Step 2: Create the mirrored database in Fabric

  1. In your workspace: New item → Mirrored Snowflake.
  2. Create a connection: the Snowflake account server name (<account>.snowflakecomputing.com), the warehouse, and the credentials for fabric_mirror_user.
  3. Choose the database, then either mirror all data or select specific tables. Selecting tables explicitly is safer. "All data" also picks up new tables automatically, which you may not want.
  4. Click Mirror database. The status page shows each table move from Running through the initial snapshot to Replicating, with row counts.

Step 3: Use it

From the SQL analytics endpoint you can query the tables with T-SQL, create views, and, because everything is in OneLake, join mirrored Snowflake data with Lakehouse or Warehouse data in the same workspace. Here's a loss ratio view by coverholder and year of account, combining mirrored facts with a planning table that the underwriting team maintains in a Fabric Lakehouse:

 1: WITH prem AS (
 2:     SELECT contract_section_key, SUM(gross_premium_gbp) AS gwp_gbp
 3:     FROM   da_mirror.gold.fact_premium_transaction
 4:     GROUP  BY contract_section_key
 5: ),
 6: clm AS (
 7:     SELECT contract_section_key,
 8:            SUM(paid_gbp + outstanding_gbp) AS incurred_gbp
 9:     FROM   da_mirror.gold.fact_claim_transaction
10:     WHERE  is_latest_position = 1
11:     GROUP  BY contract_section_key
12: )
13: SELECT  ch.coverholder_name,
14:         cs.year_of_account,
15:         cs.class_of_business,
16:         SUM(p.gwp_gbp)                                         AS gwp_gbp,
17:         SUM(COALESCE(c.incurred_gbp, 0))                       AS incurred_gbp,
18:         SUM(COALESCE(c.incurred_gbp, 0)) / NULLIF(SUM(p.gwp_gbp), 0) AS incurred_loss_ratio,
19:         MAX(pl.planned_loss_ratio)                             AS planned_loss_ratio
20: FROM    prem p
21: JOIN    da_mirror.gold.dim_contract_section AS cs ON cs.contract_section_key = p.contract_section_key
22: JOIN    da_mirror.gold.dim_coverholder      AS ch ON ch.coverholder_key = cs.coverholder_key
23: LEFT JOIN clm c                             ON c.contract_section_key = p.contract_section_key
24: LEFT JOIN uw_lakehouse.dbo.business_plan    AS pl
25:        ON pl.class_of_business = cs.class_of_business
26:       AND pl.year_of_account   = cs.year_of_account
27: GROUP   BY ch.coverholder_name, cs.year_of_account, cs.class_of_business;

For Power BI, build a semantic model on the mirrored database in Direct Lake mode. Reports now read Snowflake data with no import refresh and no DirectQuery round trips to Snowflake. If your Power BI reports currently use DirectQuery against Snowflake, this can remove a noticeable share of your Snowflake credit usage.

Things to watch for

  • It's near-real-time, not real-time. Changes usually show up within a minute or two, sometimes longer after big month-end loads. That's more than enough for bordereaux, which arrive monthly anyway.
  • Snowflake compute still runs. Picking up changes means queries against the streams on your warehouse. During month-end, when the facts change constantly, mirror_wh stays busy. Watch WAREHOUSE_METERING_HISTORY for the first couple of month-ends.
  • Prefer mirroring curated tables. Mirror the gold layer, not bordereaux staging tables that get truncated and reloaded. A full reload in Snowflake means a full re-replication.
  • Data type mapping. Most types map cleanly, but check VARIANT/semi-structured columns and very high-precision numbers and timestamps, which may arrive as strings or with reduced precision. Keep premium and claim amounts as fixed-precision NUMBER columns, and flatten any semi-structured bordereau fields into typed columns in Snowflake first.
  • There's a table limit per mirrored database (a few hundred). For big estates, split by subject area into several mirrored databases. That also keeps permissions tidier.
  • Watch out for schema changes. Adding columns is generally handled, but some DDL changes on the source can require the table to be re-synced. Coordinate with whoever owns the Snowflake models.
  • Security doesn't carry over. This is the big one in DA. Snowflake row access policies and masking policies don't come with the data. Mirroring reads with the mirror user's rights. If coverholders or class underwriters should only see their own business, or claimant details must be masked, re-apply that in Fabric (SQL endpoint permissions, RLS in the semantic model, or OneLake security), or only mirror data that's already safe for the audience.

Mirroring vs. other options

OptionGood forWatch out for
MirroringLow-maintenance, near-real-time replication of curated tablesSnowflake compute for change capture, security re-implementation
Pipeline / Copy jobScheduled batch copies, transformations on the wayYou own the pipeline, watermarks and failures
Iceberg + OneLake shortcutOne copy of the data, readable by both platformsMore setup, and you need to understand Iceberg catalogs
Power BI DirectQuery to SnowflakeAlways-current data, no copyReport speed and Snowflake cost at scale

For a classic "Snowflake for bordereaux engineering, Power BI for underwriting and oversight" setup, mirroring looks like the simplest option. It needs less code than any other way of moving data between the two platforms.

April 17, 2025

Snowflake Cost Control: Warehouses, Resource Monitors and Finding Where Credits Go

The first Snowflake invoice is often a surprise. Not because Snowflake is expensive, but because consumption pricing behaves very differently from a SQL Server licence. On-premises, a badly written query costs you nothing but time. In Snowflake it costs credits, and a warehouse left running over a weekend costs a lot of credits.

Delegated authority workloads have a particular shape: quiet for most of the month, then a rush in the first working days when coverholders' risk, premium and claims bordereaux are loaded, validated, and reported to underwriters and oversight teams. That pattern is a good fit for Snowflake's pricing, but only if the account is set up for it. Here's the checklist I go through, with the queries I use to find where the money goes.

How you're billed, in one paragraph

Compute is billed in credits per hour per running virtual warehouse. An X-Small uses 1 credit/hour, and each size up doubles that. Billing is per second, with a 60-second minimum each time a warehouse resumes. Storage is billed separately per TB per month, and serverless features (Snowpipe, automatic clustering, search optimisation, serverless tasks, and so on) have their own credit usage. In most accounts, warehouses are 70–90% of the bill, so start there.

1. Auto-suspend and auto-resume on every warehouse

1: SHOW WAREHOUSES;
2: -- check the auto_suspend column: anything NULL or > 300 deserves a question
3: 
4: ALTER WAREHOUSE bdx_reporting_wh SET AUTO_SUSPEND = 60 AUTO_RESUME = TRUE;
5: ALTER WAREHOUSE bdx_etl_wh       SET AUTO_SUSPEND = 60 AUTO_RESUME = TRUE;

60 seconds is a good default. Going lower rarely helps because of the 60-second minimum on resume. The exception is a warehouse serving a busy Power BI model, where a slightly longer suspend (2–5 minutes) keeps the local cache warm and can actually save credits by avoiding repeated cold starts.

2. Separate warehouses by workload

One giant warehouse for everything makes it impossible to see who is spending what. I split by workload:

  • bdx_etl_wh: bordereaux ingestion, validation and transformation into the risk, premium and claims facts
  • bdx_reporting_wh: Power BI, including GWP, GPI utilisation and loss ratio dashboards
  • actuarial_adhoc_wh: actuaries and analysts querying history
  • dev_wh: development, kept small

Then grant USAGE per role, so people can only use the warehouse meant for them.

3. Resource monitors as a safety net

 1: USE ROLE ACCOUNTADMIN;
 2: 
 3: CREATE OR REPLACE RESOURCE MONITOR rm_actuarial_adhoc
 4:   WITH CREDIT_QUOTA = 200
 5:        FREQUENCY = MONTHLY
 6:        START_TIMESTAMP = IMMEDIATELY
 7:   TRIGGERS
 8:     ON 75  PERCENT DO NOTIFY
 9:     ON 90  PERCENT DO NOTIFY
10:     ON 100 PERCENT DO SUSPEND
11:     ON 110 PERCENT DO SUSPEND_IMMEDIATE;
12: 
13: ALTER WAREHOUSE actuarial_adhoc_wh SET RESOURCE_MONITOR = rm_actuarial_adhoc;

SUSPEND lets running queries finish. SUSPEND_IMMEDIATE cancels them. For bordereaux ETL, I set monitors to notify only. Suspending the warehouse halfway through month-end processing just delays the reports underwriters are waiting for. For ad-hoc and dev, hard limits are fine. Make sure notifications are enabled for the account admins in their user preferences, or the alerts go nowhere.

4. Statement timeouts

The default statement timeout is 2 days. That's a long time for an accidental cartesian join between the risk and claims facts to run on a Large warehouse.

1: ALTER WAREHOUSE actuarial_adhoc_wh SET STATEMENT_TIMEOUT_IN_SECONDS = 1800;   -- 30 min
2: ALTER WAREHOUSE bdx_reporting_wh   SET STATEMENT_TIMEOUT_IN_SECONDS = 600;    -- 10 min
3: ALTER WAREHOUSE bdx_etl_wh         SET STATEMENT_TIMEOUT_IN_SECONDS = 7200;   -- 2 h
4: 
5: ALTER WAREHOUSE actuarial_adhoc_wh SET STATEMENT_QUEUED_TIMEOUT_IN_SECONDS = 600;

5. Find where the credits go

These views live in SNOWFLAKE.ACCOUNT_USAGE and lag behind real time by up to a few hours.

Credits by warehouse, last 30 days:

1: SELECT warehouse_name,
2:        ROUND(SUM(credits_used), 1)                     AS credits,
3:        ROUND(SUM(credits_used_cloud_services), 1)      AS cloud_services
4: FROM   snowflake.account_usage.warehouse_metering_history
5: WHERE  start_time >= DATEADD(day, -30, CURRENT_TIMESTAMP())
6: GROUP  BY warehouse_name
7: ORDER  BY credits DESC;

Credits by day of month, which shows the month-end bordereaux peak clearly and helps decide whether a bigger warehouse for those days is worth it:

1: SELECT warehouse_name,
2:        DAY(start_time)               AS day_of_month,
3:        ROUND(AVG(credits_used), 2)   AS avg_credits_per_hour
4: FROM   snowflake.account_usage.warehouse_metering_history
5: WHERE  start_time >= DATEADD(month, -3, CURRENT_TIMESTAMP())
6: GROUP  BY 1, 2
7: ORDER  BY 1, 2;

Run the same query with HOUR(start_time) to spot warehouses that never sleep. If a warehouse shows steady usage at 3am every night and nothing is scheduled then, something (often a BI tool's keep-alive or a forgotten refresh) is keeping it awake.

The most expensive query patterns, grouped by query hash so the same query with different literals counts once:

 1: SELECT query_parameterized_hash,
 2:        ANY_VALUE(query_text)                         AS sample_query,
 3:        ANY_VALUE(warehouse_name)                     AS warehouse,
 4:        COUNT(*)                                      AS executions,
 5:        ROUND(SUM(total_elapsed_time) / 1000 / 3600, 2) AS total_hours,
 6:        ROUND(AVG(bytes_scanned) / POWER(1024, 3), 2)   AS avg_gb_scanned
 7: FROM   snowflake.account_usage.query_history
 8: WHERE  start_time >= DATEADD(day, -7, CURRENT_TIMESTAMP())
 9:   AND  warehouse_name IS NOT NULL
10: GROUP  BY query_parameterized_hash
11: ORDER  BY total_hours DESC
12: LIMIT  25;

Elapsed time isn't the same as credits, since several queries share a running warehouse. But as a ranking it reliably points at the right places. In DA accounts the top of this list is often a Power BI query recomputing GWP or incurred claims across every year of account on every refresh. A pre-aggregated mart by coverholder, contract section and bordereau month usually fixes it. Snowflake also provides QUERY_ATTRIBUTION_HISTORY for per-query compute cost, which is worth using if your account has it.

6. Right-size warehouses

Doubling a warehouse doubles the cost per second. If the query runs twice as fast, the total cost is the same and the result arrives sooner. If it doesn't speed up, you're wasting credits. In the query profile, look for:

  • "Bytes spilled to local/remote storage": the warehouse is too small for the query's memory needs. Going up a size can be cheaper.
  • Queuing on the reporting warehouse during month-end, when every underwriter opens the dashboards at once: use a multi-cluster warehouse (Enterprise edition) with SCALING_POLICY = ECONOMY rather than a bigger size.
  • Small queries on a big warehouse: scale down. Looking up one coverholder's contract sections doesn't need a Large.

For the month-end peak, a scheduled task can resize the ETL warehouse up on working day 1 and back down on day 6. It's simple, and it beats running a Large all month.

7. Look at storage too

1: SELECT table_catalog, table_schema, table_name,
2:        ROUND(active_bytes      / POWER(1024,4), 3) AS active_tb,
3:        ROUND(time_travel_bytes / POWER(1024,4), 3) AS time_travel_tb,
4:        ROUND(failsafe_bytes    / POWER(1024,4), 3) AS failsafe_tb
5: FROM   snowflake.account_usage.table_storage_metrics
6: WHERE  deleted = FALSE
7: ORDER  BY (active_bytes + time_travel_bytes + failsafe_bytes) DESC
8: LIMIT  20;

Tables where Time Travel and Fail-safe are much larger than active storage are usually bordereaux staging tables that get fully reloaded each month. Make them TRANSIENT with short retention.

Make it a habit

None of this is a one-off. Put the queries above into a Snowsight dashboard, review it every Monday (and straight after each month-end), and set a budget alert with Snowflake's Budgets feature. A Snowflake account that's watched weekly rarely produces a bad surprise. One that's set up and forgotten nearly always does.

February 20, 2025

Incremental Loads in Microsoft Fabric Pipelines with the Watermark Pattern

If you built ETL in SSIS or Azure Data Factory, you've built a watermark pattern: remember the last ModifiedDate you loaded, pull only rows changed since then, and move the watermark forward. It's still the most reliable way to load incrementally from a relational source that doesn't have CDC enabled. I've built it many times in SSIS and Azure Data Factory, and this post shows how I build it in Microsoft Fabric Data Factory pipelines, metadata-driven so one pipeline handles many tables.

The source in this example is the SQL Server database behind a delegated authority platform: the system that holds binding authority contracts and sections, coverholders and brokers, and the processed risk, premium and claim transactions from each month's bordereaux.

The target layout

  • Source: the DA platform's SQL Server database, reached through the on-premises data gateway (an Azure SQL database works the same way).
  • Landing: a Lakehouse, with incremental extracts written as Parquet to Files/landing/<table>/.
  • Control tables: in a Fabric Warehouse, because I want T-SQL stored procedures and transactions for updating watermarks.

Step 1: The control table

 1: CREATE TABLE etl.watermark
 2: (
 3:     table_name        VARCHAR(128)  NOT NULL,
 4:     source_schema     VARCHAR(128)  NOT NULL,
 5:     watermark_column  VARCHAR(128)  NOT NULL,
 6:     key_columns       VARCHAR(400)  NOT NULL,
 7:     last_watermark    DATETIME2(6)  NOT NULL,
 8:     is_enabled        BIT           NOT NULL,
 9:     last_run_rows     BIGINT        NULL,
10:     last_run_utc      DATETIME2(6)  NULL
11: );
12: 
13: INSERT INTO etl.watermark VALUES
14:  ('RiskTransaction',    'bdx', 'ModifiedDate', 'RiskTransactionId',    '1900-01-01', 1, NULL, NULL),
15:  ('PremiumTransaction', 'bdx', 'ModifiedDate', 'PremiumTransactionId', '1900-01-01', 1, NULL, NULL),
16:  ('ClaimTransaction',   'bdx', 'ModifiedDate', 'ClaimTransactionId',   '1900-01-01', 1, NULL, NULL),
17:  ('ContractSection',    'ref', 'ModifiedDate', 'UMR,SectionNo',        '1900-01-01', 1, NULL, NULL),
18:  ('Coverholder',        'ref', 'ModifiedDate', 'CoverholderId',        '1900-01-01', 1, NULL, NULL),
19:  ('Broker',             'ref', 'ModifiedDate', 'BrokerCode',           '1900-01-01', 1, NULL, NULL);

And the procedure that moves the watermark forward:

 1: CREATE PROCEDURE etl.usp_update_watermark
 2:     @table_name     VARCHAR(128),
 3:     @new_watermark  DATETIME2(6),
 4:     @rows_copied    BIGINT
 5: AS
 6: BEGIN
 7:     UPDATE etl.watermark
 8:     SET    last_watermark = @new_watermark,
 9:            last_run_rows  = @rows_copied,
10:            last_run_utc   = SYSUTCDATETIME()
11:     WHERE  table_name = @table_name;
12: END;

Step 2: The outer pipeline: look up the table list

Create a pipeline pl_da_incremental_master:

  1. Lookup activity LkpTables against the Warehouse, with First row only unticked:
    1: SELECT table_name, source_schema, watermark_column, key_columns, last_watermark
    2: FROM   etl.watermark
    3: WHERE  is_enabled = 1;
  2. ForEach activity with Items = @activity('LkpTables').output.value, sequential off, and a batch count of 4 so the DA platform's database isn't overloaded during the working day.
  3. Inside the ForEach, an Invoke pipeline activity calls the child pipeline and passes @item().table_name, @item().source_schema, @item().watermark_column and @item().last_watermark as parameters.

I prefer a child pipeline over putting everything inside the ForEach, because the child can be run and debugged on its own for one table.

Step 3: The child pipeline: one table, one increment

Activity 1: Lookup LkpNewWatermark (source SQL Server, first row only):

1: @concat('SELECT MAX(', pipeline().parameters.watermark_column,
2:         ') AS new_watermark FROM ', pipeline().parameters.source_schema,
3:         '.', pipeline().parameters.table_name)

Capture the new high-water mark before copying, so rows written during the copy aren't skipped. This matters in the first week of the month, when bordereaux are being processed on the DA platform while the pipeline runs. Those rows are picked up next run instead.

Activity 2: Copy data CopyIncrement

Source query (dynamic content):

1: @concat('SELECT * FROM ', pipeline().parameters.source_schema, '.',
2:         pipeline().parameters.table_name,
3:         ' WHERE ', pipeline().parameters.watermark_column,
4:         ' > ''', pipeline().parameters.last_watermark, '''',
5:         ' AND ',  pipeline().parameters.watermark_column,
6:         ' <= ''', activity('LkpNewWatermark').output.firstRow.new_watermark, '''')

Destination: the Lakehouse, Files, Parquet, with the folder path:

1: @concat('landing/', pipeline().parameters.table_name, '/',
2:         formatDateTime(utcnow(), 'yyyy/MM/dd/HHmmss'))

Activity 3: Stored procedure SpUpdateWatermark (Warehouse), which runs only on Copy success:

  • @table_name = @pipeline().parameters.table_name
  • @new_watermark = @activity('LkpNewWatermark').output.firstRow.new_watermark
  • @rows_copied = @activity('CopyIncrement').output.rowsCopied

If the copy fails, the watermark doesn't move, and the next run retries the same window. That property is what makes this pattern safe to re-run.

Step 4: Merge into Delta tables

The landing files are raw increments. A notebook, called from the master pipeline after the ForEach, merges each one into a silver Delta table using the keys from the control table:

 1: from delta.tables import DeltaTable
 2: 
 3: def merge_increment(table: str, key_columns: str, run_folder: str):
 4:     target_name = f"silver_{table.lower()}"
 5:     inc = spark.read.parquet(f"Files/landing/{table}/{run_folder}")
 6: 
 7:     if not spark.catalog.tableExists(target_name):
 8:         inc.write.format("delta").saveAsTable(target_name)
 9:         return
10: 
11:     condition = " AND ".join(f"t.{k} = s.{k}" for k in key_columns.split(","))
12:     (DeltaTable.forName(spark, target_name).alias("t")
13:         .merge(inc.alias("s"), condition)
14:         .whenMatchedUpdateAll()
15:         .whenNotMatchedInsertAll()
16:         .execute())
17: 
18: merge_increment("PremiumTransaction", "PremiumTransactionId", run_folder)
19: merge_increment("ContractSection", "UMR,SectionNo", run_folder)

In practice the notebook reads the control table and loops over the tables that loaded rows in this run. From silver, the gold star schema (fact_risk_transaction, fact_premium_transaction and fact_claim_transaction, with dim_coverholder, dim_broker and dim_contract_section) is built in the Warehouse.

Gotchas

  • Time zones and precision. If the source column is DATETIME (3.33 ms precision) and you pass the watermark around as a string, format it consistently, for example yyyy-MM-dd HH:mm:ss.fff. Otherwise rounding can cause rows to be skipped or loaded twice.
  • Hard deletes aren't captured. A watermark only sees inserts and updates. If the DA platform physically deletes transactions when a bordereau is rejected and reloaded, you need CDC, soft-delete flags, or a periodic full key comparison for those tables.
  • Unindexed watermark columns. MAX(ModifiedDate) and the range filter will scan the whole table if there's no index. On years of risk transactions, that hurts. Ask the source DBA for an index. It's a small change that makes a big difference.
  • Empty increments. Add an If Condition on rowsCopied > 0 before the merge step so you don't start Spark sessions for nothing. Mid-month, most claims and reference tables won't have changed, and Spark start-up time uses capacity.

What about Copy job?

Fabric's newer Copy job item can do watermark-based incremental copies for simple cases without building any of this. For a handful of tables going straight into a Lakehouse or Warehouse, try it first. I still build the pipeline version when I need custom logic, logging into my own audit tables (useful when Lloyd's or an auditor asks how a figure was loaded), or a merge step I control. The pattern is old, but it still works well in a lakehouse.