PySpark 12 — Thao tác Delta & Iceberg trong PySpark

14 thg 7, 2026 4 lượt xem
#data-engineering
#lakehouse
#delta-lake
#pyspark
#iceberg

Vì sao lakehouse cần bảng giao dịch, không chỉ Parquet trần

Suốt series này, khi ghi kết quả xuống kho phân tích ta hay nói "ghi Delta/Iceberg" mà chưa mổ xẻ. Bài này bù đúng góc đó: dùng PySpark thao tác bảng lakehouse có giao dịch (ACID) — bổ sung góc Python cho ba bài đọc kèm là Delta lakehouse trên Spark, tổng quan IcebergDelta Lake trên Databricks.

Vấn đề với Parquet trần (thư mục Parquet không có lớp giao dịch) lộ ra ngay khi dữ liệu cần thay đổi. Parquet là file bất biến: muốn sửa một dòng phải đọc cả partition, ghi lại toàn bộ. Không có khái niệm "cập nhật một khách hàng". Tệ hơn:

  • Không upsert/xoá/cập nhật: nghiệp vụ ngân hàng luôn có dữ liệu thay đổi — số dư biến động, khách đổi địa chỉ, giao dịch bị đảo (reversal). Với Parquet trần bạn phải tự viết logic overwrite partition, rất dễ sai và không nguyên tử.
  • Không nguyên tử (atomicity): nếu job ghi nửa chừng rồi chết, thư mục còn file rác, người đọc thấy dữ liệu dở dang. Không có commit "được ăn cả, ngã về không".
  • Không time travel: không đọc lại được trạng thái bảng tại thời điểm quá khứ — trong ngân hàng đây là yêu cầu audit bắt buộc, không phải tính năng cho vui.
  • Schema cứng: thêm một cột mới vào luồng đang chạy là cơn ác mộng thủ công.

Bảng giao dịch (Delta Lake, Apache Iceberg, Hudi) giải quyết bằng một transaction log (Delta: thư mục _delta_log; Iceberg: các file metadata + manifest) ghi lại từng commit. Nhờ đó có ACID: mỗi lần ghi là một version mới, nguyên tử, có thể đọc lại version cũ, và hỗ trợ MERGE/UPDATE/DELETE chuẩn SQL. Đó là điều biến "data lake" thành "lakehouse" (xem thêm kiến trúc lakehouse).

Delta Lake trong PySpark

Cấu hình

Delta không nằm sẵn trong Spark OSS, phải nạp package và bật hai extension. Khi chạy spark-submit/pyspark:

spark-submit \
  --packages io.delta:delta-spark_2.12:3.2.0 \
  --conf spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension \
  --conf spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog \
  job.py

Hoặc dựng SparkSession bằng helper của Delta (Python — minh hoạ):

from delta import configure_spark_with_delta_pip
from pyspark.sql import SparkSession

builder = (SparkSession.builder.appName("delta-demo")
    .config("spark.sql.extensions",
            "io.delta.sql.DeltaSparkSessionExtension")
    .config("spark.sql.catalog.spark_catalog",
            "org.apache.spark.sql.delta.catalog.DeltaCatalog"))
spark = configure_spark_with_delta_pip(builder).getOrCreate()

Chú ý phiên bản delta-spark phải khớp với phiên bản Spark (Delta 3.x đi với Spark 3.5.x). Sai cặp version là lỗi hay gặp nhất — nối tiếp chủ đề cấu hình & triển khai.

Ghi và đọc format "delta"

# ghi
(df.write.format("delta").mode("overwrite")
   .save("s3a://lake/silver/customer"))

# đọc theo path
d = spark.read.format("delta").load("s3a://lake/silver/customer")

# hoặc quản lý qua catalog (managed/external table)
df.write.format("delta").saveAsTable("silver.customer")
spark.table("silver.customer")

Xem thêm cách chọn nguồn/định dạng ở đọc/ghi nguồn dữ liệu.

