Idempotent Data Pipelines: Handling Partial Failures with dbt and Airbyte

Business Context and the Importance of Idempotency

In distributed data systems, failures are a mathematical certainty. APIs drop connections, external databases experience lockups, and compute nodes hit memory limits. When a pipeline fails midway—a partial failure—it leaves the data warehouse in an unknown state.

If a pipeline lacks idempotency, resolving a partial failure requires manual intervention. Engineers must identify which records were written, execute custom deletion scripts to remove partial data, and clean up intermediate tables before restarting the process. This manual rollback delays reporting, consumes engineering hours, and introduces the risk of human error, such as accidental data deletion.

Idempotency is a property of an operation that guarantees the same final state regardless of whether it is executed once or multiple times. In data engineering, an idempotent pipeline ensures that a task can be safely retried from any point of failure without causing data duplication or corruption. The operational standard becomes “restart with confidence,” enabling automated retries and predictable SLA fulfillment.

8 Real-World Cases of Partial Failures

Below are common scenarios where data pipelines break idempotency, the standard mistakes made during implementation, and the correct architectural approaches to solve them using the modern data stack (BigQuery, dbt).

Case 1: Retrying Failed Ingestions and Duplication

A pipeline extracts data from a CRM system and writes it to a BigQuery table. A network timeout causes the task to fail after inserting 50% of the rows. The orchestrator automatically retries the task.

  • Standard Mistake: Executing a naive INSERT INTO target SELECT * FROM source. On retry, the initial 50% of rows are inserted again, duplicating data.
  • Correct Approach: Use the MERGE statement (UPSERT) based on a primary key, or utilize dbt’s incremental materialization with a defined unique_key. If the record exists, it updates; if not, it inserts.

Case 2: API Pagination Failures

An extraction script pulls 10,000 records from an API that returns 100 records per page. The connection drops on page 95.

  • Standard Mistake: Storing data in memory during extraction and attempting a single batch insert at the end. The failure destroys the progress of the first 94 pages, forcing a complete restart.
  • Correct Approach: Manage extraction state using cursors. Load data incrementally and store the cursor value (e.g., updated_at timestamp) after each successful page. Upon retry, extraction resumes from page 95. Modern ELT tools handle this state management natively.

Case 3: Streaming At-Least-Once Delivery

Real-time events flow from application backends through Pub/Sub into BigQuery. Pub/Sub guarantees “at-least-once” delivery, meaning network retries inherently introduce duplicate messages into the raw tables.

  • Standard Mistake: Relying on the downstream BI tool to filter duplicates, which slows down dashboard rendering, or using complex nested subqueries (GROUP BY and max timestamps) to clean the data during transformation.
  • Correct Approach: Accept raw duplicates in the extraction layer, but enforce idempotency in the staging layer using BigQuery’s native window functions to isolate the latest event.

Case 4: Historical Backfills

A change in business logic requires recalculating marketing attribution for the past six months.

  • Standard Mistake: Running massive UPDATE statements directly on production tables. If the query times out halfway, the table contains a mix of old and new logic, making the state inconsistent and difficult to audit.
  • Correct Approach: Treat data as immutable. Calculate the historical data into a temporary table, verify the output, and then execute an atomic table swap or partition overwrite (WRITE_TRUNCATE).

Case 5: Generating Surrogate Keys

A transformation pipeline generates unique IDs for new dimensional records, such as newly registered product categories.

  • Standard Mistake: Using database sequence generators, AUTO_INCREMENT, or GENERATE_UUID(). If the pipeline is executed twice, GENERATE_UUID() creates entirely new IDs for the exact same records, breaking historical joins in the fact tables.
  • Correct Approach: Use deterministic hashing based on the natural keys of the record.SQL-- The hash remains identical across multiple runs for the same input SELECT TO_HEX(MD5(CONCAT(LOWER(product_name), '|', LOWER(category)))) AS product_key

Case 6: Mixed DML Operations (Side-effects)

A script executes a SQL transformation and subsequently triggers an API call to send an email notification to the analytics team.

  • Standard Mistake: Coupling data mutations with external side-effects in the same block of code. If the email API times out, the task fails. Upon retry, the SQL executes again, wasting BigQuery compute resources.
  • Correct Approach: Decouple transformations from side-effects. SQL execution should reside in one isolated node, and notifications should be handled by the orchestrator (e.g., Airflow/Prefect on_success_callback).

Case 7: File Path Collisions in ELT

