SECTION 1: LEARNING OBJECTIVES
By the end of this lesson, you will be able to:
-
Define Data Lakes and Data Lakehouses and explain their role in modern banking data architecture.
-
Compare Data Lakes vs. Data Warehouses, understanding when each is appropriate for banking workloads.
-
Understand the Data Lakehouse concept and why banks are adopting it as a unified architecture.
-
Explain the Medallion Architecture (Bronze, Silver, Gold layers) used in modern banking data platforms.
-
Identify real-world banking use cases for Data Lakes and Lakehouses (fraud detection, risk modeling, customer analytics).
-
Understand the key technologies used in modern banking data platforms (Apache Spark, Delta Lake, Iceberg, Hudi).
SECTION 2: UNDERSTANDING DATA LAKES
2.1 What Is a Data Lake?
A Data Lake is a centralized repository that stores vast amounts of raw data in its native format. Unlike data warehouses (which store cleaned, structured data), data lakes store data “as-is” with no transformation.
Key Characteristics of Data Lakes:
| Characteristic | Description | Banking Example |
|---|---|---|
| Schema-on-Read | Data structure applied when read, not when written | Store raw ATM logs; parse later |
| Store Everything | All data, regardless of structure | Ledgers + emails + call transcripts |
| Raw Storage | No transformation required | Save exactly what source provides |
| Massive Scale | Handles petabytes | All market tick data + transaction history |
| Cost-Effective | Low-cost storage | Cloud object storage (S3, Azure Blob) |
2.2 What Goes Into a Banking Data Lake?
BANKING DATA LAKE CONTENTS: ┌─────────────────────────────────────────────────────────────────┐ │ STRUCTURED DATA │ │ • Transaction logs (OLTP extracts) │ │ • Account balances (daily snapshots) │ │ • Customer master data │ │ • Loan portfolio data │ │ • Reference data (currency codes, branch lists) │ ├─────────────────────────────────────────────────────────────────┤ │ SEMI-STRUCTURED DATA │ │ • JSON from APIs (credit bureau responses) │ │ • XML from SWIFT messages │ │ • Log files from ATM/POS systems │ │ • Web server logs from online banking │ ├─────────────────────────────────────────────────────────────────┤ │ UNSTRUCTURED DATA │ │ • Customer emails and chat transcripts │ │ • Call center recordings (audio) │ │ • PDF statements and loan documents │ │ • SEC filings and regulatory submissions │ │ • Social media posts about the bank │ │ • Market news feeds │ ├─────────────────────────────────────────────────────────────────┤ │ STREAMING DATA │ │ • Real-time market data (stock ticks) │ │ • ATM transaction feeds │ │ • POS transaction streams │ │ • Fraud alerts (real-time) │ └─────────────────────────────────────────────────────────────────┘
2.3 Data Lake Use Cases in Banking
Use Case 1: Fraud Detection
Scenario: A bank needs to detect fraud patterns in real-time using multiple data sources.
# Example: Streaming fraud detection using data lake # In production, this would use Apache Flink or Spark Streaming def detect_fraud(transaction, historical_patterns): """Check transaction against fraud patterns from data lake.""" # Load customer's typical behavior from data lake customer_history = data_lake.read( path=f"customer_patterns/{transaction.customer_id}", format="parquet" ) # Load global fraud patterns from data lake fraud_patterns = data_lake.read( path="fraud_patterns/latest", format="delta" ) # Check for anomalies if transaction.amount > customer_history.avg_amount * 3: return "FRAUD_ALERT: Amount deviation" if transaction.location not in customer_history.typical_locations: return "FRAUD_ALERT: Location anomaly" if transaction.merchant in fraud_patterns.high_risk_merchants: return "FRAUD_ALERT: High-risk merchant" return "CLEAR" # Store raw transaction for future analysis data_lake.write( path=f"raw_transactions/{transaction.date}/", data=transaction, format="parquet", partition_by=["date", "status"] )
Use Case 2: Risk Modeling
Scenario: A data scientist needs to build a new credit risk model using 10 years of historical data plus external economic indicators.
# Example: Building a risk model using data lake def prepare_risk_modeling_data(): """Extract and combine data from multiple sources in the lake.""" # 1. Loan performance data (10 years) loan_data = data_lake.read( path="structured/loan_performance/", format="parquet", filter={"date": {"gte": "2014-01-01"}} ) # 2. Customer credit history credit_data = data_lake.read( path="structured/credit_bureau/", format="parquet" ) # 3. Macroeconomic indicators macro_data = data_lake.read( path="structured/economic_indicators/", format="parquet" ) # 4. Unstructured loan documents (NLP processed) document_features = data_lake.read( path="processed/document_features/", format="parquet" ) # Combine all data combined = ( loan_data .join(credit_data, "customer_id") .join(macro_data, "quarter") .join(document_features, "loan_id") ) return combined # Data scientist can now build model using PySpark data = prepare_risk_modeling_data() features = ["credit_score", "dti_ratio", "gdp_growth", "document_sentiment"] model = train_xgboost(data[features], data["default"])
Use Case 3: Regulatory Compliance
Scenario: A regulator requests all data related to a specific customer for the past 5 years.
# Example: Data lake query for regulatory request def fulfill_regulatory_request(customer_id, start_date, end_date): """Retrieve all data for a customer from the data lake.""" request_data = {} # 1. Transaction data (structured) request_data['transactions'] = data_lake.read( path="structured/transactions/", format="parquet", filter={ "customer_id": customer_id, "date": {"gte": start_date, "lte": end_date} } ) # 2. Account data request_data['accounts'] = data_lake.read( path="structured/accounts/", format="parquet", filter={"customer_id": customer_id} ) # 3. Loan data request_data['loans'] = data_lake.read( path="structured/loans/", format="parquet", filter={"customer_id": customer_id} ) # 4. Communication records (unstructured) request_data['emails'] = data_lake.read( path="unstructured/emails/", format="text", filter={"customer_id": customer_id} ) # 5. Call transcripts request_data['calls'] = data_lake.read( path="unstructured/call_transcripts/", format="parquet", filter={"customer_id": customer_id} ) # 6. Compliance notes request_data['notes'] = data_lake.read( path="structured/compliance_notes/", format="parquet", filter={"customer_id": customer_id} ) # Generate comprehensive report generate_regulatory_report(request_data) return request_data
SECTION 3: THE DATA LAKEHOUSE
3.1 What Is a Data Lakehouse?
A Data Lakehouse combines the best of data lakes and data warehouses. It provides the flexibility and scale of a data lake with the performance and features of a data warehouse.
Data Lakehouse Architecture:
┌─────────────────────────────────────────────────────────────────┐ │ DATA LAKEHOUSE │ │ │ │ ┌─────────────────────────────────────────────────────────┐ │ │ │ LAYER 1: BRONZE │ │ │ │ Raw Data (source of truth) │ │ │ │ • All data stored as-is │ │ │ │ • Immutable (no modifications) │ │ │ │ • Partitioned for performance │ │ │ └─────────────────────────────────────────────────────────┘ │ │ ↓ │ │ ┌─────────────────────────────────────────────────────────┐ │ │ │ LAYER 2: SILVER │ │ │ │ Cleaned Data │ │ │ │ • Data quality applied │ │ │ │ • Deduplication │ │ │ │ • Standardized formats │ │ │ │ • Business rules applied │ │ │ └─────────────────────────────────────────────────────────┘ │ │ ↓ │ │ ┌─────────────────────────────────────────────────────────┐ │ │ │ LAYER 3: GOLD │ │ │ │ Business-Ready Data │ │ │ │ • Aggregated and modeled │ │ │ │ • Dimensional modeling │ │ │ │ • Pre-calculated metrics │ │ │ │ • Query-optimized │ │ │ └─────────────────────────────────────────────────────────┘ │ └─────────────────────────────────────────────────────────────────┘
3.2 The Medallion Architecture (Bronze, Silver, Gold)
Bronze Layer (Raw):
-
Immutable snapshot of source data
-
Minimal processing (just enough to handle partitions)
-
Complete audit trail of what came from source
# Example: Bronze layer - raw data def load_bronze(source_path, target_path): """Load raw data to Bronze layer.""" # Read raw data df = spark.read.format(source_format).load(source_path) # Minimal processing: add audit columns df = df.withColumn("ingestion_date", current_date()) df = df.withColumn("ingestion_timestamp", current_timestamp()) df = df.withColumn("source_system", lit(source_system_name)) df = df.withColumn("record_hash", md5(concat_ws("||", *df.columns))) # Write to Bronze df.write.format("delta") \ .mode("append") \ .partitionBy("ingestion_date") \ .save(f"{target_path}/bronze/")
Silver Layer (Cleaned):
-
Data quality applied
-
Deduplication
-
Standardized schemas
-
Light transformations
# Example: Silver layer - cleaned data def transform_to_silver(bronze_path, silver_path): """Transform Bronze data to Silver.""" # Read from Bronze df = spark.read.format("delta").load(bronze_path) # Apply transformations transformed = df \ .filter(col("status") != "TEST") \ .dropDuplicates(["transaction_id"]) \ .withColumn("date", to_date(col("transaction_date"))) \ .withColumn("year", year(col("date"))) \ .withColumn("month", month(col("date"))) \ .na.fill({"amount": 0, "customer_id": "UNKNOWN"}) # Validate data quality quality_df = transformed \ .filter(col("amount") < 0) \ .select("transaction_id") \ .count() if quality_df > 0: # Log negative amounts spark.createDataFrame([{"issue": "negative_amount", "count": quality_df}]) \ .write.format("delta") \ .mode("append") \ .save(f"{silver_path}/quality_issues/") # Write to Silver transformed.write.format("delta") \ .mode("overwrite") \ .partitionBy("year", "month") \ .save(f"{silver_path}/transactions/")
Gold Layer (Business-Ready):
-
Dimensional modeling (Star Schemas)
-
Aggregated metrics
-
Pre-calculated KPIs
-
Query-optimized
# Example: Gold layer - business-ready data def build_gold_dimensional(silver_path, gold_path): """Build Gold layer dimensional model.""" # Read Silver data transactions = spark.read.format("delta").load(f"{silver_path}/transactions/") customers = spark.read.format("delta").load(f"{silver_path}/customers/") accounts = spark.read.format("delta").load(f"{silver_path}/accounts/") # Build fact table fact_transactions = transactions \ .join(accounts.select("account_id", "customer_id"), "account_id") \ .join(customers.select("customer_id", "segment", "state"), "customer_id") \ .select( col("transaction_id"), col("date").alias("transaction_date"), col("year"), col("month"), col("amount"), col("transaction_type"), col("account_id"), col("customer_id"), col("segment").alias("customer_segment"), col("state").alias("customer_state"), col("amount_category") ) # Build dimensions dim_customer = customers.select( col("customer_id").alias("customer_key"), col("customer_id"), col("name"), col("segment"), col("state"), col("credit_score"), col("annual_income") ) dim_date = fact_transactions.select( col("transaction_date").alias("date_key"), col("transaction_date"), col("year"), col("month"), dayofmonth(col("transaction_date")).alias("day"), dayofweek(col("transaction_date")).alias("day_of_week"), when(dayofweek(col("transaction_date")).isin([6,7]), True).otherwise(False).alias("is_weekend") ).distinct() # Write Gold layer fact_transactions.write.format("delta") \ .mode("overwrite") \ .partitionBy("year", "month") \ .save(f"{gold_path}/fact_transactions/") dim_customer.write.format("delta") \ .mode("overwrite") \ .save(f"{gold_path}/dim_customer/") dim_date.write.format("delta") \ .mode("overwrite") \ .save(f"{gold_path}/dim_date/")
3.3 Why Banks Are Adopting Lakehouses
| Benefit | Why It Matters to Banks |
|---|---|
| Single Platform | One place for all data (structured + unstructured) |
| Lower Costs | Storage on cheap object storage, not expensive DW |
| Real-Time Analytics | Can handle streaming data natively |
| Schema Evolution | Easy to add new data types without breaking pipelines |
| Regulatory Compliance | Full data lineage and audit capability |
| Machine Learning Ready | Data scientists can work with raw and transformed data |
| Time Travel | Can query historical data at any point in time |
SECTION 4: KEY TECHNOLOGIES FOR DATA LAKEHOUSES
4.1 Apache Spark
Spark is the foundation for most modern data lakehouses. It provides distributed processing at scale.
Why Banking Uses Spark:
-
Process petabytes of data
-
Run complex transformations
-
Machine learning at scale
-
Streaming capabilities
# Example: Spark in banking analytics from pyspark.sql import SparkSession from pyspark.sql.functions import * spark = SparkSession.builder \ .appName("Banking Analytics") \ .config("spark.sql.adaptive.enabled", "true") \ .config("spark.sql.adaptive.coalescePartitions.enabled", "true") \ .getOrCreate() # Read from data lake transactions = spark.read.format("parquet") \ .load("/mnt/datalake/transactions/") # Complex aggregation daily_summary = transactions \ .filter(col("date") >= "2024-01-01") \ .groupBy("date", "customer_segment") \ .agg( sum("amount").alias("total_amount"), count("*").alias("transaction_count"), avg("amount").alias("avg_amount") ) # Write back to lakehouse daily_summary.write.format("delta") \ .mode("overwrite") \ .partitionBy("date") \ .save("/mnt/lakehouse/gold/daily_summary/")
4.2 Delta Lake
Delta Lake adds transactional features to data lakes:
-
ACID transactions
-
Time travel
-
Schema enforcement
-
Version history
# Example: Delta Lake features def delta_example(): """Demonstrate Delta Lake capabilities.""" # 1. ACID transactions try: spark.sql(""" INSERT INTO transactions SELECT * FROM staging_transactions WHERE transaction_date >= '2024-01-01' """) except Exception as e: # Atomic: all or nothing print(f"Transaction failed: {e}") # 2. Time travel historical_data = spark.read.format("delta") \ .option("versionAsOf", 10) \ .load("/mnt/lakehouse/transactions/") # 3. Schema enforcement df = spark.read.format("delta").load("/mnt/lakehouse/transactions/") df.write.format("delta") \ .mode("append") \ .option("mergeSchema", "false") # Rejects new columns # 4. Vacuum old versions (clean up) spark.sql("VACUUM transactions RETAIN 168 HOURS")
4.3 Apache Iceberg and Apache Hudi
Apache Iceberg: High-performance table format with features similar to Delta Lake.
Apache Hudi: Focuses on incremental processing and upsert capabilities.
Banking Use Case Comparison:
| Feature | Delta Lake | Iceberg | Hudi |
|---|---|---|---|
| ACID Transactions | ✅ | ✅ | ✅ |
| Time Travel | ✅ | ✅ | ✅ |
| Schema Evolution | ✅ | ✅ | ✅ |
| Streaming Support | ✅ | Limited | ✅ |
| Incremental Processing | Limited | Limited | ✅ |
| Open Format | ✅ | ✅ | ✅ |
| Cloud Support | All | All | All |
SECTION 5: REAL-WORLD BANKING DATA LAKEHOUSE EXAMPLES
5.1 Fraud Detection System
# Example: Real-time fraud detection on Lakehouse def fraud_detection_pipeline(): """Streaming fraud detection using Lakehouse architecture.""" # 1. Bronze: Ingest raw transactions raw_transactions = spark.readStream.format("kafka") \ .option("kafka.bootstrap.servers", "broker:9092") \ .option("subscribe", "transactions") \ .load() \ .select(from_json(col("value").cast("string"), transaction_schema).alias("data")) \ .select("data.*") raw_transactions.writeStream.format("delta") \ .outputMode("append") \ .option("checkpointLocation", "/mnt/checkpoints/bronze") \ .start("/mnt/lakehouse/bronze/transactions") # 2. Silver: Validate and clean clean_transactions = raw_transactions \ .filter(col("amount") > 0) \ .filter(col("customer_id").isNotNull()) clean_transactions.writeStream.format("delta") \ .outputMode("append") \ .option("checkpointLocation", "/mnt/checkpoints/silver") \ .start("/mnt/lakehouse/silver/transactions") # 3. Gold: Detect fraud (real-time) fraud_alerts = clean_transactions \ .join(customer_patterns, "customer_id") \ .where(col("amount") > col("avg_amount") * 3) \ .select("*", lit("HIGH_AMOUNT_DEVIATION").alias("fraud_reason")) fraud_alerts.writeStream.format("console") \ .outputMode("append") \ .start() \ .awaitTermination()
5.2 Customer 360 View
# Example: Building Customer 360 view from Lakehouse def build_customer_360(): """Build complete customer view from all sources.""" # 1. Read from different sources transactions = spark.read.format("delta").load("/mnt/lakehouse/gold/fact_transactions/") accounts = spark.read.format("delta").load("/mnt/lakehouse/gold/dim_accounts/") loans = spark.read.format("delta").load("/mnt/lakehouse/gold/dim_loans/") interactions = spark.read.format("parquet").load("/mnt/lakehouse/unstructured/interactions/") # 2. Calculate customer metrics customer_metrics = transactions \ .groupBy("customer_id") \ .agg( count("*").alias("total_transactions"), sum("amount").alias("total_spent"), avg("amount").alias("avg_transaction"), countDistinct("merchant").alias("unique_merchants"), max("date").alias("last_activity") ) # 3. Add customer profile customer_profile = accounts \ .groupBy("customer_id") \ .agg( sum("balance").alias("total_balance"), count("account_id").alias("account_count"), collect_list("account_type").alias("account_types") ) # 4. Add sentiment from interactions customer_sentiment = interactions \ .groupBy("customer_id") \ .agg( avg("sentiment_score").alias("avg_sentiment"), max("sentiment_score").alias("max_sentiment"), count("*").alias("interaction_count") ) # 5. Combine into 360 view customer_360 = customer_metrics \ .join(customer_profile, "customer_id", "full") \ .join(customer_sentiment, "customer_id", "left") \ .fillna({"total_balance": 0, "avg_sentiment": 0, "interaction_count": 0}) # 6. Write result customer_360.write.format("delta") \ .mode("overwrite") \ .save("/mnt/lakehouse/gold/customer_360/") return customer_360
SECTION 6: BUSINESS RISK & FINANCIAL IMPACT
6.1 Benefits of Modern Data Architecture
| Benefit | Financial Impact |
|---|---|
| Unified Data Platform | $20M/year savings on infrastructure |
| Faster Time-to-Market | 50% faster ML model deployment |
| Better Fraud Detection | $500M/year reduction in fraud losses |
| Improved Compliance | $100M/year reduction in regulatory fines |
| Better Customer Insights | $300M/year revenue increase from cross-selling |
6.2 Risks and Mitigations
| Risk | Mitigation |
|---|---|
| Data Governance Failure | Implement strict metadata management |
| Security Breach | Encrypt all data, implement access controls |
| Performance Issues | Partition data, optimize queries |
| Vendor Lock-in | Use open formats (Delta, Iceberg, Parquet) |
| Cost Overrun | Implement auto-scaling and cost monitoring |
SECTION 7: SUMMARY FOR THE DATA PRACTITIONER
7.1 The 1-Minute Elevator Pitch
“Data Lakehouses are the future of banking data architecture. They combine the scale and flexibility of data lakes with the performance and reliability of data warehouses. Using the Medallion Architecture (Bronze, Silver, Gold), banks can store all their data—structured and unstructured—in one place, ensuring data quality and enabling real-time analytics. Technologies like Apache Spark, Delta Lake, and Iceberg make this possible. The result: faster fraud detection, better customer insights, and lower infrastructure costs.”
7.2 Key Takeaways
-
Data Lakes store raw data from all sources at low cost.
-
Data Lakehouses add transactional features and performance optimization.
-
Medallion Architecture provides a clear progression from raw to business-ready data.
-
Bronze Layer (raw) → Silver Layer (cleaned) → Gold Layer (business-ready).
-
Apache Spark is the compute engine for most lakehouse implementations.
-
Delta Lake provides ACID transactions and time travel.
-
Unstructured data (emails, PDFs, audio) lives naturally in data lakes.
-
Real-time analytics is possible with streaming ingestion.
-
Compliance is easier with full data lineage and audit trails.
-
Cost savings come from using cheap object storage and open formats.
7.3 Recommended Next Steps
-
Experiment with Delta Lake using local Spark
-
Build a simple Medallion Architecture pipeline
-
Explore how your bank uses data lakes
-
Learn about cloud data platforms (AWS S3, Azure Data Lake)