MERGE INTO — upsert, trái tim của CDC

MERGE khớp bảng đích với dữ liệu nguồn theo một điều kiện, rồi tuỳ trạng thái mà UPDATE, DELETE hoặc INSERT. Đây là toán tử quan trọng nhất khi làm CDC. Delta cho hai đường: SQL hoặc DeltaTable API (Python).

DeltaTable API (Python — minh hoạ):

from delta.tables import DeltaTable

tgt = DeltaTable.forName(spark, "silver.customer")
(tgt.alias("t")
   .merge(source=changes.alias("s"), condition="t.cust_id = s.cust_id")
   .whenMatchedDelete(condition="s.op = 'D'")          # bản ghi bị xoá ở nguồn
   .whenMatchedUpdateAll(condition="s.op != 'D'")      # update toàn cột
   .whenNotMatchedInsertAll(condition="s.op != 'D'")   # bản ghi mới
   .execute())

Cùng logic viết bằng SQL:

MERGE INTO silver.customer t
USING changes s ON t.cust_id = s.cust_id
WHEN MATCHED AND s.op = 'D' THEN DELETE
WHEN MATCHED AND s.op <> 'D' THEN UPDATE SET *
WHEN NOT MATCHED AND s.op <> 'D' THEN INSERT *

Lưu ý quan trọng về đúng đắn: nếu một khoá xuất hiện nhiều lần trong changes mà khớp một dòng đích, MERGE báo lỗi (không xác định lấy bản nào). Vì thế phải khử trùng nguồn trước — giữ đúng một bản mới nhất mỗi khoá (xem phần idempotent bên dưới).

UPDATE / DELETE

tgt.update(condition="status = 'DORMANT'",
           set={"risk_flag": "'REVIEW'"})
tgt.delete("closed_at < '2020-01-01'")

Hoặc UPDATE silver.customer SET ... WHERE ... / DELETE FROM ... bằng SQL. Đây là DML thật trên bảng phân tích — thứ Parquet trần không có.

Time travel

Mỗi commit tạo một version tăng dần. Đọc lại quá khứ bằng versionAsOf hoặc timestampAsOf (Python — minh hoạ):

v3 = spark.read.format("delta").option("versionAsOf", 3) \
        .table("silver.customer")
snap = spark.read.format("delta") \
        .option("timestampAsOf", "2026-06-30 23:59:59") \
        .table("silver.customer")

# xem lịch sử commit
DeltaTable.forName(spark, "silver.customer").history().show()

Time travel dựa vào file còn tồn tại — nếu đã VACUUM xoá file cũ thì không lùi quá xa được nữa (mặc định giữ 7 ngày).

Schema evolution

Khi nguồn thêm cột, bật mergeSchema để bảng tự nới schema thay vì văng lỗi:

(df.write.format("delta").mode("append")
   .option("mergeSchema", "true")
   .saveAsTable("silver.customer"))

Với MERGE, bật spark.databricks.delta.schema.autoMerge.enabled=true để whenMatchedUpdateAll/insertAll tự thêm cột mới. Dùng có kiểm soát — schema tự nới không kiểm soát dễ nuốt phải cột rác.

OPTIMIZE, Z-ORDER, VACUUM

MERGE và streaming sinh rất nhiều file nhỏ (small files problem) làm chậm đọc. Bảo trì định kỳ:

OPTIMIZE silver.customer;                       -- gộp file nhỏ (bin-packing)
OPTIMIZE silver.customer ZORDER BY (cust_id);   -- gom dữ liệu theo cột lọc
VACUUM  silver.customer RETAIN 168 HOURS;       -- xoá file mồ côi > 7 ngày

OPTIMIZE gộp file; ZORDER sắp xếp đồng địa phương theo cột hay lọc để data skipping hiệu quả hơn (bổ trợ tối ưu & debug nếu có trong series); VACUUM dọn file không còn được version nào tham chiếu — nhưng chạy VACUUM sẽ cắt khả năng time travel về trước mốc giữ lại, cân nhắc kỹ với bảng cần audit.

