Lakeflow data engineering

azure4 min read

Lakeflow covers data ingestion, transformation, and job orchestration on Azure Databricks. Use this guide to choose a connector, handle source changes, and recover failed runs. The platform team owns access, compute controls, and deployment standards. The data team owns transformation logic, data-quality rules, and the meaning of the published data.

The three responsibilities#

Responsibility Databricks surface Own explicitly
Ingest Lakeflow Connect or custom ingestion Source cursor, schema drift, retries, quarantine, and reconciliation
Transform Lakeflow Spark Declarative Pipelines or jobs Data contracts, expectations, state, backfills, and table ownership
Orchestrate Lakeflow Jobs Dependencies, parameters, identities, notifications, and recovery

Do not collapse all three into one notebook. A source extraction failure, a data-quality failure, and a publish failure need different recovery paths and different evidence.

Default flow#

  1. Land source-faithful data with ingestion metadata and a repeatable cursor.
  2. Quarantine malformed or contract-breaking records instead of silently dropping them.
  3. Standardize types, identifiers, deduplication, and late-arriving behavior in a durable layer.
  4. Publish business tables only after reconciliation and data-quality checks pass.
  5. Make every stage rerunnable for a bounded time range without duplicating results.
  6. Record source counts, target counts, rejected rows, freshness, and the code version for each run.

Managed connector or custom code?#

Prefer a managed connector when it supports the required source objects, authentication, incremental semantics, network path, and recovery behavior. Use custom ingestion when one of those is missing or when the source contract requires transformations before landing. The decision is operational. Choose the path that the team can reconcile and restore during an incident.

Recovery and incremental processing#

Recover from an Auto Loader schema change#

An added source column can stop an Auto Loader stream. Choose the expected response with cloudFiles.schemaEvolutionMode before production. The Auto Loader schema guide defines the available modes:

Mode New column behavior Recovery decision
addNewColumns Updates the saved schema, then stops with UnknownFieldException Configure job retries or restart after the schema update. Test downstream compatibility.
rescue Stores unexpected fields in rescued data without adding schema columns Alert on rescued records so the data team can review each source change.
failOnNewColumns Stops without updating the schema Approve the schema change or remove the unexpected input before restart.
none Does not evolve the schema Configure a rescued-data column if unexpected fields must be retained.

Without an explicit schema, the default is addNewColumns. With an explicit schema, the default is none, and addNewColumns is not allowed. Schema hints do not have that restriction. The rescued-data column retains fields that do not fit the schema. It is separate from the rescue evolution mode. Monitor type mismatches as well as new columns. Auto Loader type widening is Public Preview on Runtime 16.4 or later. Test its supported conversions before you enable it.

Before production, test a new column and an incompatible type with the selected mode. Include a downstream consumer that rejects the new schema in the recovery test.

Distinguish DataFrame and streaming checkpoints#

A DataFrame checkpoint shortens a query plan. A streaming checkpoint records progress and state so a query can recover after a failure. Check which API the code uses before you choose a storage location or investigate permissions.

Checkpoint Unity Catalog volume use
DataFrame Requires Runtime 18.1 or later with standard or dedicated access mode. Serverless does not support it. Standard mode uses spark.checkpoint.dir.
Structured Streaming Set checkpointLocation to a volume path. Give each query a separate location.

See DataFrame checkpoints in volumes and streaming checkpoints. Do not delete a streaming checkpoint as a routine error fix. A new or deleted checkpoint starts the query afresh and can require reprocessing. Check whether the proposed query change supports recovery from the existing checkpoint.

Process daily files on a schedule#

For files that arrive once a day, schedule Auto Loader with AvailableNow. The query processes files present before it starts, then stops. It can use several micro-batches. Files that arrive later wait for the next run. See Auto Loader production guidance.

For a daily file feed:

  1. If the files are encrypted, decrypt them into a governed staging location first.
  2. Schedule incremental ingestion after file preparation succeeds.
  3. Reuse the query's checkpoint across scheduled runs.
  4. Exclude checkpoint storage from lifecycle deletion.
  5. Reconcile source files, accepted records, and quarantined records before you publish the data.

The data team defines decryption and reconciliation rules for the source. The platform team supplies the storage access and compute controls that those steps require.

Production checklist#

  • Service principals or workload identities own production runs; personal identities do not.
  • Schemas and data-quality expectations are versioned with the pipeline.
  • Backfill, replay, and schema-change procedures are tested before cutover.
  • Notifications identify the failed stage and link to useful evidence.
  • Compute policies, cost tags or serverless usage policies, and retention are defined.
  • Deployment uses Declarative Automation Bundles or another reviewed promotion path.

Official sources#

Community examples#