Skip to content
Reliable Data Engineering
Lesson
Open in the interactive app

Databricks 4: Ingestion, Declarative Pipelines, Jobs and Streaming

This module covers how data gets into Databricks and how pipelines are built and orchestrated. Expect scenario questions such as “ingest 50,000 JSON files per hour with evolving schemas”, “implement SCD2 from a CDC feed” or “how do you rerun only the failed part of a job?”.

Naming note: Delta Live Tables (DLT) was renamed Lakeflow Declarative Pipelines, and its engine was open-sourced into Apache Spark as Spark Declarative Pipelines; Workflows became Lakeflow Jobs; managed connectors are Lakeflow Connect. Interviewers use old and new names interchangeably.


1. File ingestion: Auto Loader vs COPY INTO

Auto Loader (cloudFiles source) incrementally ingests new files from cloud storage as a stream:

(spark.readStream.format("cloudFiles")
    .option("cloudFiles.format", "json")
    .option("cloudFiles.schemaLocation", "/Volumes/prod/raw/_schemas/orders")   # inferred schema + evolution history
    .option("cloudFiles.schemaEvolutionMode", "addNewColumns")                  # or rescue / failOnNewColumns / none
    .option("cloudFiles.inferColumnTypes", "true")
    .load("s3://landing/orders/")
  .writeStream
    .option("checkpointLocation", "/Volumes/prod/raw/_checkpoints/orders")
    .trigger(availableNow=True)                                                 # process the backlog, then stop
    .toTable("prod.bronze.orders"))

COPY INTO is an idempotent SQL command that loads files not loaded before. It’s simple for thousands of files and SQL-first teams, but less scalable than Auto Loader for millions of files and continuous ingestion.

Auto LoaderCOPY INTO
InterfaceStructured Streaming (Python/SQL, pipelines)SQL command
ScaleMillions of files, continuousThousands of files, scheduled
Schema evolutionRich (inference, modes, rescued column)Basic
UseDefault for production file ingestionSimple, ad-hoc or SQL warehouse loads

Lakeflow Connect provides managed connectors for SaaS apps (Salesforce, Workday, ServiceNow…) and databases (SQL Server, Postgres, etc., via CDC), writing into streaming tables with incremental sync, so there’s no custom ingestion code for common sources.


2. Lakeflow Declarative Pipelines (formerly DLT)

You declare what tables should contain; the framework handles dependency ordering, incremental processing, retries, checkpoints, infrastructure and data quality.

import dlt                      # newer API: from pyspark import pipelines as dp  (@dp.table, @dp.materialized_view)
from pyspark.sql import functions as F

@dlt.table(comment="Raw orders from Auto Loader")
def bronze_orders():
    return (spark.readStream.format("cloudFiles")
            .option("cloudFiles.format", "json").load("/Volumes/prod/landing/orders"))

@dlt.table
@dlt.expect_or_drop("valid_amount", "amount >= 0")
@dlt.expect_or_fail("has_key", "order_id IS NOT NULL")
@dlt.expect("recent", "order_ts > '2020-01-01'")            # warn only: recorded in metrics
def silver_orders():
    return (dlt.read_stream("bronze_orders")
            .withColumn("order_date", F.to_date("order_ts")))

@dlt.table                                                     # batch read → materialized view semantics
def gold_daily_revenue():
    return (dlt.read("silver_orders").groupBy("order_date")
            .agg(F.sum("amount").alias("revenue")))

Dataset types:

Expectations (data quality): expect (warn: keep rows, record metrics), expect_or_drop (drop violating rows), expect_or_fail (stop the update). Metrics land in the pipeline event log, a queryable Delta table, for quality dashboards.

CDC and SCD with AUTO CDC (formerly APPLY CHANGES INTO):

dlt.create_streaming_table("silver_customers")