Iceberg trong PySpark

Cấu hình catalog

Iceberg tổ chức quanh catalog. Khai báo một catalog tên (ví dụ lake) trỏ tới kho metadata và kho lưu trữ:

spark-submit \
  --packages org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.6.0 \
  --conf spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtension \
  --conf spark.sql.catalog.lake=org.apache.iceberg.spark.SparkCatalog \
  --conf spark.sql.catalog.lake.type=hive \
  --conf spark.sql.catalog.lake.warehouse=s3a://lake/warehouse \
  job.py

Tên bảng khi đó có 3 phần: lake.silver.customer (catalog.namespace.table). Iceberg hỗ trợ nhiều loại catalog (Hive, REST, JDBC, Nessie, Glue) — chi tiết ở catalog & engine của Iceberg.

Tạo, đọc, MERGE/UPDATE/DELETE

spark.sql("""
  CREATE TABLE lake.silver.customer (
    cust_id BIGINT, full_name STRING, city STRING,
    balance DECIMAL(18,2), updated_at TIMESTAMP)
  USING iceberg
  PARTITIONED BY (days(updated_at))
""")

changes.writeTo("lake.silver.customer").append()   # DataFrameV2 API
spark.table("lake.silver.customer")

MERGE/UPDATE/DELETE bằng SQL — cú pháp gần như y hệt Delta:

MERGE INTO lake.silver.customer t
USING changes s ON t.cust_id = s.cust_id
WHEN MATCHED AND s.op = 'D' THEN DELETE
WHEN MATCHED THEN UPDATE SET *
WHEN NOT MATCHED THEN INSERT *

Time travel

Iceberg dùng snapshot (mỗi commit là một snapshot có ID). Cú pháp SQL chuẩn:

SELECT * FROM lake.silver.customer FOR VERSION AS OF 3921089283028;  -- snapshot-id
SELECT * FROM lake.silver.customer FOR TIMESTAMP AS OF '2026-06-30 23:59:59';

Trong DataFrame API dùng option snapshot-id hoặc as-of-timestamp khi spark.read.

Hidden partitioning

Điểm mạnh riêng của Iceberg: hidden partitioning. Bạn khai PARTITIONED BY (days(updated_at)), Iceberg tự tính giá trị partition từ cột — người viết query không cần biết cột partition, chỉ cần lọc WHERE updated_at >= ... và Iceberg tự prune. Khác Hive (phải có cột partition riêng và tự điền đúng), tránh cả lớp lỗi partition sai. Còn cho phép partition evolution — đổi cách chia partition mà không viết lại dữ liệu cũ (xem tổng quan Iceberg).

Bảo trì qua Spark procedures

Iceberg cũng có small files và snapshot tích tụ. Bảo trì gọi bằng stored procedure trên catalog:

CALL lake.system.rewrite_data_files(table => 'silver.customer');
CALL lake.system.expire_snapshots(
       table => 'silver.customer',
       older_than => TIMESTAMP '2026-07-07 00:00:00');
CALL lake.system.remove_orphan_files(table => 'silver.customer');

rewrite_data_files tương đương OPTIMIZE (gộp/sắp xếp lại file); expire_snapshots là "VACUUM của Iceberg" — xoá snapshot cũ và cắt time travel về trước mốc đó; remove_orphan_files dọn file không được metadata nào tham chiếu.

Delta vs Iceberg khi dùng từ PySpark