A process downloads a large CSV export from an SFTP server and uploads it to Google Cloud Storage (GCS) before loading it into BigQuery.

  • Standard Mistake: Appending execution timestamps to filenames (e.g., sales_20261008_143000.csv). Three retries result in three distinct files. The downstream load task processes all three, tripling the data.
  • Correct Approach: Use static, deterministic filenames based on the logical execution date (sales_20261008.csv). A retry simply overwrites the existing file in GCS.

Case 8: Late Arriving Dimensions

A fact table receives a transaction record, but the corresponding user data (dimension) has not yet been processed due to an upstream failure.

  • Standard Mistake: Using an INNER JOIN during mart construction, which permanently drops the transaction record. If the dimension arrives later, the pipeline must be fully recomputed to recover the dropped fact.
  • Correct Approach: Use a LEFT JOIN and assign a default “Unknown” surrogate key (e.g., -1) for missing dimensions. When the dimension eventually arrives, the idempotent pipeline updates the fact table with the correct key during the next run.

Implementing an Idempotent Pipeline from Scratch

To demonstrate the correct architectural approach, we will build an extraction, loading, and transformation pipeline using Airbyte, BigQuery, and dbt. This stack shifts the responsibility of state management and idempotency from custom Python scripts to specialized serverless tools.

Step 1: Extraction and Loading (Airbyte)

Writing custom Python scripts to pull data from APIs often leads to poor pagination handling and lack of cursor state storage. Airbyte solves this natively.

When configuring an Airbyte connection from a source (e.g., Stripe, Shopify) to BigQuery:

  1. Sync Mode: Select Incremental - Append.
  2. Cursor Field: Define the cursor field (e.g., updated_at).

How it handles partial failures: Airbyte maintains a state file. If a sync fails after extracting 50,000 out of 100,000 records, Airbyte records the updated_at timestamp of the last successful batch. Upon retry, it queries the source API starting strictly from that timestamp. Data already loaded into the BigQuery raw schema (_airbyte_raw_tables) remains intact, and no API quotas are wasted redownloading the first 50,000 records.

Step 2: Transformation and Deduplication (dbt + BigQuery)

The data in the raw schema may contain duplicates (due to overlapping cursors or source-level updates). We use dbt to transform this data into a clean staging layer idempotently.

We configure the dbt model using the incremental materialization and specify a unique_key. Under the hood, dbt compiles this into a BigQuery MERGE statement.

SQL

-- models/staging/stg_payments.sql

{{ config(
    materialized='incremental',
    unique_key='payment_id',
    on_schema_change='fail'
) }}

WITH source_data AS (
    SELECT 
        JSON_EXTRACT_SCALAR(_airbyte_data, '$.id') AS payment_id,
        CAST(JSON_EXTRACT_SCALAR(_airbyte_data, '$.amount') AS NUMERIC) AS amount,
        CAST(JSON_EXTRACT_SCALAR(_airbyte_data, '$.updated_at') AS TIMESTAMP) AS updated_at
    FROM {{ source('raw_data', '_airbyte_raw_payments') }}
    
    {% if is_incremental() %}
        -- Only process rows extracted since the last successful dbt run
        WHERE CAST(JSON_EXTRACT_SCALAR(_airbyte_data, '$.updated_at') AS TIMESTAMP) > (SELECT MAX(updated_at) FROM {{ this }})
    {% endif %}
),

deduplicated_data AS (
    SELECT *
    FROM source_data
    -- BigQuery native function to handle within-batch duplicates
    QUALIFY ROW_NUMBER() OVER(PARTITION BY payment_id ORDER BY updated_at DESC) = 1
)

SELECT * FROM deduplicated_data

Execution Flow and Result

  1. Run 1 (Success): Airbyte loads records 1-100. dbt reads records 1-100, deduplicates them, and INSERTs them into stg_payments.
  2. Run 2 (Partial Failure): Airbyte attempts to load records 101-200. It fails at 150. Records 101-150 are in the raw table. dbt does not run.
  3. Run 3 (Retry): Airbyte resumes from record 150 and successfully loads 150-200.
  4. Run 4 (dbt execution): dbt queries the raw table for all records where updated_at is greater than the maximum updated_at currently in stg_payments (which covers records 101-200).
  5. Idempotent resolution: The QUALIFY clause ensures that if record 150 was pulled twice during the Airbyte retry, only the latest version is kept. dbt generates a MERGE statement based on payment_id.

If the engineer manually triggers the dbt run five times in a row without new data, the table state remains completely unaltered. The pipeline is mathematically idempotent, fully automated, and resilient to partial failures.

Similar Posts