Skip to content

About

A db2-to-spark-sql-modernization-engine project

Resources

Stars

0 stars

Watchers

0 watching

Forks

Latest commit

ย 

History

1 Commit

Folders and files

NameName
Last commit message
Last commit date
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 
ย 

Repository files navigation

โšก DB2 โ†’ Spark SQL Modernization Engine

A hands-on learning project that accepts IBM DB2 SQL statements and generates equivalent Apache Spark SQL and PySpark DataFrame code โ€” complete with query analysis, complexity scoring, data lineage, and optimization recommendations.


๐ŸŽฏ What This Project Teaches

Concept What You Learn
DB2 Architecture Tables, indexes, packages, plans, tablespaces, buffer pools
SQL Parsing How text becomes an AST (Abstract Syntax Tree)
Dialect Translation Why DB2 SQL differs from Spark SQL and how to convert
PySpark DataFrame API How SQL maps to transformation chains
Query Optimization Broadcast joins, partitioning, predicate pushdown
Data Lineage Tracking data flow from source to output
Migration Patterns How organizations move DB2 workloads to Spark

๐Ÿ—๏ธ Project Structure

db2-spark-modernizer/
โ”œโ”€โ”€ src/
โ”‚   โ”œโ”€โ”€ ast/                    # AST node definitions
โ”‚   โ”‚   โ””โ”€โ”€ nodes.py            # SelectNode, TableNode, JoinNode, etc.
โ”‚   โ”œโ”€โ”€ parser/                 # SQL Parser
โ”‚   โ”‚   โ””โ”€โ”€ sql_parser.py       # DB2SQLParser (sqlglot-backed)
โ”‚   โ”œโ”€โ”€ analyzer/               # Query Analysis
โ”‚   โ”‚   โ”œโ”€โ”€ query_analyzer.py   # Metadata extraction
โ”‚   โ”‚   โ””โ”€โ”€ complexity_analyzer.py  # 1-10 complexity scoring
โ”‚   โ”œโ”€โ”€ spark_sql_generator/    # Spark SQL Generation
โ”‚   โ”‚   โ””โ”€โ”€ generator.py        # sqlglot dialect transpilation
โ”‚   โ”œโ”€โ”€ dataframe_generator/    # PySpark Code Generation
โ”‚   โ”‚   โ””โ”€โ”€ generator.py        # DataFrame API code synthesis
โ”‚   โ”œโ”€โ”€ lineage/                # Data Lineage
โ”‚   โ”‚   โ””โ”€โ”€ lineage_engine.py   # DAG + Mermaid diagram
โ”‚   โ”œโ”€โ”€ optimization/           # Spark Optimization
โ”‚   โ”‚   โ””โ”€โ”€ optimizer.py        # Recommendations engine
โ”‚   โ””โ”€โ”€ docs_generator/         # Migration Documentation
โ”‚       โ””โ”€โ”€ generator.py        # Auto-generated migration docs
โ”œโ”€โ”€ ui/
โ”‚   โ””โ”€โ”€ app.py                  # Streamlit application
โ”œโ”€โ”€ examples/
โ”‚   โ””โ”€โ”€ scenarios.py            # 12 real-world migration scenarios
โ”œโ”€โ”€ tests/
โ”‚   โ”œโ”€โ”€ test_parser.py          # 30 parser tests
โ”‚   โ”œโ”€โ”€ test_analyzer.py        # 22 analyzer tests
โ”‚   โ””โ”€โ”€ test_spark_generator.py # 39 generator + lineage tests
โ”œโ”€โ”€ docs/
โ”‚   โ””โ”€โ”€ mapping_catalog.md      # DB2 โ†’ Spark concept mapping
โ”œโ”€โ”€ requirements.txt
โ””โ”€โ”€ README.md

๐Ÿš€ Quick Start

1. Install Dependencies

pip install -r requirements.txt

requirements.txt:

sqlglot>=25.0.0
streamlit>=1.35.0
graphviz>=0.20
pytest>=8.0.0
pandas>=2.0.0

2. Run the Streamlit UI

streamlit run ui/app.py

Open http://localhost:8501 in your browser.

3. Run Tests

pytest tests/ -v
# 91 tests, all passing

4. Use the Python API Directly