Khía cạnhDelta LakeIceberg
Cách nạppackage + extension, catalog spark_catalogpackage + extension + khai báo catalog riêng
Đặt têndb.table hoặc pathcatalog.namespace.table (3 phần)
MERGE/UPDATE/DELETESQL DeltaTable API (Python)chủ yếu SQL (DataFrameV2 để ghi)
Time travelversionAsOf / timestampAsOfFOR VERSION AS OF (snapshot-id) / FOR TIMESTAMP AS OF
Partitioncột partition tường minh, ZORDERhidden partitioning + partition evolution
Bảo trìOPTIMIZE, ZORDER, VACUUMrewrite_data_files, expire_snapshots (procedures)
Hệ sinh tháimạnh nhất trên Databricks/Sparktrung lập engine (Spark, Trino, Flink, BigQuery)

Khi nào chọn gì (từ góc PySpark):

  • Stack xoay quanh Databricks/Spark thuần, muốn API Python trực tiếp (DeltaTable) và ZORDER → Delta.
  • Cần nhiều engine đọc chung một bảng (Spark ghi, Trino/Flink/BigQuery đọc), muốn partition evolution và hidden partitioning → Iceberg.
  • Cả hai đều ACID, đều MERGE/time travel tốt. Ở NCB lựa chọn thường do engine ecosystem và định hướng catalog trung lập quyết định, không phải do thiếu tính năng.

Mẫu thực chiến

CDC upsert vào silver

Luồng phổ biến nhất: nguồn core banking phát change events (qua Debezium → Kafka), mỗi event có op (c=insert, u=update, d=delete), khoá nghiệp vụ, và timestamp. Ta MERGE chúng vào bảng silver.

Xử lý trùng + idempotent write

Trước khi MERGE, khử trùng: một micro-batch có thể chứa nhiều thay đổi cho cùng một khoá; chỉ giữ bản mới nhất theo updated_at (Python — minh hoạ):

from pyspark.sql import Window, functions as F

w = Window.partitionBy("cust_id").orderBy(F.col("updated_at").desc())
latest = (changes
          .withColumn("rn", F.row_number().over(w))
          .filter("rn = 1").drop("rn"))

MERGE vốn idempotent theo khoá: chạy lại cùng batch cho ra cùng kết quả (khác với append — chạy lại là nhân đôi dữ liệu). Đó là lý do CDC luôn dùng MERGE thay vì append vào silver.

SCD (Slowly Changing Dimension)

  • SCD Type 1 — ghi đè: chính là whenMatchedUpdateAll ở trên, chỉ giữ giá trị hiện tại.
  • SCD Type 2 — giữ lịch sử: mỗi thay đổi mở một dòng version mới với valid_from/valid_to/is_current. Làm bằng MERGE hai bước — bước một UPDATE đóng dòng hiện tại (set valid_to, is_current=false), bước hai INSERT dòng mới is_current=true. Delta/Iceberg đều làm được vì có UPDATE + INSERT nguyên tử.

Kết hợp streaming: foreachBatch + MERGE

Structured Streaming không MERGE trực tiếp vào sink được, nhưng foreachBatch cho ta một DataFrame tĩnh mỗi micro-batch để chạy MERGE như batch (Python — minh hoạ):

def upsert_batch(batch_df, batch_id):
    w = Window.partitionBy("cust_id").orderBy(F.col("updated_at").desc())
    latest = (batch_df.withColumn("rn", F.row_number().over(w))
                      .filter("rn = 1").drop("rn"))
    tgt = DeltaTable.forName(spark, "silver.customer")
    (tgt.alias("t")
        .merge(latest.alias("s"), "t.cust_id = s.cust_id")
        .whenMatchedDelete(condition="s.op = 'd'")
        .whenMatchedUpdateAll(condition="s.op <> 'd'")
        .whenNotMatchedInsertAll(condition="s.op <> 'd'")
        .execute())

(spark.readStream.format("kafka").option(...).load()
   .selectExpr(...)                      # parse Debezium payload
   .writeStream
   .foreachBatch(upsert_batch)
   .option("checkpointLocation", "s3a://lake/_ckpt/customer")
   .trigger(processingTime="1 minute")
   .start())

