PySpark 12 — Thao tác Delta & Iceberg trong PySpark
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 Iceberg và Delta 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ạnh | Delta Lake | Iceberg |
|---|---|---|
| Cách nạp | package + extension, catalog spark_catalog | package + extension + khai báo catalog riêng |
| Đặt tên | db.table hoặc path | catalog.namespace.table (3 phần) |
| MERGE/UPDATE/DELETE | SQL và DeltaTable API (Python) | chủ yếu SQL (DataFrameV2 để ghi) |
| Time travel | versionAsOf / timestampAsOf | FOR VERSION AS OF (snapshot-id) / FOR TIMESTAMP AS OF |
| Partition | cột partition tường minh, ZORDER | hidden partitioning + partition evolution |
| Bảo trì | OPTIMIZE, ZORDER, VACUUM | rewrite_data_files, expire_snapshots (procedures) |
| Hệ sinh thái | mạnh nhất trên Databricks/Spark | trung 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ộtUPDATEđóng dòng hiện tại (setvalid_to,is_current=false), bước haiINSERTdòng mớiis_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:
- Đọc change log ngày từ bronze, khử trùng theo
acct_idgiữ bảntsmới nhất. 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ở.OPTIMIZE ZORDER BY (acct_id)cuối tuần,VACUUM RETAIN 720 HOURSgiữ 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; IcebergFOR 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; Icebergrewrite_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
- Delta Lake Documentation — mục Table Deletes, Updates, and Merges; Table Utility Commands (OPTIMIZE, VACUUM, Z-ORDER)
- Delta Lake — Table Batch Reads and Writes — mergeSchema, time travel (versionAsOf / timestampAsOf)
- Apache Iceberg Documentation — mục Spark Writes (MERGE INTO), Spark Queries (time travel), Maintenance Procedures
- Apache Iceberg — Partitioning — hidden partitioning và partition evolution
- Apache Spark — Structured Streaming Programming Guide — foreachBatch và checkpointing
- Delta Lake — Table Streaming Reads and Writes — foreachBatch + MERGE cho upsert streaming
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.
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.
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.
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.
Cảm nhận của bạn
Bình luận
Chưa có bình luận. Hãy là người đầu tiên chia sẻ!