Data Ingestion — Auto Loader & Lakeflow Declarative Pipelines
Auto Loader (cloudFiles), COPY INTO, Lakeflow Declarative Pipelines (trước đây là Delta Live Tables), data quality expectations.
Nạp dữ liệu đúng cách — incremental, đáng tin cậy, tự phục hồi khi lỗi — là nền móng của mọi pipeline Bronze/Silver/Gold. Đây là công việc chiếm nhiều thời gian nhất của Data Engineer trong thực tế, và Databricks cung cấp bộ công cụ chuyên biệt (Auto Loader, Lakeflow) giải quyết tốt hơn nhiều so với tự viết pipeline nạp dữ liệu thủ công.
🎯 Mục tiêu học tập
- Sử dụng Auto Loader để nạp dữ liệu incremental từ cloud storage một cách đáng tin cậy
- Phân biệt COPY INTO và Auto Loader, biết khi nào dùng loại nào
- Xây dựng pipeline khai báo bằng Lakeflow Declarative Pipelines (DLT) kèm data quality expectations
- Hiểu sự khác nhau giữa streaming table và materialized view trong Lakeflow
Auto Loader (cloudFiles)
Auto Loader tự động phát hiện file mới xuất hiện trong cloud storage (S3/ADLS/GCS) và nạp incremental — không cần tự viết logic theo dõi "file nào đã xử lý rồi". Nó dùng cơ chế thông báo sự kiện của cloud (event notification) thay vì liệt kê lại toàn bộ thư mục mỗi lần, nên hiệu quả kể cả khi thư mục nguồn có hàng triệu file.
df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", "/mnt/schema/orders")
.load("/mnt/raw/orders/"))
(df.writeStream
.option("checkpointLocation", "/mnt/checkpoints/orders")
.trigger(availableNow=True) # chạy 1 lần xử lý hết dữ liệu mới rồi dừng — phù hợp job theo lịch
.toTable("bronze.raw_orders"))
trigger(availableNow=True) là mẫu phổ biến nhất trên Databricks: viết code streaming
nhưng chạy theo lịch định kỳ (không cần cluster chạy 24/7 chờ dữ liệu) — kết hợp được sự đơn giản của
streaming API với chi phí của batch job.
COPY INTO — khi nào dùng thay Auto Loader
COPY INTO là lệnh SQL đơn giản hơn, phù hợp khi số lượng file nguồn không quá lớn (hàng
nghìn, không phải hàng triệu) và không cần độ trễ thấp — dễ viết, dễ hiểu cho người quen SQL hơn là thiết
lập streaming. Auto Loader vượt trội khi khối lượng file lớn, cần độ trễ thấp, hoặc cần xử lý schema
evolution phức tạp tự động.
Lakeflow Declarative Pipelines (trước đây gọi là Delta Live Tables)
Thay vì viết từng bước ETL rời rạc và tự quản lý thứ tự chạy, Lakeflow cho phép khai báo bảng đích và Databricks tự suy ra đồ thị phụ thuộc, tự quản lý thứ tự thực thi, tự xử lý retry, và tự động áp dụng incremental processing khi có thể:
import dlt
from pyspark.sql import functions as F
@dlt.table(comment="Dữ liệu đơn hàng thô, nạp từ cloud storage")
def bronze_orders():
return (spark.readStream.format("cloudFiles")
.option("cloudFiles.format", "json")
.load("/mnt/raw/orders/"))
@dlt.table(comment="Đơn hàng đã làm sạch")
@dlt.expect_or_drop("valid_amount", "amount > 0") # data quality constraint
@dlt.expect("has_customer", "customer_id IS NOT NULL") # cảnh báo nhưng không drop
def silver_orders():
return (dlt.read_stream("bronze_orders")
.withColumn("order_date", F.to_date("created_at")))
Expectations (@dlt.expect*) là tính năng nổi bật nhất: định nghĩa ràng
buộc chất lượng dữ liệu ngay trong code pipeline — expect_or_drop tự loại bỏ dòng vi phạm,
expect_or_fail dừng cả pipeline nếu vi phạm, expect chỉ ghi nhận số liệu để theo
dõi. Toàn bộ số liệu chất lượng dữ liệu được Databricks tự động hiển thị trên UI giám sát pipeline.
Streaming Table vs Materialized View
Trong Lakeflow, streaming table chỉ xử lý dữ liệu mới kể từ lần chạy trước (incremental thuần) — phù hợp cho tầng Bronze/Silver nạp liên tục. Materialized view tính toán lại toàn bộ hoặc một phần kết quả dựa trên input hiện tại, tối ưu cho các phép tổng hợp phức tạp (Gold layer) mà Databricks tự quyết định có thể tính incremental được phần nào để tránh tính lại toàn bộ mỗi lần nếu có thể.
🏋️ Bài tập thực hành
Xây 1 pipeline Lakeflow Declarative Pipelines đầy đủ 3 tầng: Bronze (Auto Loader nạp file JSON/CSV mẫu), Silver (làm sạch + 2 expectations kiểm tra chất lượng), Gold (materialized view tổng hợp theo ngày). Chạy pipeline, cố tình đưa vào một số dòng dữ liệu lỗi (amount âm, thiếu customer_id) để quan sát cách expectations xử lý và số liệu hiển thị trên UI giám sát.
📚 Tài nguyên học tập
-
Databricks — Auto LoaderTài liệu chính thức Auto LoaderDocs
-
Databricks — Lakeflow Declarative PipelinesTài liệu chính thức Lakeflow (DLT)Docs
-
Databricks — Expectations for data qualityHướng dẫn data quality expectationsDocs
✅ Tự đánh giá hoàn thành
- Viết được pipeline Auto Loader nạp incremental có checkpoint đúng cách
- Phân biệt rõ khi nào dùng COPY INTO, khi nào dùng Auto Loader
- Xây được pipeline Lakeflow 3 tầng với ít nhất 2 expectations
- Giải thích được sự khác nhau giữa streaming table và materialized view