checkpointLocation + MERGE cho exactly-once về mặt trạng thái bảng: nếu batch chạy lại sau sự cố, MERGE theo khoá không tạo bản trùng. Cách test luồng streaming này xem streaming & testing.

Ví dụ ghép: tạo → MERGE → time travel → OPTIMIZE

Trình tự end-to-end một chu kỳ CDC ngày (Python — minh hoạ):

from delta.tables import DeltaTable
from pyspark.sql import functions as F, Window

# 1) tạo bảng Delta lần đầu (bootstrap từ full snapshot)
(seed_df.write.format("delta").mode("overwrite")
   .saveAsTable("silver.account_balance"))

# 2) MERGE upsert từ batch CDC hằng ngày
w = Window.partitionBy("acct_id").orderBy(F.col("ts").desc())
latest = (cdc_df.withColumn("rn", F.row_number().over(w))
                .filter("rn = 1").drop("rn"))

tgt = DeltaTable.forName(spark, "silver.account_balance")
(tgt.alias("t")
    .merge(latest.alias("s"), "t.acct_id = s.acct_id")
    .whenMatchedDelete(condition="s.op = 'd'")
    .whenMatchedUpdateAll(condition="s.op <> 'd'")
    .whenNotMatchedInsertAll(condition="s.op <> 'd'")
    .execute())

# 3) time travel: đọc lại trạng thái cuối tháng để tái tạo báo cáo
snap = (spark.read.format("delta")
        .option("timestampAsOf", "2026-06-30 23:59:59")
        .table("silver.account_balance"))

# 4) bảo trì cuối chu kỳ
spark.sql("OPTIMIZE silver.account_balance ZORDER BY (acct_id)")
spark.sql("VACUUM silver.account_balance RETAIN 720 HOURS")  # giữ 30 ngày cho audit

Use case thực tế

Bối cảnh NCB. Đội dữ liệu duy trì bảng silver silver.account_balance — số dư và trạng thái của khoảng 6 triệu tài khoản khách hàng — cùng silver.customer hồ sơ khách. Trước đây các bảng này được overwrite Parquet toàn phần mỗi đêm: đọc full dump core banking cỡ ~40–60 GB, ghi đè sạch. Cách đó tốn khoảng 35–40 phút mỗi đêm, không có time travel, và mỗi khi kiểm toán hỏi "số dư tài khoản X ngày 30/6 là bao nhiêu?" thì không ai trả lời được nếu không phục hồi backup.

Chuyển sang Delta + CDC MERGE. Core banking bật CDC qua Debezium → Kafka; mỗi ngày chỉ khoảng 1,5–2 triệu bản ghi thay đổi (số dư biến động, mở/đóng tài khoản) thay vì đọc lại toàn bộ 6 triệu. Job đêm:

  1. Đọc change log ngày từ bronze, khử trùng theo acct_id giữ bản ts mới nhất.
  2. MERGE INTO silver.account_balance: matched-delete cho tài khoản đóng, matched-update cho số dư mới, not-matched-insert cho tài khoản mới mở.
  3. OPTIMIZE ZORDER BY (acct_id) cuối tuần, VACUUM RETAIN 720 HOURS giữ 30 ngày time travel.

Kết quả ước lượng. Thời gian job giảm còn khoảng 8–12 phút (chỉ xử lý delta, không ghi lại 6 triệu dòng). Quan trọng hơn về nghiệp vụ: khi kiểm toán hoặc đội rủi ro cần tái tạo báo cáo cuối kỳ, chỉ cần đọc timestampAsOf '2026-06-30 23:59:59' là ra đúng ảnh chụp bảng lúc đó — phục vụ audit và đối soát mà không cần restore backup. Bảng silver.customer áp dụng SCD Type 2 để lưu lịch sử đổi địa chỉ/nhóm rủi ro, phục vụ điều tra AML (xem tổng quan AML) và chất lượng dữ liệu. Ghép mọi mảnh này thành pipeline hoàn chỉnh là nội dung ETL ngân hàng end-to-end.

