- Mark as New
- Bookmark
- Subscribe
- Mute
- Subscribe to RSS Feed
- Permalink
- Report Inappropriate Content
03-26-2026 10:04 AM
Great question — and since you already have DABs and numbered SQL files, you're most of the way there. You do not need Alembic or SQLAlchemy. Here's a concrete implementation of the migration runner pattern that plugs directly into your existing DABs setup.
The Pattern
The idea is simple:
- Keep your numbered SQL migration files as-is (001, 002, etc.)
- Add a migration history table per environment to track what's been applied
- Add a single migration runner task in your DABs bundle that runs all unapplied migrations in order
- Each migration runs exactly once — no editing old files, new changes go in new files
Step 1: Migration History Table
This gets created automatically by the runner, but here's what it looks like:
CREATE TABLE IF NOT EXISTS ${catalog}.admin.schema_migrations (
version STRING NOT NULL,
file_name STRING,
applied_at TIMESTAMP DEFAULT current_timestamp(),
checksum STRING
);
Step 2: Migration Runner (Python Task)
Create a file migrations/run_migrations.py in your DABs project:
import os
import hashlib
from pyspark.sql import SparkSession
spark = SparkSession.builder.getOrCreate()
# These come from DABs variable overrides per environment
CATALOG = spark.conf.get("spark.databricks.migration.catalog")
MIGRATIONS_DIR = spark.conf.get("spark.databricks.migration.dir", "/Workspace/migrations/sql")
def get_applied_versions():
"""Read which migrations have already been applied."""
spark.sql(f"""
CREATE TABLE IF NOT EXISTS {CATALOG}.admin.schema_migrations (
version STRING NOT NULL,
file_name STRING,
applied_at TIMESTAMP,
checksum STRING
)
""")
rows = spark.sql(
f"SELECT version FROM {CATALOG}.admin.schema_migrations"
).collect()
return {row.version for row in rows}
def get_pending_migrations(applied):
"""Find SQL files that haven't been applied yet, sorted by version prefix."""
files = []
for f in sorted(os.listdir(MIGRATIONS_DIR)):
if not f.endswith(".sql"):
continue
version = f.split("_")[0] # e.g. "001" from "001_create_base_schemas.sql"
if version not in applied:
files.append((version, f))
return files
def run_migration(version, file_name):
"""Execute a single migration file and record it."""
path = os.path.join(MIGRATIONS_DIR, file_name)
with open(path, "r") as fh:
sql_content = fh.read()
checksum = hashlib.md5(sql_content.encode()).hexdigest()
# Split on semicolons to handle multi-statement files
statements = [s.strip() for s in sql_content.split(";") if s.strip()]
for stmt in statements:
# Replace ${catalog} placeholder with actual catalog
resolved = stmt.replace("${catalog}", CATALOG)
print(f" Executing: {resolved[:80]}...")
spark.sql(resolved)
# Record successful migration
spark.sql(f"""
INSERT INTO {CATALOG}.admin.schema_migrations
VALUES ('{version}', '{file_name}', current_timestamp(), '{checksum}')
""")
print(f" Recorded migration {version}: {file_name}")
def main():
applied = get_applied_versions()
print(f"Already applied: {sorted(applied)}")
pending = get_pending_migrations(applied)
if not pending:
print("No new migrations to apply.")
return
print(f"Applying {len(pending)} migration(s)...")
for version, file_name in pending:
print(f"\n--- Migration {version}: {file_name} ---")
run_migration(version, file_name)
print("\nAll migrations applied successfully.")
main()
Step 3: DABs Bundle Configuration
In your databricks.yml, add the migration runner as a job with environment-specific catalog overrides:
variables:
catalog:
default: dev_catalog
resources:
jobs:
schema_migrations:
name: "schema-migrations-${bundle.environment}"
tasks:
- task_key: run_migrations
existing_cluster_id: ${var.cluster_id}
spark_python_task:
python_file: ./migrations/run_migrations.py
parameters: []
spark_conf:
spark.databricks.migration.catalog: ${var.catalog}
spark.databricks.migration.dir: /Workspace/${workspace.root_path}/migrations/sql
environments:
dev:
variables:
catalog: dev_catalog
test:
variables:
catalog: test_catalog
prod:
variables:
catalog: prod_catalog
Step 4: Your SQL Migration Files
Keep them exactly as you have them — numbered, one per change, never edited after creation:
migrations/sql/
001_create_base_schemas.sql
002_create_orders_table.sql
003_seed_reference_data.sql
004_add_status_column.sql ← new changes = new file
Example migration file (001_create_base_schemas.sql😞
CREATE SCHEMA IF NOT EXISTS ${catalog}.analytics;
CREATE SCHEMA IF NOT EXISTS ${catalog}.admin;
Example seed file (003_seed_reference_data.sql) — use MERGE for idempotency:
MERGE INTO ${catalog}.analytics.order_status AS target
USING (
SELECT * FROM VALUES
('NEW', 'New Order'),
('SHIPPED', 'Order Shipped'),
('DELIVERED', 'Order Delivered')
AS source(code, description)
) AS source
ON target.code = source.code
WHEN NOT MATCHED THEN INSERT *
How It Works in Practice
Adding a new migration:
- Create 005_add_customer_email.sql in migrations/sql/
- Commit and push
- databricks bundle deploy -e test → runs the migration job → runner sees 005 is not in history → applies it
- databricks bundle deploy -e prod → same thing for prod
The runner is safe to re-run — it always checks the history table first. Already-applied migrations are skipped.
To Answer Your Specific Questions
- Version = the prefix (001, 002, etc.). Exactly right — never edit old migrations, always add new files.
- One migration job, not one task per file — the single Python runner task handles all files. No need to edit YAML when adding migrations.
- Version check is in the runner, not in each SQL file — the runner reads the history table once, then only executes files whose version prefix isn't recorded yet.
- No need for Alembic — this pattern gives you the same ordered, idempotent, environment-aware migrations without adding Python ORM complexity. Your migrations stay as plain SQL, which is easier for the whole team to work with.
Solutions Engineer @ Databricks