dlt.create_auto_cdc_flow(                  # older name: dlt.apply_changes(...)
    target="silver_customers",
    source="bronze_customer_changes",
    keys=["customer_id"],
    sequence_by=F.col("lsn"),              # ordering column: handles out-of-order events
    apply_as_deletes=F.expr("op = 'D'"),
    except_column_list=["op", "lsn"],
    stored_as_scd_type=2,                  # 1 = overwrite in place, 2 = keep history (__START_AT/__END_AT)
)

It handles out-of-order records by sequence_by, deletes, and SCD1/SCD2 bookkeeping, replacing hundreds of lines of hand-written MERGE logic.

Pipeline settings to know:

When not to use Declarative Pipelines: very custom low-level Spark logic or tuning, complex arbitrary stateful processing, or teams standardised on another orchestration/transformation framework (e.g. dbt + Jobs).


3. Lakeflow Jobs (orchestration)

A job is a DAG of tasks: notebook, Python script/wheel, SQL (queries, dashboards, alerts), dbt, pipeline, JAR, Spark submit, “run job” (call another job), and control-flow tasks.

FeatureWhat it gives you
Task dependenciesDAG with run-if conditions (all succeeded, at least one failed, all done…)
If/else and for-each tasksBranching on task values; looping over a list (e.g. per-country tasks) with bounded concurrency
Task valuesPass small values between tasks (dbutils.jobs.taskValues.set/get)
Job parametersParameterise runs ({{job.start_time}}, custom params) for backfills
TriggersSchedule (cron), file arrival (new files in a location), table update (when upstream UC tables change: data-aware), continuous
Retries & timeoutsPer task, with alerts on failure/duration thresholds
Repair runRe-run only failed and downstream tasks of a failed run, keeping successful ones
Compute per taskJob clusters shared across tasks, serverless, or SQL warehouses
NotificationsEmail, Slack, webhooks; system tables for run history

Interview points:


4. Structured Streaming on Databricks


Interview questions

Auto Loader vs COPY INTO: when do you use each?

Auto Loader for production, high-volume or continuous file ingestion: scalable discovery (incremental listing or file notifications), exactly-once file tracking in a checkpoint, schema inference/evolution and the rescued data column, and both streaming and availableNow batch modes. COPY INTO for simple, idempotent SQL loads of thousands of files, ad-hoc loads, or SQL-warehouse-only teams.

Streaming table vs materialized view in Declarative Pipelines?

A streaming table processes each new input record once (append-only, incremental), ideal for ingestion and transformations of append-only sources. A materialized view represents a query’s full result, refreshed incrementally when possible or fully recomputed otherwise, ideal for aggregations, joins with changing dimensions and anything where inputs are updated or deleted.

How do you implement SCD Type 2 from a CDC feed on Databricks with minimal code?

In Lakeflow Declarative Pipelines, create a streaming table target and use AUTO CDC (formerly APPLY CHANGES) with keys, sequence_by (the CDC ordering column, e.g. LSN), apply_as_deletes for delete ops and stored_as_scd_type=2. It handles out-of-order events, deletes and validity columns automatically. Without pipelines, write a MERGE that closes the current version and inserts the new one, ordering by the sequence column and deduplicating per key per batch.

Expectations: warn vs drop vs fail. How do you choose?

expect (warn) for soft rules where you want metrics but not data loss; expect_or_drop for record-level validity rules where bad rows should be excluded (and ideally captured elsewhere for review); expect_or_fail for invariants whose violation means the data or pipeline is fundamentally wrong (missing keys, impossible states), so the update stops before publishing bad data.

A 12-task job failed at task 9. How do you recover efficiently?

Fix the cause, then use Repair run to re-run only the failed task and its downstream dependents, reusing the successful tasks’ results. Tasks must be idempotent (partition overwrite/MERGE) so re-running is safe. Add retries with backoff for transient failures and alerts on duration so you hear about it before consumers do.

How do you make a Databricks pipeline run only when its upstream data is ready?

Use a table-update trigger on the upstream Unity Catalog tables or a file-arrival trigger on the landing location (or chain jobs with “run job” tasks), instead of fixed schedules. For cross-system dependencies orchestrated elsewhere, have the upstream publish a completion signal that triggers the job via the API.