from src.parser.sql_parser import parse_db2_sql
from src.analyzer.query_analyzer import analyze_query
from src.analyzer.complexity_analyzer import analyze_complexity
from src.spark_sql_generator.generator import generate_spark_sql
from src.dataframe_generator.generator import generate_pyspark_code

db2_sql = """
SELECT   C.CUSTOMER_ID,
         C.CUSTOMER_NAME,
         SUM(T.AMOUNT) AS TOTAL_AMOUNT
FROM     BANKDB.CUSTOMER    C
         INNER JOIN BANKDB.TRANSACTION T ON C.CUSTOMER_ID = T.CUSTOMER_ID
WHERE    C.STATUS = 'ACTIVE'
GROUP BY C.CUSTOMER_ID, C.CUSTOMER_NAME
HAVING   SUM(T.AMOUNT) > 10000
ORDER BY TOTAL_AMOUNT DESC
FETCH FIRST 100 ROWS ONLY
"""

# Parse
result     = parse_db2_sql(db2_sql)

# Analyze
metadata   = analyze_query(result)
complexity = analyze_complexity(result, metadata)

# Generate
spark_sql  = generate_spark_sql(db2_sql)
pyspark    = generate_pyspark_code(result)

print(spark_sql.spark_sql)
print(pyspark.code)
print(f"Complexity: {complexity.score}/10 ({complexity.level})")

๐Ÿ–ฅ๏ธ Streamlit UI

The UI has 8 tabs for each migration artifact:

Tab Content
๐Ÿ” Parsed AST Visual tree of query nodes
๐Ÿ“‹ Metadata Tables, columns, filters, features
โšก Spark SQL Side-by-side DB2 vs Spark SQL
๐Ÿ PySpark Code DataFrame transformation chain
๐Ÿ“Š Complexity Score (1-10), factors, challenges
๐Ÿ”— Lineage Data flow graph + Mermaid diagram
๐Ÿš€ Optimizations Ranked recommendations
๐Ÿ“„ Migration Doc Auto-generated migration document

Sidebar features:

  • Load 12 pre-built industry examples (Banking, Insurance, Retail, Telecom, Healthcare)
  • Toggle AST JSON view
  • Toggle DB2 concept explanations
  • Download migration documents (.md)

๐Ÿ“š Supported DB2 SQL Features

Feature Status Notes
SELECT columns, aliases โœ… Including AS aliases
SELECT * โœ…
SELECT DISTINCT โœ…
FROM with schema qualification โœ… SCHEMA.TABLE
Table aliases โœ…
INNER JOIN โœ…
LEFT JOIN โœ…
RIGHT JOIN โœ…
FULL OUTER JOIN โœ…
CROSS JOIN โœ…
WHERE with AND/OR โœ…
IN, BETWEEN, LIKE โœ…
GROUP BY โœ…
HAVING โœ…
ORDER BY ASC/DESC โœ…
COUNT, SUM, AVG, MIN, MAX โœ…
UNION / UNION ALL โœ…
FETCH FIRST n ROWS ONLY โœ… โ†’ LIMIT n
CURRENT DATE / CURRENT TIMESTAMP โœ… DB2 special registers
WITH UR/CS/RS/RR isolation hints โœ… Removed with warning
Subqueries โš ๏ธ Basic support
CASE expressions โš ๏ธ Parsed, basic codegen
Window functions (OVER) โš ๏ธ Detected, flagged
Stored procedures โŒ Documentation only

๐Ÿญ Industry Scenarios

12 pre-built examples across 5 industries:

๐Ÿฆ Banking

  • bank_001 โ€” Customer Balance Lookup (Simple)
  • bank_002 โ€” Daily Transaction Summary (Medium)
  • bank_003 โ€” Customer Account Portfolio Analysis (Medium)
  • bank_004 โ€” Risk Portfolio Summary (Complex)
  • bank_005 โ€” UNION ALL Account Report (Medium)

๐Ÿฅ Insurance

  • ins_001 โ€” Active Policy Report by Agent (Medium)
  • ins_002 โ€” Claims Loss Ratio Analysis (Complex)

๐Ÿ›’ Retail

  • ret_001 โ€” Product Sales Analysis (Medium)
  • ret_002 โ€” Low Inventory Alert (Medium)

