Skip to content

MCP Approach

Skill: databricks-spark-declarative-pipelines

Skip Asset Bundles entirely when you just want to test pipeline logic. Your AI coding assistant uploads your local files, calls manage_pipeline to create the pipeline, runs it, and verifies the output tables — all in one conversation. This is the fastest path from idea to a running pipeline. Use it for prototyping, demos, and quick experiments; reach for Asset Bundles only when you need multi-environment deployment or CI/CD.

“Create a serverless pipeline called my_orders_pipeline from local SQL files in /tmp/my_pipeline, run it with full refresh, and validate the output tables”

# Step 1: Upload pipeline files to the workspace
manage_workspace_files(
action="upload",
local_path="/tmp/my_pipeline",
workspace_path="/Workspace/Users/user@example.com/my_pipeline"
)
# Step 2: Create the pipeline and run it in one call
result = manage_pipeline(
action="create_or_update",
name="my_orders_pipeline",
root_path="/Workspace/Users/user@example.com/my_pipeline",
catalog="my_catalog",
schema="my_schema",
workspace_file_paths=[
"/Workspace/Users/user@example.com/my_pipeline/bronze/ingest_orders.sql",
"/Workspace/Users/user@example.com/my_pipeline/silver/clean_orders.sql",
"/Workspace/Users/user@example.com/my_pipeline/gold/daily_summary.sql"
],
start_run=True,
wait_for_completion=True,
full_refresh=True
)
# Step 3: Validate output data
if result["success"]:
stats = get_table_stats_and_schema(
catalog="my_catalog",
schema="my_schema",
table_names=["bronze_*", "silver_*", "gold_*"]
)

Key decisions:

  • manage_pipeline(action="create_or_update") is idempotent — it creates the pipeline if missing, updates it in place if it exists. Using the same name twice never duplicates.
  • get_table_stats_and_schema over manual SQL counts — returns schema, row counts, and column stats in one call. Don’t use execute_sql with COUNT queries for validation.
  • workspace_file_paths lists every file — unlike Asset Bundles’ glob patterns, MCP requires explicit paths. Add or remove files by updating this list.
  • full_refresh=True reprocesses everything — for first runs and schema changes. Skip for incremental runs.

“Set up a local folder structure with SQL and Python files before uploading to the workspace”

my_pipeline/
├── bronze/
│ ├── ingest_orders.sql
│ └── ingest_events.py
├── silver/
│ └── clean_orders.sql
└── gold/
└── daily_summary.sql

SQL file (bronze/ingest_orders.sql):

CREATE OR REFRESH STREAMING TABLE bronze_orders
CLUSTER BY (order_date)
AS
SELECT
*,
current_timestamp() AS _ingested_at,
_metadata.file_path AS _source_file
FROM STREAM read_files(
'/Volumes/catalog/schema/raw/orders/',
format => 'json',
schemaHints => 'order_id STRING, customer_id STRING, amount DECIMAL(10,2), order_date DATE'
);

Python file (bronze/ingest_events.py):

from pyspark import pipelines as dp
from pyspark.sql.functions import col, current_timestamp
@dp.table(name="bronze_events", cluster_by=["event_date"])
def bronze_events():
return (
spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("/Volumes/catalog/schema/raw/events/")
.withColumn("_ingested_at", current_timestamp())
.withColumn("_source_file", col("_metadata.file_path"))
)

Edit and test files locally, then upload the entire folder in one call. The folder structure does not affect the pipeline — only the explicit workspace_file_paths list matters.

“The pipeline failed. Find the error, then re-upload my fixed file and re-run.”

# Get detailed pipeline state and recent error events
manage_pipeline(action="get", pipeline_id=result["pipeline_id"])
# Or get events directly with severity filter
manage_pipeline_run(
action="get_events",
pipeline_id=result["pipeline_id"],
event_log_level="ERROR",
max_results=10
)
# Fix the file locally, then re-upload and re-run
manage_workspace_files(
action="upload",
local_path="/tmp/my_pipeline/silver/clean_orders.sql",
workspace_path="/Workspace/Users/user@example.com/my_pipeline/silver/clean_orders.sql"
)
manage_pipeline(
action="create_or_update",
name="my_orders_pipeline",
root_path="/Workspace/Users/user@example.com/my_pipeline",
catalog="my_catalog",
schema="my_schema",
workspace_file_paths=[...],
start_run=True
)

The iteration loop is: read the error from result["message"] or manage_pipeline(action="get"), fix the file locally, re-upload, re-run. The same name keeps updating in place — no need to delete and recreate.

“I already have a pipeline. Just trigger a new run with full refresh and wait for it.”

manage_pipeline_run(
action="start",
pipeline_id="abc-123",
full_refresh=True,
wait=True,
timeout=1800
)
# Or refresh only specific tables
manage_pipeline_run(
action="start",
pipeline_id="abc-123",
refresh_selection=["silver_orders", "gold_daily_summary"]
)
# Validate without writing data
manage_pipeline_run(
action="start",
pipeline_id="abc-123",
validate_only=True
)

manage_pipeline_run handles run lifecycle independent of the pipeline definition. Use validate_only=True to dry-run a config check without consuming compute, and refresh_selection to reprocess specific tables only.

“Create a pipeline with development mode, continuous execution, and pipeline-level Spark configuration”

manage_pipeline(
action="create_or_update",
name="my_streaming_pipeline",
root_path="/Workspace/Users/user@example.com/my_pipeline",
catalog="my_catalog",
schema="my_schema",
workspace_file_paths=[...],
extra_settings={
"development": True,
"continuous": True,
"configuration": {
"spark.sql.shuffle.partitions": "auto"
},
"tags": {"environment": "development", "owner": "data-team"}
}
)

The extra_settings dict accepts any pipeline configuration parameter — development mode, continuous execution, custom tags, notifications, cluster overrides, Python dependencies. See the Advanced Configuration page for the full reference.

  • No version control for pipeline settings — MCP creates pipelines imperatively. If you lose the conversation, you lose the configuration. Export to Asset Bundles before going to production.
  • workspace_file_paths must be absolute — relative paths or local paths fail. Always use the full /Workspace/Users/... path that matches where manage_workspace_files placed the files.
  • Pipeline ID changes on delete-and-recreate — if you delete and recreate, external references (job triggers, monitoring dashboards) break. Update in place with manage_pipeline(action="create_or_update") using the same name instead.
  • Don’t validate with execute_sql COUNT queries — get_table_stats_and_schema returns row counts and schema in one call and is meaningfully faster.