π₯ Source Systems & The Ingestion Layer
Source Systems & The Ingestion Layer
ποΈ Supported Data Sources
Before processing data, it must be brought into the Lakehouse. Databricks acts as a massive sponge, capable of ingesting data from nearly any source system:
- Structured Data: RDBMS systems like Oracle, SQL Server, PostgreSQL, MySQL, and SAP.
- Semi-Structured Data: File-based data (JSON, XML, CSV, Log files).
- Unstructured Data: Complex data (Documents, Images, Videos, raw IoT sensor data).
- SaaS Applications: Enterprise COTS apps (Salesforce, Workday, ServiceNow, NetSuite).
- Streaming Systems: Event brokers (Apache Kafka, Azure Event Hubs, Confluent, Kinesis).
π Change Data Capture (CDC) via Lakeflow Connect
What Existed Previously: To keep a Data Warehouse up to date, engineers would run nightly batch jobs that performed a Full Data Reloadβpulling millions of rows from the source database, truncating the target table, and rewriting the entire dataset.
Problems Faced:
- Extreme Compute Costs: Re-reading and rewriting unchanged historical data over and over wastes massive amounts of compute power.
- Data Latency: Because full reloads are so heavy, they can only run during off-peak hours (nightly). Business reports are always 24 hours out of date.
How Present Technology Solves It (CDC): Databricks handles this natively using Change Data Capture (CDC) integrated via Lakeflow Connect. Instead of full reloads, CDC continuously tracks and extracts only the incremental changes (Inserts, Updates, Deletes) occurring in the transactional system.
- Lakeflow Connect provides highly scalable, point-and-click built-in connectors for databases and SaaS apps to stream these changes directly into your Bronze and Silver Delta layers.
- It automatically handles complex logic like upserts, out-of-order events, and SCD (Slowly Changing Dimensions) Type 2 tracking.
Real-World Analogy Mapping:
Imagine a global Order Management System. When an order's status changes from "Placed" β "Shipped" β "Delivered", CDC captures just that specific status update event and syncs it immediately to the Lakehouse. The dashboard reflects the "Delivered" status instantly, without reloading the 50 million older orders that haven't changed.
π File-Based Ingestion: Databricks Auto Loader
When raw data (like JSON or CSV logs) is continuously dumped into cloud object storage (AWS S3, Azure ADLS), we need a way to efficiently read it into Delta tables.
The Long Way:
# Traditional Spark Batch Read
df = spark.read.csv("s3://telecom-bucket/call-records/")
Why this is bad: Every time this job runs, Spark must scan the entire directory, looking at millions of old files just to figure out what was added since the last run.
The Smart Way (Databricks Auto Loader):
# Databricks Auto Loader (cloudFiles)
df = spark.readStream.format("cloudFiles") \
.option("cloudFiles.format", "csv") \
.load("s3://telecom-bucket/call-records/")
Why this is brilliant: Auto Loader (cloudFiles) acts as a highly optimized incremental ingestor. It keeps an internal tracking state of exactly which files have already been processed. When it runs, it reads only the new files that arrived since the last trigger.
Auto Loader Superpowers:
- Schema Inference & Evolution: It automatically infers the data schema and gracefully adapts if a new column is suddenly added to the raw files without breaking the pipeline.
- Cost Efficiency: Perfect for telecom or IoT scenarios where millions of tiny log files are dumped every hour.
π§ͺ Practice Drill
Q1. A database administrator complains that their Oracle database is heavily burdened by Databricks nightly data extraction jobs. What ingestion strategy should you switch to?
Q2. What is the fundamental problem with using standard spark.read.json() to ingest logs that are being continuously dumped into a cloud bucket?
Q3. What specific format keyword in PySpark triggers the Databricks Auto Loader?
π‘ Click for Solutions
A1. Change Data Capture (CDC) via Lakeflow Connect. It will only pull the incremental changes rather than burdening the database with a full nightly reload.
A2. Standard Spark reads must scan the entire directory tree to find new files, which becomes incredibly slow and expensive as the number of historical files grows.
A3. "cloudFiles" (e.g., spark.readStream.format("cloudFiles")).
β π Environment Setup & Workspace Provisioning | Next Topic β π The Lakehouse Storage Architecture