Search⌘KHide sidebarOpen menu
Switch to dark mode
Copy page content as Markdown⌘⌥C
Login
Pipelines

Quarantine bad rows

Splitting rows that fail a check into their own dataset.

This recipe sorts every row into one of 2 datasets. Valid rows go to a clean dataset that reports and dashboards read. Rows that fail a check, such as a blank customer ID or a date nobody can parse, go to a rejected dataset with a reject_reason column, where you can count them and trace them back to the source.

When you need this

Use it when a source you do not control, such as uploaded files or a partner's API, sometimes sends incomplete or malformed rows. A WHERE clause that quietly drops bad rows hides the problem, and a cast that fails stops the whole run, so the good rows wait on the bad ones.

How the recipe fits together

A pipeline writes exactly one destination, so the recipe uses 2 pipelines over the same source. Both share their first 2 stages and differ only in the last.

PipelineStagesDestination
Cleantyped → checked → cleanorders_clean, Silver or Gold
Rejectedtyped → checked → rejectedorders_rejected

The examples read a source with the alias orders_raw whose columns are all text. Build the stages in the pipeline builder and preview each one before adding the next.

Parse the raw columns

The first stage converts text to real types. TRY_CAST returns NULL when a value cannot be converted, where CAST would fail the run. Keep the raw text next to the parsed value so the rejected dataset shows what actually arrived.

-- stage: typed
SELECT
  order_id,
  customer_id,
  customer_email,
  currency,
  order_date AS order_date_raw,
  amount AS amount_raw,
  TRY_CAST(order_date AS DATE) AS order_date,
  TRY_CAST(amount AS DECIMAL(12, 2)) AS amount
FROM orders_raw

TRY_CAST(... AS DATE) reads ISO dates such as 2026-03-01. The pitfalls below show how to handle other layouts.

Assign a reject reason

The second stage adds reject_reason, which stays NULL for a valid row. CASE returns the first branch that matches, so put the checks in the order you want them reported.

-- stage: checked
SELECT
  *,
  CASE
    WHEN order_id IS NULL OR TRIM(CAST(order_id AS VARCHAR)) = '' THEN 'missing order_id'
    WHEN customer_id IS NULL THEN 'missing customer_id'
    WHEN order_date_raw IS NULL THEN 'missing order_date'
    WHEN order_date IS NULL THEN 'unparseable order_date'
    WHEN order_date > CURRENT_DATE THEN 'order_date in the future'
    WHEN amount IS NULL THEN 'missing or unparseable amount'
    WHEN amount <= 0 OR amount > 1000000 THEN 'amount out of range'
    WHEN currency IS NULL OR currency NOT IN ('EUR', 'USD', 'GBP') THEN 'unknown currency'
    WHEN customer_email IS NOT NULL
      AND NOT regexp_like(customer_email, '^[^@ ]+@[^@ ]+[.][^@ ]+$') THEN 'malformed customer_email'
  END AS reject_reason
FROM typed

Limits such as the maximum amount can come from a pipeline variable instead of a literal (see Template variables). If you use one, give it the same value in both pipelines.

Finish the clean pipeline

The last stage keeps the valid rows and only the typed columns.

-- stage: clean
SELECT order_id, customer_id, customer_email, currency, order_date, amount
FROM checked
WHERE reject_reason IS NULL

Build the rejected pipeline

Create a second pipeline with the same source, copy the typed and checked stages unchanged, and end it with the stage below. You can also ask the assistant to build it from the first pipeline.

-- stage: rejected
SELECT
  order_id,
  customer_id,
  customer_email,
  currency,
  order_date_raw,
  amount_raw,
  reject_reason,
  now() AS checked_at
FROM checked
WHERE reject_reason IS NOT NULL

Give each pipeline its own destination name. Turn on Trigger on update for the source in both pipelines so that each new delivery runs both. They run independently, and one can finish a little before the other.

Choose write modes

Give both destinations the same write mode, matched to how the source arrives. See Write modes and partitioning.

Source arrives asWrite mode on bothEffect
A full replacementOverwrite, full tableBoth datasets show the current state. A row fixed at the source leaves orders_rejected on the next run.
One partition per runOverwrite, predicate such as dt = '{{ run.dt }}'A rerun replaces that partition in both. Filter typed to the same partition.

Append on the rejected pipeline keeps a history of rejections, but a rerun of the same data adds the same rejects again.

Check the result in Query

Count the rejects by reason:

SELECT reject_reason, count(*) AS row_count
FROM orders_rejected
GROUP BY reject_reason
ORDER BY row_count DESC

With full-table overwrites, every source row should land in exactly one of the 2 datasets:

SELECT
  (SELECT count(*) FROM orders_raw) AS source_rows,
  (SELECT count(*) FROM orders_clean) AS clean_rows,
  (SELECT count(*) FROM orders_rejected) AS rejected_rows

If clean_rows + rejected_rows differs from source_rows, either the 2 pipelines no longer share the same checked stage or one of them has not run since the last delivery.

Watch the rejected dataset

An automation can check orders_rejected on a schedule, with a goal such as "Report needs attention if any row in orders_rejected has a checked_at in the last 24 hours, grouped by reject_reason." Beetl does not send notifications, so make the Automations list part of someone's routine.

A one-pipeline variant

To keep the checks in one place, write the output of checked to a Silver dataset such as orders_checked. A second pipeline reads it and writes Gold with WHERE reject_reason IS NULL, and you query orders_checked to inspect rejects. You trade the duplicated stages for one more dataset and run.

Pitfalls

  • currency NOT IN (...) is neither true nor false when currency is NULL, so on its own that branch lets the row through. Test IS NULL explicitly, as the example does.
  • TRIM rejects integer columns. Cast to VARCHAR first, as in the order_id check.
  • to_date and to_timestamp fail the run on a single bad value. For dates such as 31.12.2026, rewrite the text to ISO and use TRY_CAST, which also returns NULL for impossible dates like 31.02.2026: TRY_CAST(regexp_replace(order_date, '^([0-9]{2})[.]([0-9]{2})[.]([0-9]{4})$', '\3-\2-\1') AS DATE).
  • The example records one reason per row. To record every failed check, join several CASE expressions with concat_ws, which skips NULL values. It returns an empty string when nothing failed, so wrap it: NULLIF(concat_ws('; ', CASE WHEN ... END, CASE WHEN ... END), '').
  • The 2 pipelines can drift apart. When you change a check, change it in both and deploy both; the count query above catches a mismatch.

Last updated on

On this page