| name | databricks-migration-deep-dive |
| description | Execute comprehensive platform migrations to Databricks from legacy systems.
Use when migrating from on-premises Hadoop, other cloud platforms,
or legacy data warehouses to Databricks.
Trigger with phrases like "migrate to databricks", "hadoop migration",
"snowflake to databricks", "legacy migration", "data warehouse migration".
|
| allowed-tools | Read, Write, Edit, Bash(databricks:*), Grep |
| version | 1.0.0 |
| license | MIT |
| author | Jeremy Longshore <jeremy@intentsolutions.io> |
Databricks Migration Deep Dive
Overview
Comprehensive migration strategies for moving to Databricks from legacy systems.
Prerequisites
- Access to source and target systems
- Understanding of current data architecture
- Migration timeline and requirements
- Stakeholder alignment
Migration Patterns
| Source System | Migration Pattern | Complexity | Timeline |
|---|
| On-prem Hadoop | Lift-and-shift + modernize | High | 6-12 months |
| Snowflake | Parallel run + cutover | Medium | 3-6 months |
| AWS Redshift | ETL rewrite + data copy | Medium | 3-6 months |
| Azure Synapse | Delta Lake conversion | Low | 1-3 months |
| Legacy DW (Oracle/Teradata) | Full rebuild | High | 12-18 months |
Instructions
Step 1: Discovery and Assessment
from dataclasses import dataclass
from typing import List, Dict
import pandas as pd
@dataclass
class SourceTableInfo:
"""Information about source table for migration planning."""
database: str
schema: str
table: str
row_count: int
size_gb: float
column_count: int
partition_columns: List[str]
dependencies: List[str]
access_frequency: str
data_classification: str
def assess_hadoop_cluster(spark, hive_metastore: str) -> List[SourceTableInfo]:
"""
Assess Hadoop/Hive cluster for migration planning.
Returns inventory of all tables with metadata.
"""
tables = []
databases = spark.sql("SHOW DATABASES").collect()
for db_row in databases:
db = db_row.databaseName
if db in ['default', 'sys']:
continue
spark.sql(f"USE {db}")
table_list = spark.sql().collect()
table_row table_list:
table_name = table_row.tableName
:
desc = spark.sql()
detail_df = desc.toPandas()
partition_cols = []
in_partition_section =
_, row detail_df.iterrows():
row[] == :
in_partition_section =
in_partition_section row[] row[].startswith():
partition_cols.append(row[])
stats = spark.sql()
tables.append(SourceTableInfo(
database=db,
schema=db,
table=table_name,
row_count=,
size_gb=,
column_count=(detail_df),
partition_columns=partition_cols,
dependencies=[],
access_frequency=,
data_classification=,
))
Exception e:
()
tables
() -> pd.DataFrame:
plan_data = []
table tables:
complexity =
complexity += table.size_gb >
complexity += (table.partition_columns) >
complexity += (table.dependencies) >
complexity += table.data_classification ==
priority =
priority += table.access_frequency ==
priority += table.data_classification ==
plan_data.append({
: ,
: ,
: table.size_gb,
: complexity,
: priority,
: (, table.size_gb / + complexity * ),
: priority > ( priority > ),
})
pd.DataFrame(plan_data).sort_values([, ], ascending=[, ])
Step 2: Schema Migration
from pyspark.sql import SparkSession
from pyspark.sql.types import *
def convert_hive_to_delta_schema(spark: SparkSession, hive_table: str) -> StructType:
"""
Convert Hive table schema to Delta Lake compatible schema.
Handles type conversions and incompatibilities.
"""
hive_schema = spark.table(hive_table).schema
type_conversions = {
'decimal(38,0)': DecimalType(38, 10),
'char': StringType(),
'varchar': StringType(),
'tinyint': IntegerType(),
}
new_fields = []
for field in hive_schema.fields:
new_type = field.dataType
type_str = str(field.dataType).lower()
for pattern, replacement in type_conversions.items():
if pattern in type_str:
new_type = replacement
break
new_fields.append(StructField(
field.name,
new_type,
field.nullable,
field.metadata
))
return StructType(new_fields)
def migrate_table_schema(
spark: SparkSession,
source_table: str,
target_table: str,
catalog: str = "migrated",
) -> dict:
"""
Migrate table schema from Hive to Delta Lake.
Returns migration result with any schema changes.
"""
source_df = spark.table(source_table)
source_schema = source_df.schema
target_schema = convert_hive_to_delta_schema(spark, source_table)
schema_ddl = .join([
f target_schema.fields
])
spark.sql()
changes = []
i, (src, tgt) ((source_schema.fields, target_schema.fields)):
(src.dataType) != (tgt.dataType):
changes.append({
: src.name,
: (src.dataType),
: (tgt.dataType),
})
{
: source_table,
: ,
: (target_schema.fields),
: changes,
}
Step 3: Data Migration
from pyspark.sql import SparkSession, DataFrame
from datetime import datetime
import time
class DataMigrator:
"""Handle data migration from legacy systems to Delta Lake."""
def __init__(self, spark: SparkSession, target_catalog: str):
self.spark = spark
self.target_catalog = target_catalog
def migrate_table(
self,
source_table: str,
target_table: str,
batch_size: int = 1000000,
partition_columns: list[str] = None,
incremental_column: str = None,
) -> dict:
"""
Migrate table data with batching and checkpointing.
Args:
source_table: Source table (Hive, JDBC, etc.)
target_table: Target Delta table
batch_size: Rows per batch
partition_columns: Columns for partitioning
incremental_column: Column for incremental loads
Returns:
Migration statistics
"""
start_time = time.time()
stats = {
'source_table': source_table,
'target_table': f"{self.target_catalog}.{target_table}",
'batches': 0,
'total_rows': 0,
'errors': [],
}
source_df = .spark.table(source_table)
partition_columns (partition_columns) > :
partitions = source_df.select(partition_columns).distinct().collect()
partition_row partitions:
conditions = [
col partition_columns
]
filter_expr = .join(conditions)
batch_df = source_df.(filter_expr)
._write_batch(batch_df, target_table, partition_columns)
stats[] +=
stats[] += batch_df.count()
:
._write_batch(source_df, target_table, partition_columns)
stats[] =
stats[] = source_df.count()
stats[] = time.time() - start_time
stats[] = stats[] / stats[]
stats
():
writer = df.write.().mode()
partition_columns:
writer = writer.partitionBy(*partition_columns)
writer.saveAsTable()
() -> :
source_df = .spark.table(source_table)
target_df = .spark.table()
validation = {
: source_df.count(),
: target_df.count(),
: ,
: ,
: ,
}
validation[] = (
validation[] == validation[]
)
source_cols = (source_df.columns)
target_cols = (target_df.columns)
validation[] = source_cols == target_cols
source_sample = source_df.limit().toPandas()
target_sample = target_df.limit().toPandas()
validation[] = (source_sample) == (target_sample)
validation
Step 4: ETL/Pipeline Migration
def convert_spark_job_to_databricks(
source_code: str,
source_type: str = "spark-submit",
) -> str:
"""
Convert legacy Spark job to Databricks job.
Handles common patterns from spark-submit, Oozie, Airflow.
"""
replacements = {
'SparkSession.builder.master("yarn")': 'SparkSession.builder',
'.master("local[*]")': '',
'hdfs://namenode:8020/': '/mnt/data/',
's3a://': 's3://',
'.enableHiveSupport()': '',
'hive_metastore.': '',
'.config("spark.sql.warehouse.dir"': '# Removed: .config("spark.sql.warehouse.dir"',
}
converted = source_code
for old, new in replacements.items():
converted = converted.replace(old, new)
header = '''
# Converted for Databricks
# Original source: {source_type}
# Conversion date: {date}
from pyspark.sql import SparkSession
# SparkSession is pre-configured in Databricks
spark = SparkSession.builder.getOrCreate()
'''.format(source_type=source_type, date=datetime.now().isoformat())
return header + converted
def () -> :
xml.etree.ElementTree ET
root = ET.fromstring(oozie_xml)
tasks = []
action root.findall():
action_name = action.get()
spark_action = action.find()
spark_action :
jar = spark_action.find().text
main_class = spark_action.find().text
tasks.append({
: action_name,
: {
: main_class,
: [],
},
: [{: jar}],
})
shell_action = action.find()
shell_action :
job_definition = {
: ,
: tasks,
: [{
: ,
: {
: ,
: ,
: ,
}
}],
}
job_definition
Step 5: Cutover Planning
from dataclasses import dataclass
from datetime import datetime, timedelta
from typing import List
@dataclass
class CutoverStep:
"""Individual step in cutover plan."""
order: int
name: str
duration_minutes: int
owner: str
rollback_procedure: str
verification: str
def generate_cutover_plan(
migration_wave: int,
tables: List[str],
cutover_date: datetime,
) -> List[CutoverStep]:
"""Generate detailed cutover plan."""
steps = [
CutoverStep(
order=1,
name="Pre-cutover validation",
duration_minutes=60,
owner="Data Engineer",
rollback_procedure="N/A - no changes made",
verification="Run validation queries on all tables",
),
CutoverStep(
order=2,
name="Disable source data pipelines",
duration_minutes=15,
owner="Platform Admin",
rollback_procedure="Re-enable pipelines in source system",
verification="Verify no new data in source",
),
CutoverStep(
order=3,
name="Final incremental sync",
duration_minutes=120,
owner=,
rollback_procedure=,
verification=,
),
CutoverStep(
order=,
name=,
duration_minutes=,
owner=,
rollback_procedure=,
verification=,
),
CutoverStep(
order=,
name=,
duration_minutes=,
owner=,
rollback_procedure=,
verification=,
),
CutoverStep(
order=,
name=,
duration_minutes=,
owner=,
rollback_procedure=,
verification=,
),
]
current_time = cutover_date
step steps:
step.start_time = current_time
step.end_time = current_time + timedelta(minutes=step.duration_minutes)
current_time = step.end_time
steps
Output
- Migration assessment complete
- Schema migration automated
- Data migration pipeline ready
- ETL conversion scripts
- Cutover plan documented
Error Handling
| Issue | Cause | Solution |
|---|
| Schema incompatibility | Unsupported types | Use type conversion mappings |
| Data loss | Truncation | Validate counts at each step |
| Performance issues | Large tables | Use partitioned migration |
| Dependency conflicts | Wrong migration order | Analyze dependencies first |
Examples
Quick Migration Validation
SELECT
'source' as system, COUNT(*) as row_count
FROM hive_metastore.db.table
UNION ALL
SELECT
'target' as system, COUNT(*) as row_count
FROM migrated.db.table;
Resources
Completion
This skill pack provides comprehensive coverage for Databricks platform operations.