Chuỗi sản xuất với workflows

Chuyển đổi dữ liệu với Spark SQL trong Databricks

Disha Mukherjee

Lead Data Engineer

Vì sao dùng Delta Lake?

Kho dữ liệu số phát sáng an toàn với các bảng dữ liệu được sắp xếp và một lá chắn bảo vệ, phong cách phẳng hiện đại

 

$$

  • Giao dịch ACID → hoàn tác ghi lỗi
  • Cưỡng chế schema → chặn kiểu không khớp
  • Phiên bản hóa → truy vấn trạng thái trước đó
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Ghi vào Delta

df_valid.write.format("delta") \
    .mode("overwrite") \
    .saveAsTable("transactions_clean")

print(f"Rows written: {df_valid.count():,}")
Rows written: 33,223

$$

$$

  • Một bảng Delta mới xuất hiện trong Unity Catalog
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Nhiệm vụ trong notebook

task1_ingest: nạp và làm sạch

  • Nạp CSV → làm sạch → ghi thành bảng Delta
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Nhiệm vụ trong notebook

task2_metrics: doanh thu theo danh mục

  • Đọc bảng đã làm sạch → tính chỉ số → ghi ra bảng mới
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Nhiệm vụ trong notebook

task3_customers: xếp hạng theo chi tiêu

  • Đọc bảng đã làm sạch → xếp hạng khách hàng → lưu ra bảng mới
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Tạo job

Giao diện Jobs và Pipelines hiển thị DAG job ba nhiệm vụ đã hoàn tất với mũi tên phụ thuộc từ task1_ingest đến task2_metrics đến task3_customers

  • Có thể chạy thủ công, lên lịch và đặt trigger
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Chạy job

DAG chạy job với task1_ingest và task2_metrics xanh (thành công) và task3_customers đỏ (thất bại)

Chuyển đổi dữ liệu với Spark SQL trong Databricks

Chạy job

Kiểm tra lỗi ở nhiệm vụ ba

  • Lỗi: customer-id phải là Customer_ID
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Chạy job

Chế độ xem đồ thị với trạng thái thành công

Chế độ xem dòng thời gian

Chuyển đổi dữ liệu với Spark SQL trong Databricks

Lakeflow là gì?

$$

$$

so sánh: Mệnh lệnh (Jobs) | Khai báo (Lakeflow)

 

  • Jobs → bạn quản lý từng bước
  • Lakeflow → khai báo nội dung các bảng
  • Databricks xử lý thứ tự, thử lại và tài nguyên tính toán
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Mẫu @dlt.table

@dlt.table(name="transactions_bronze")
def transactions_bronze():
    return spark.read.format("csv").schema(schema).load(FILE_PATH)

@dlt.table(name="transactions_silver") def transactions_silver(): return dlt.read("transactions_bronze").na.drop(...).filter(...)
@dlt.table(name="category_revenue_gold") def category_revenue_gold(): return dlt.read("transactions_silver").groupBy("Category").agg(...)
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Chạy pipeline

DAG pipeline Lakeflow hiển thị transactions_bronze 100K dòng rồi transactions_silver 33K dòng rồi category_revenue_gold 6 dòng, tất cả materialized view màu xanh

  • Bronze (100K dòng) → Silver (33K dòng) → Gold (6 dòng)
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Chọn Notebooks, Jobs hay Lakeflow?

$$

các lớp: Notebooks, Databricks Jobs, Lakeflow Pipelines

 

$$

  • Notebooks → khám phá, dựng thử
  • Databricks Jobs → pipeline nhiều bước, có lịch chạy
  • Lakeflow → quản lý trọn vẹn, khai báo
Chuyển đổi dữ liệu với Spark SQL trong Databricks

Cùng luyện tập nào!

Chuyển đổi dữ liệu với Spark SQL trong Databricks

Preparing Video For Download...