Automated Data Quality in BigQuery: A Guide to Dataform Assertions
Flawed data in BI dashboards leads to incorrect management decisions and undermines business trust in the engineering team. Discovering a data defect at the dashboard level means the incorrect data has already passed through the entire infrastructure, consuming BigQuery compute resources along the way.
The solution is the shift-left data quality concept. Checks (assertions) are integrated directly into the data transformation Directed Acyclic Graph (DAG). If a critical check fails, the pipeline halts, preventing corrupted data from reaching the final data marts.
The technical stack for this approach includes:
- Dataform executes SQL transformations and provides a built-in assertion mechanism. If a check fails, dependent downstream nodes are blocked.
- BigQuery acts as the compute engine processing these checks.
- Cloud Run and Cloud Scheduler orchestrate the Dataform API triggers, poll for execution status, and route failure notifications to corporate messengers.
8 Real-World Data Quality Cases
Below are frequent anomalies in data warehouses, methods to intercept them using Dataform assertions, standard implementation mistakes, and the correct engineering approaches. In Dataform, an assertion is a SQL query that returns rows violating a rule. If the query returns zero rows, the test passes.
Case 1: Duplicates from Source Retries
Web analytics events or payment system webhooks often duplicate due to network latency and at-least-once delivery mechanisms.
- The Problem: Duplicate
event_idrecords inflate core metrics. - Standard Mistake: Using nested subqueries for deduplication (e.g.,
SELECT * FROM (SELECT ..., ROW_NUMBER() OVER(...) as rn) WHERE rn = 1). This adds unnecessary nesting and complicates code readability. - Correct Approach: Utilize BigQuery’s native
QUALIFYclause with window functions to streamline the logic within a singleSELECTstatement.
SQL
-- definitions/staging/stg_events.sqlx
config {
type: "table",
assertions: {
uniqueKey: ["event_id"]
}
}
SELECT
event_id,
user_id,
event_timestamp
FROM ${ref("raw_events")}
QUALIFY ROW_NUMBER() OVER(PARTITION BY event_id ORDER BY event_timestamp DESC) = 1
Case 2: Referential Integrity Violations
Transactions appear with a user_id that does not exist in the primary users table, often due to CRM synchronization failures.
- The Problem: Downstream JOINs drop revenue records because the user dimension is missing.
- Standard Mistake: Writing an assertion with a direct INNER JOIN while ignoring soft-deletes. This generates false positive alerts for legitimately deleted user accounts.
- Correct Approach: Write a custom assertion that accounts for deletion flags.
SQL
-- definitions/assertions/assert_orders_users_fk.sqlx
config { type: "assertion" }
SELECT o.order_id, o.user_id
FROM ${ref("stg_orders")} o
LEFT JOIN ${ref("stg_users")} u
ON o.user_id = u.user_id AND u.is_deleted = FALSE
WHERE u.user_id IS NULL
Case 3: Stale Data
The pipeline executes without technical errors, but a third-party API returned an empty array or data from the previous week.
- The Problem: BI reports show zero revenue for yesterday, triggering business panic.
- Standard Mistake: Hardcoding local timezones and ignoring weekend schedules when source systems might legitimately not update.
- Correct Approach: Use UTC timestamps and incorporate business day logic.
SQL
-- definitions/assertions/assert_freshness_sales.sqlx
config { type: "assertion" }
SELECT 1
FROM ${ref("stg_sales")}
HAVING MAX(created_at) < TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 24 HOUR)
-- Apply filters for business days if the source system pauses on weekends
Case 4: Missing Critical Values (NULLs)
The amount or currency field arrives empty due to JSON parsing errors during the extraction phase.
- The Problem: Financial aggregations fail or return skewed results.
- Standard Mistake: Applying Dataform’s built-in
nonNullcheck to every column in the table. This blocks the pipeline when optional parameters (like apromo_code) are naturally missing. - Correct Approach: Restrict
nonNullassertions exclusively to Primary Keys and critical business metrics.
JavaScript
// definitions/staging/stg_transactions.sqlx
config {
type: "table",
assertions: {
nonNull: ["transaction_id", "amount", "currency"]
}
}
Case 5: Anomalous Volume Drops
The daily ingestion volume drops by 50%, but table schemas and data types remain intact.
- The Problem: Partial data loss goes unnoticed because the pipeline technically succeeds.
- Standard Mistake: Setting static row count thresholds (e.g.,
< 1000 rows). These require constant manual adjustment as the business scales. - Correct Approach: Dynamically calculate the baseline using rolling averages.
SQL
-- definitions/assertions/assert_volume_drop.sqlx
config { type: "assertion" }
WITH daily_counts AS (
SELECT DATE(event_timestamp) as dt, COUNT(*) as cnt
FROM ${ref("stg_events")}
GROUP BY 1
),
stats AS (
SELECT
cnt as current_cnt,
AVG(cnt) OVER(ORDER BY dt ROWS BETWEEN 7 PRECEDING AND 1 PRECEDING) as avg_7d
FROM daily_counts
WHERE dt = CURRENT_DATE()
)
SELECT * FROM stats WHERE current_cnt < (avg_7d * 0.5)
Case 6: Unexpected Categorical Values
The frontend team introduces a new order status like partially_refunded, breaking hardcoded SQL logic in downstream marts.
- The Problem: Unknown statuses are mapped to NULL or cause case statement failures.
- Standard Mistake: Using Dataform’s
acceptedValuesarray for dynamic business entities. Every new valid category will require an engineering commit to update the assertion. - Correct Approach: Use
acceptedValuesonly for strictly static flags. For dynamic values, validate against a centralized dictionary table.
JavaScript
// For strictly static flags only
config {
assertions: {
acceptedValues: {
transaction_type: ["purchase", "refund"]
}
}
}
Case 7: Logical Date Conflicts
A subscription’s end_date is recorded chronologically earlier than its start_date.
- The Problem: Calculating Monthly Recurring Revenue (MRR) or user lifetime yields negative intervals.
- Standard Mistake: Comparing the dates directly without accounting for active subscriptions where
end_dateis naturally NULL. - Correct Approach: Explicitly handle NULL values in the validation logic.
SQL
-- definitions/assertions/assert_dates_logic.sqlx
config { type: "assertion" }
SELECT subscription_id
FROM ${ref("stg_subscriptions")}
WHERE end_date IS NOT NULL AND end_date < start_date
Case 8: Impossible Negative Values
Item quantities or transaction amounts drop below zero.
- The Problem: Negative values skew total revenue calculations.
- Standard Mistake: Applying a blanket
>= 0rule to all amount columns, forgetting that tables often containrefundorchargebackevents where negative values are required. - Correct Approach: Bind the numeric condition strictly to the relevant transaction type.
SQL
-- definitions/assertions/assert_positive_revenue.sqlx
config { type: "assertion" }
SELECT transaction_id, amount, transaction_type
FROM ${ref("stg_transactions")}
WHERE transaction_type = 'purchase' AND amount <= 0
Implementing the Pipeline from Scratch
Deploying this shift-left architecture requires configuring the Dataform repository, defining the dependency graph, and setting up an external orchestrator to handle failures.
Step 1: Dataform Repository Configuration
Initialize a Dataform repository in the Google Cloud Console, connect it to your version control system, and define the target datasets in dataform.json. Dedicate a separate schema for assertions to keep the warehouse clean.
JSON
{
"defaultSchema": "data_warehouse",
"assertionSchema": "data_quality_logs",
"defaultDatabase": "your-gcp-project-id"
}
Step 2: Defining the Graph and Blocking Downstream Execution
Define the staging models with built-in assertions. Then, explicitly configure the downstream data marts to depend on those specific assertions. If the uniqueKey test on stg_users fails, Dataform will cancel the execution of mart_active_users, protecting the BI layer.
SQL
-- definitions/staging/stg_users.sqlx
config {
type: "table",
assertions: {
uniqueKey: ["user_id"]
}
}
SELECT user_id, email FROM ${ref("raw_users")}
-- definitions/marts/mart_active_users.sqlx
config {
type: "table",
dependencies: ["stg_users_assertions_uniqueKey_0"]
}
SELECT * FROM ${ref("stg_users")} WHERE is_active = true
Step 3: Orchestration and Alerting Setup
Relying on scheduled executions directly within Dataform limits alerting capabilities. Instead, deploy a Cloud Run service that triggers the Dataform API, monitors the execution state, and dispatches webhook notifications if an assertion fails. Cloud Scheduler acts as the cron trigger for this service.
Python implementation for the Cloud Run container:
Python
from google.cloud import dataform_v1beta1
import time
import requests
def run_dataform_pipeline():
client = dataform_v1beta1.DataformClient()
project_id = "your-gcp-project-id"
location = "europe-west1"
repo = "dataform-repo"
parent = f"projects/{project_id}/locations/{location}/repositories/{repo}"
# Trigger the workflow
request = dataform_v1beta1.CreateWorkflowInvocationRequest(
parent=parent,
workflow_invocation=dataform_v1beta1.WorkflowInvocation(
compilation_result=f"{parent}/compilationResults/latest"
)
)
invocation = client.create_workflow_invocation(request=request)
invocation_name = invocation.name
# Poll execution status
while True:
status_req = dataform_v1beta1.GetWorkflowInvocationRequest(name=invocation_name)
status = client.get_workflow_invocation(request=status_req)
state = status.state.name
if state == "SUCCEEDED":
break
elif state in ["FAILED", "CANCELLED"]:
send_slack_alert(invocation_name)
break
time.sleep(30)
def send_slack_alert(invocation_name):
webhook_url = "https://hooks.slack.com/services/T000/B000/XXX"
message = {
"text": f"Data Quality Check Failed. BI updates blocked. Invocation: {invocation_name}"
}
requests.post(webhook_url, json=message)
Configure Cloud Scheduler to send an HTTP GET request to this Cloud Run endpoint daily, prior to core business hours. This ensures that any data anomalies are caught, pipeline execution is halted, and the engineering team is notified before business users access their reports.