๐Ÿ“ก Telecom

  • tel_001 โ€” Subscriber Usage for Billing (Medium)

๐Ÿฅ Healthcare

  • hc_001 โ€” Patient Billing Summary (Complex)

๐Ÿ”ฎ COBOL

  • cobol_001 โ€” Embedded SQL Extraction Example

๐Ÿง  Learning Concepts

Phase 1: DB2 Architecture

โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”
โ”‚             z/OS Mainframe               โ”‚
โ”‚  โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”  โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”  โ”‚
โ”‚  โ”‚  CICS/IMS    โ”‚  โ”‚  DB2 Subsystem  โ”‚  โ”‚
โ”‚  โ”‚  (Online TP) โ”‚โ”€โ–ถโ”‚  โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”  โ”‚  โ”‚
โ”‚  โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜  โ”‚  โ”‚ Optimizer โ”‚  โ”‚  โ”‚
โ”‚  โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”  โ”‚  โ”‚ Buffer    โ”‚  โ”‚  โ”‚
โ”‚  โ”‚  Batch JCL   โ”‚โ”€โ–ถโ”‚  โ”‚ Pools     โ”‚  โ”‚  โ”‚
โ”‚  โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜  โ”‚  โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜  โ”‚  โ”‚
โ”‚                    โ”‚  Tablespaces     โ”‚  โ”‚
โ”‚                    โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜  โ”‚
โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜
             โฌ‡ Migration
โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”
โ”‚          Apache Spark / Delta Lake        โ”‚
โ”‚  Driver โ”€ Catalyst Optimizer             โ”‚
โ”‚  โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ” โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ” โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ” โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ”   โ”‚
โ”‚  โ”‚Exec 1โ”‚ โ”‚Exec 2โ”‚ โ”‚Exec 3โ”‚ โ”‚Exec 4โ”‚   โ”‚
โ”‚  โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”˜ โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”˜ โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”˜ โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”˜   โ”‚
โ”‚         Delta Lake (cloud storage)        โ”‚
โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜

Phase 2: How the Parser Works

DB2 SQL Text
    โ†“ Lexer (sqlglot)
Tokens: [SELECT, CUSTOMER_ID, FROM, CUSTOMER, WHERE, ...]
    โ†“ Parser
sqlglot AST (internal representation)
    โ†“ Our converter
Educational AST nodes:
    SelectNode
    โ”œโ”€โ”€ TableNode(CUSTOMER)
    โ”œโ”€โ”€ ColumnNode(CUSTOMER_ID)
    โ”œโ”€โ”€ WhereNode(BALANCE > 1000)
    โ””โ”€โ”€ OrderByNode(CUSTOMER_ID ASC)

Phase 3: DB2 โ†’ Spark Concept Mapping

DB2 Concept Spark Equivalent
Table DataFrame / Delta Lake Table
View Temp View / DataFrame
Index Partitioning + Z-ORDER
Tablespace Database / Schema
Buffer Pool Spark Memory Cache
Package/Plan Catalyst Physical Plan
RUNSTATS ANALYZE TABLE (Delta)
Stored Procedure Spark UDF / Python Function
Cursor DataFrame iterator / collect()
COMMIT/ROLLBACK Delta Lake ACID Transactions
FETCH FIRST n ROWS LIMIT n
JCL Batch Job Spark Job / Databricks Job

Phase 4-5: SQL Translation Examples

Simple SELECT:

-- DB2
SELECT CUSTOMER_ID, NAME, BALANCE
FROM   BANKDB.CUSTOMER
WHERE  STATUS = 'ACTIVE'
ORDER BY NAME

-- Spark SQL (identical โ€” ANSI SQL)
SELECT CUSTOMER_ID, NAME, BALANCE
FROM   BANKDB.CUSTOMER
WHERE  STATUS = 'ACTIVE'
ORDER BY NAME

-- PySpark DataFrame API
customer_df = spark.table("BANKDB.CUSTOMER")
customer_df = customer_df.filter(col("STATUS") == "ACTIVE")
customer_df = customer_df.select("CUSTOMER_ID", "NAME", "BALANCE")
customer_df = customer_df.orderBy(col("NAME").asc())

DB2-Specific Syntax:

-- DB2: FETCH FIRST
SELECT * FROM CUSTOMER FETCH FIRST 100 ROWS ONLY