Ghi nhớ

  • Parquet trần không có ACID: không upsert/xoá/cập nhật nguyên tử, không time travel, không schema evolution. Bảng giao dịch (Delta/Iceberg) mới biến data lake thành lakehouse.
  • MERGE INTO là toán tử cốt lõi của CDC: matched-update/delete, not-matched-insert. Delta có cả DeltaTable API (Python) lẫn SQL; Iceberg chủ yếu SQL. Luôn dùng MERGE (idempotent theo khoá) thay vì append vào silver.
  • Khử trùng nguồn trước MERGE: một khoá khớp nhiều dòng nguồn sẽ làm MERGE lỗi; dùng row_number() giữ bản mới nhất mỗi khoá.
  • Time travel: Delta versionAsOf/timestampAsOf; Iceberg FOR VERSION AS OF (snapshot-id) / FOR TIMESTAMP AS OF. Bắt buộc cho audit ngân hàng.
  • VACUUM / expire_snapshots cắt time travel về trước mốc giữ lại — chỉnh retention theo yêu cầu audit trước khi dọn file.
  • Bảo trì file nhỏ: Delta OPTIMIZE/ZORDER; Iceberg rewrite_data_files. MERGE và streaming sinh nhiều file nhỏ nên phải chạy định kỳ.
  • Iceberg khác biệt ở hidden partitioning + partition evolution và tính trung lập engine; Delta mạnh ở API Python và hệ Databricks/Spark. Chọn theo ecosystem, không phải theo thiếu tính năng.
  • Streaming upsert = foreachBatch + MERGE + checkpoint → exactly-once về trạng thái bảng, chạy lại batch không nhân đôi dữ liệu.

Nguồn tham khảo

Bài viết liên quan

So sánh các định dạng dữ liệu (CSV, JSON, XML, Avro, Parquet, ORC) và lý do lưu theo cột nhanh hơn cho phân tích. Bài đi sâu vào row vs columnar storage, nén (Snappy/gzip/zstd), schema evolution, OLTP vs OLAP, object storage và partitioning để tối ưu chi phí lẫn tốc độ truy vấn.

13 thg 7, 2026 10

Data Engineering là ngành xây dựng và vận hành hệ thống biến dữ liệu thô thành dữ liệu sạch, tin cậy, sẵn sàng cho phân tích và AI. Bài giới thiệu vai trò Data Engineer trong vòng đời dữ liệu (nguồn → ingestion → storage → transformation → serving), phân biệt với Analyst/Scientist/ML Engineer, bức tranh hệ sinh thái công cụ và bài toán đưa dữ liệu core banking sang kho phân tích.

13 thg 7, 2026 9

Stream processing là gì, khác biệt batch vs stream (bounded/unbounded), micro-batch (Spark) vs true streaming, Apache Flink là gì và định vị so với Spark Structured Streaming và Kafka Streams. Kiến trúc runtime JobManager/TaskManager, các tầng API, triết lý streaming-first và bối cảnh phát hiện gian lận ngân hàng.

13 thg 7, 2026 8

Vì sao một máy không đủ và cần xử lý phân tán: từ MapReduce, Hadoop/HDFS đến Apache Spark in-memory. Kiến trúc driver–executor–cluster manager, các mức trừu tượng RDD/DataFrame/Dataset, cơ chế lazy evaluation với DAG, và vì sao shuffle (wide dependency) là phần tốn kém nhất. Kèm PySpark, Spark SQL và các kỹ thuật tối ưu (partition, broadcast join, cache, chống skew) cùng khi nào KHÔNG nên dùng Spark.

13 thg 7, 2026 7

Cảm nhận của bạn

Bình luận

Bạn cần để viết bình luận.

Chưa có bình luận. Hãy là người đầu tiên chia sẻ!