+1 (415) 943-1448

Orchestrating BigQuery: Scheduled Queries, Workflows, Dataform, and Cloud Composer

Almost every BigQuery estate starts the same way: someone clicks Schedule on a query, and it works. Eighteen months later there are forty scheduled queries, three of them depend on each other by being scheduled fifteen minutes apart, and nobody knows which one failed last night. Orchestration is the part of a warehouse that quietly decides whether your data is trustworthy at 8 a.m.

Google Cloud gives you four credible ways to run BigQuery work on a schedule. This tutorial walks through each one, shows the minimum viable implementation, and ends with a decision table you can take to a planning meeting.

The four options at a glance

Scheduled queriesCloud WorkflowsDataformCloud Composer (Airflow)
Lives inBigQuery / Data Transfer ServiceServerless GCP serviceBigQueryManaged Airflow cluster
Dependency modelNone (time only)Explicit steps, you write the graphAutomatic DAG from ref()Explicit DAG in Python
Handles non-BigQuery tasksNoYes (any HTTP/GCP API)NoYes (huge operator library)
BackfillsManual rerunsYou code itBuilt in, date-parameterisedBuilt in (catchup, clear)
Baseline costFree (you pay for the query)Cents per million stepsFree (you pay for the query)Environment runs ~24/7
Ops burdenNear zeroLowLowReal — versions, workers, upgrades

1. Scheduled queries: fine until they are not

A scheduled query is a Data Transfer Service job that runs one SQL statement on a cron-like schedule. It costs nothing beyond the query itself and takes a minute to create:

-- Run from the BigQuery console or bq CLI; @run_date is supplied by the scheduler
CREATE OR REPLACE TABLE `acme.mart.daily_orders`
PARTITION BY order_date AS
SELECT
  DATE(created_at) AS order_date,
  customer_id,
  SUM(amount) AS revenue
FROM `acme.raw.orders`
WHERE DATE(created_at) = @run_date
GROUP BY 1, 2;
bq mk --transfer_config \
  --data_source=scheduled_query \
  --display_name="daily_orders" \
  --target_dataset=mart \
  --schedule="every day 05:30" \
  --params='{"query":"...","destination_table_name_template":"daily_orders","write_disposition":"WRITE_TRUNCATE"}' \
  --service_account_name=bq-scheduler@acme.iam.gserviceaccount.com

Two details people miss. First, always attach a service account (--service_account_name); a scheduled query owned by an employee's user credentials dies the day that employee leaves. Second, scheduled queries support @run_date and @run_time parameters, which is what makes a rerun of a past date produce the right partition instead of today's.

Where they break: there is no dependency graph. If job B needs job A's output, you are encoding that relationship in wall-clock time, and the first time A takes twenty minutes longer than usual, B reads stale data and no one is alerted. Treat more than about ten interdependent scheduled queries as a warning sign.

2. Cloud Workflows: serverless sequencing without a cluster

Cloud Workflows executes a YAML state machine and bills per step, which in practice means a few cents a month for a nightly pipeline. It is the natural upgrade when you need ordering and error handling but do not want an Airflow environment.

main:
  params: [args]
  steps:
    - init:
        assign:
          - project: ${sys.get_env("GOOGLE_CLOUD_PROJECT_ID")}
          - run_date: ${default(map.get(args, "run_date"), text.substring(time.format(sys.now()), 0, 10))}
    - load_staging:
        call: googleapis.bigquery.v2.jobs.query
        args:
          projectId: ${project}
          body:
            useLegacySql: false
            query: ${"CALL acme.sp_load_staging('" + run_date + "')"}
        result: staging
    - build_marts:
        try:
          call: googleapis.bigquery.v2.jobs.query
          args:
            projectId: ${project}
            body:
              useLegacySql: false
              query: ${"CALL acme.sp_build_marts('" + run_date + "')"}
        retry:
          predicate: ${http.default_retry_predicate}
          max_retries: 3
          backoff:
            initial_delay: 30
            max_delay: 300
            multiplier: 2
        except:
          as: e
          steps:
            - alert:
                call: http.post
                args:
                  url: ${sys.get_env("SLACK_WEBHOOK")}
                  body:
                    text: ${"Mart build failed for " + run_date + ": " + json.encode_to_string(e)}
            - rethrow:
                raise: ${e}
    - done:
        return: ${run_date}

Trigger it with Cloud Scheduler, or from an Eventarc rule when a file lands in Cloud Storage — event-driven beats cron whenever the upstream is a file drop. Keep the SQL itself in stored procedures inside BigQuery so the workflow stays a thin controller and your logic stays version-controlled and testable.

3. Dataform: dependencies you never have to declare twice

If your pipeline is mostly SQL transforming tables already in BigQuery, Dataform is usually the right answer. You write SELECT statements; Dataform parses ref() calls, derives the DAG, and creates tables in the correct order.

-- definitions/daily_orders.sqlx
config {
  type: "incremental",
  schema: "mart",
  bigquery: {
    partitionBy: "order_date",
    clusterBy: ["customer_id"],
    updatePartitionFilter: "order_date >= CURRENT_DATE() - 3"
  },
  assertions: {
    uniqueKey: ["order_date", "customer_id"],
    nonNull: ["order_date", "revenue"]
  }
}

SELECT
  DATE(created_at) AS order_date,
  customer_id,
  SUM(amount) AS revenue
FROM ${ref("raw", "orders")}
${when(incremental(), `WHERE DATE(created_at) >= CURRENT_DATE() - 3`)}
GROUP BY 1, 2

A release configuration compiles the project; a workflow configuration runs it on a schedule, with tags so that tags: ["hourly"] and tags: ["daily"] can run on different cadences from one repository. Assertions fail the run when data violates them, which turns "the numbers looked wrong on Tuesday" into an alert instead of a discovery.