-- Spark SQL
SELECT * FROM CUSTOMER LIMIT 100

-- PySpark
customer_df.limit(100)
-- DB2: CURRENT DATE (special register)
WHERE TXN_DATE = CURRENT DATE

-- Spark SQL
WHERE TXN_DATE = CURRENT_DATE()

-- PySpark
.filter(col("TXN_DATE") == current_date())

Phase 7: Complexity Scoring

Score Level Migration Effort Auto-Migration
1โ€“3 ๐ŸŸข Simple Hours 85โ€“95%
4โ€“5 ๐ŸŸก Medium Days 65โ€“85%
6โ€“7 ๐ŸŸ  Complex Weeks 40โ€“65%
8โ€“10 ๐Ÿ”ด Critical Months <40%

Phase 8: Spark Optimization Patterns

# 1. BROADCAST JOIN (for small dimension tables)
from pyspark.sql.functions import broadcast
result = large_fact.join(broadcast(small_dim), "key")

# 2. PARTITION PUSHDOWN (for large tables)
spark.conf.set("spark.sql.adaptive.enabled", "true")
df = spark.table("fact_table").filter(col("region") == "NORTH")

# 3. AGGREGATE ONCE (all functions in one .agg() call)
result = df.groupBy("region").agg(
    sum("amount").alias("total"),
    count("id").alias("cnt"),
    avg("balance").alias("avg_bal")
)

# 4. SEMI-JOIN (replace IN subquery)
# DB2: WHERE id IN (SELECT id FROM vip_table)
# Spark:
result = df.join(vip_df.select("id"), "id", "left_semi")

๐Ÿงช Running Tests

# All tests
pytest tests/ -v

# With coverage
pytest tests/ -v --cov=src --cov-report=term-missing

# Specific test file
pytest tests/test_parser.py -v
pytest tests/test_analyzer.py -v
pytest tests/test_spark_generator.py -v

Test coverage:

  • test_parser.py โ€” 30 tests: SELECT, WHERE, JOIN, UNION, GROUP BY, HAVING, ORDER BY, preprocessing
  • test_analyzer.py โ€” 22 tests: metadata extraction, complexity scoring
  • test_spark_generator.py โ€” 39 tests: Spark SQL gen, PySpark gen, lineage engine

๐Ÿ—บ๏ธ Learning Roadmap

MVP (Complete โœ…)

  • DB2 SQL Parser (sqlglot-backed)
  • Educational AST nodes
  • Query metadata extraction
  • Complexity scoring (1-10)
  • Spark SQL generation (dialect transpilation)
  • PySpark DataFrame code generation
  • Data lineage graph + Mermaid diagrams
  • Spark optimization recommendations
  • Migration documentation generator
  • Streamlit UI (8 tabs)
  • 12 industry scenarios
  • 91 pytest unit tests

Stretch Goals (Next Steps)

  • Subquery support (full correlated subqueries)
  • Common Table Expressions (CTEs / WITH clause)
  • Window functions (ROW_NUMBER, RANK, LEAD, LAG)
  • Stored procedure analysis
  • COBOL embedded SQL extractor
  • Delta Lake DDL generation
  • Query performance estimator
  • AI-generated migration explanations (Claude API)
  • SQL dialect comparison dashboard
  • Batch migration: analyze a folder of SQL files

๐Ÿข Real-World Context

In large-scale mainframe modernization projects:

  1. IBM DataStage / AWS DMS โ€” extract DB2 data to cloud storage
  2. Tools like this โ€” analyze and translate SQL workloads
  3. Apache Spark / Databricks โ€” run translated workloads at scale
  4. Delta Lake โ€” replace DB2 tables with ACID-compliant lake tables
  5. Apache Atlas / OpenLineage โ€” enterprise lineage tracking

This project teaches the analytical and translation phase that sits between extraction and execution.


๐Ÿ› ๏ธ Tech Stack

Component Technology
SQL Parsing sqlglot
UI Streamlit
Diagrams Mermaid (text-based)
Testing pytest
Data manipulation pandas
Language Python 3.9+

Built as a learning-focused educational tool for DB2 developers moving to Apache Spark and modern data lake architectures.

About

A db2-to-spark-sql-modernization-engine project

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages