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.
| Pipeline | Stages | Destination |
|---|---|---|
| Clean | typed → checked → clean | orders_clean, Silver or Gold |
| Rejected | typed → checked → rejected | orders_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_rawTRY_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 typedLimits 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 NULLBuild 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 NULLGive 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 as | Write mode on both | Effect |
|---|---|---|
| A full replacement | Overwrite, full table | Both datasets show the current state. A row fixed at the source leaves orders_rejected on the next run. |
| One partition per run | Overwrite, 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 DESCWith 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_rowsIf 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 whencurrencyisNULL, so on its own that branch lets the row through. TestIS NULLexplicitly, as the example does.TRIMrejects integer columns. Cast toVARCHARfirst, as in theorder_idcheck.to_dateandto_timestampfail the run on a single bad value. For dates such as31.12.2026, rewrite the text to ISO and useTRY_CAST, which also returnsNULLfor impossible dates like31.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
CASEexpressions withconcat_ws, which skipsNULLvalues. 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