Dataform's limit is scope: it orchestrates SQL inside BigQuery and nothing else. It will not call an API, run a Dataflow job, or wait on an SFTP drop. Many teams run Dataform inside a larger orchestrator for exactly that reason — Composer triggers the Dataform workflow invocation and waits for it.

4. Cloud Composer: when the pipeline leaves BigQuery

Cloud Composer is managed Apache Airflow. Reach for it when your pipeline genuinely spans systems: pull from an API, land in GCS, run a Dataflow job, transform in BigQuery, push to a CRM, notify a team.

from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.google.cloud.operators.bigquery import (
    BigQueryInsertJobOperator,
)
from airflow.providers.google.cloud.transfers.gcs_to_bigquery import (
    GCSToBigQueryOperator,
)
from airflow.providers.google.cloud.sensors.gcs import GCSObjectExistenceSensor

default_args = {
    "retries": 3,
    "retry_delay": timedelta(minutes=5),
    "retry_exponential_backoff": True,
}

with DAG(
    dag_id="daily_orders",
    start_date=datetime(2026, 1, 1),
    schedule="30 5 * * *",
    catchup=True,
    max_active_runs=1,
    default_args=default_args,
    tags=["bigquery", "mart"],
) as dag:

    wait_for_drop = GCSObjectExistenceSensor(
        task_id="wait_for_drop",
        bucket="acme-landing",
        object="orders/{{ ds }}/orders.csv",
        mode="reschedule",          # free the worker slot while waiting
        timeout=60 * 60 * 3,
    )

    load = GCSToBigQueryOperator(
        task_id="load_raw",
        bucket="acme-landing",
        source_objects=["orders/{{ ds }}/orders.csv"],
        destination_project_dataset_table="acme.raw.orders${{ ds_nodash }}",
        write_disposition="WRITE_TRUNCATE",
        autodetect=True,
    )

    build_mart = BigQueryInsertJobOperator(
        task_id="build_mart",
        configuration={
            "query": {
                "query": "CALL acme.sp_build_marts('{{ ds }}')",
                "useLegacySql": False,
                "maximumBytesBilled": 2 * 1024 ** 4,  # 2 TiB guardrail
            }
        },
        location="US",
    )

    wait_for_drop >> load >> build_mart

Three habits keep Composer pipelines cheap and sane:

  • Use deferrable operators and mode="reschedule" sensors. A sensor that blocks a worker for three hours is a worker you are paying for and cannot use.
  • Push compute into BigQuery. The operator should submit a job, not process rows in Python on an Airflow worker.
  • Set maximumBytesBilled on every query job. One bad WHERE clause in an automated DAG that runs 365 times a year is a budget event.

The cost to weigh is that a Composer environment runs continuously — a small one is on the order of a few hundred dollars a month before any query cost — and it needs someone to own upgrades. That is fine for a data platform team, and heavy for a three-person analytics function.

Retries, idempotency, and backfills

Whichever tool you choose, orchestration only works if tasks are idempotent: running the same task twice for the same date produces the same result. In BigQuery that means:

  • Write to a specific partition (WRITE_TRUNCATE against table$20260115, or MERGE on a partition filter) instead of appending to the whole table.
  • Parameterise every job by run date, never CURRENT_DATE() inside the SQL — otherwise a rerun of last Tuesday rebuilds today.
  • Make retries safe before you turn retries on. A retry on a non-idempotent append is how you double revenue.

A backfill is then just the same task across a date range: airflow dags backfill in Composer, a date-parameterised release in Dataform, a loop step in Workflows, or a bq mk --transfer_run --run_time per day for scheduled queries.

Alerting on failure — and on silence

Failure alerts are the easy half. Composer and Workflows emit errors to Cloud Logging; create a log-based metric and an alerting policy, and route it to whatever your on-call uses. Scheduled queries can email on failure, which is better than nothing and worse than a pager.

The half teams forget is the silent failure: the job succeeded but produced no rows, or ran an hour late and the dashboard is stale. Catch it with a freshness check that runs after the pipeline and queries your own metadata:

SELECT
  table_name,
  TIMESTAMP_MILLIS(last_modified_time) AS last_modified,
  TIMESTAMP_DIFF(CURRENT_TIMESTAMP(), TIMESTAMP_MILLIS(last_modified_time), HOUR) AS hours_stale,
  row_count
FROM `acme.mart.__TABLES__`
WHERE TIMESTAMP_DIFF(CURRENT_TIMESTAMP(), TIMESTAMP_MILLIS(last_modified_time), HOUR) > 26
   OR row_count = 0;

Any row returned is an incident. Wire that query to an alert and you have covered both "it broke" and "it quietly did nothing".

Choosing

Start at the top of this list and stop at the first line that describes you:

  1. A handful of independent SQL statements on a timer — scheduled queries with a service account and run-date parameters.
  2. SQL transformations with real dependencies between tables — Dataform, with assertions and tags.
  3. A short sequence that touches a couple of non-BigQuery APIs, or is event-driven — Cloud Workflows plus Cloud Scheduler or Eventarc.
  4. Dozens of pipelines spanning several systems, with a team to own them — Cloud Composer, with Dataform invoked from it for the SQL layer.

The expensive mistake is not picking the "wrong" tool; it is staying on scheduled queries three years past the point where dependencies became real, and paying for it in silent data quality failures.

If you would like a second opinion on which tier your pipelines belong in — or help migrating a tangle of scheduled queries into something with a dependency graph — our BigQuery data engineering team does this work weekly. Get in touch with a short description of your current setup.