Databricks for Data & AI Engineers
Trang chủ / Module 6
📥 Module 6 / 14 · 2 tuần

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.

Auto LoaderDLTStreaming
📌 Vì sao module này quan trọng

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

Pipeline Lakeflow hoàn chỉnh Bronze→Silver→Gold

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

✅ 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