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.

No comments:

Post a Comment