PySpark 8 — ETL ngân hàng end-to-end & Tối ưu
Ghép mọi thứ lại: một pipeline ETL ngân hàng thật
Bảy bài trước của series đã tách từng mảnh: tổng quan PySpark, DataFrame API, đọc/ghi nguồn dữ liệu, UDF & pandas_udf, pandas API on Spark, MLlib, và streaming + testing. Bài kết này ghép tất cả thành một pipeline ETL ngân hàng end-to-end chạy hằng ngày trên cluster — thứ mà một data engineer NCB thực sự phải xây và vận hành.
Bài toán rất cụ thể. Mỗi đêm, hệ thống nhận dữ liệu từ nhiều nguồn: core banking (tài khoản, số dư, sổ cái), hệ thống thẻ (giao dịch POS/ATM/e-commerce), các kênh (Internet Banking, Mobile, ví liên kết). Khối lượng cỡ hàng trăm triệu bản ghi/ngày. Nhiệm vụ: nạp dữ liệu thô, làm sạch và chuẩn hoá, làm giàu và tổng hợp, rồi ghi vào kho phân tích (lakehouse) để sáng hôm sau các đội BI, rủi ro, marketing và mô hình ML dùng được ngay. Job phải xong trong cửa sổ đêm, phải tái chạy được khi lỗi, và không được làm hỏng số liệu đã có.
Chìa khoá kiến trúc là medallion (bronze → silver → gold): chia pipeline thành ba tầng chất lượng tăng dần, mỗi tầng là một bảng vật lý trong lakehouse. Cách chia này không phải trang trí — nó cho phép tái chạy từng tầng độc lập, debug rõ ràng, và tách trách nhiệm giữa "nạp thô" với "làm sạch" với "phục vụ".
Bronze — nạp thô, giữ nguyên hiện trạng
Tầng bronze chỉ có một việc: đưa dữ liệu từ nguồn vào lakehouse càng gần nguyên bản càng tốt, không biến đổi logic nghiệp vụ. Đọc từ file landing (CSV/JSON từ core), từ Parquet do hệ thống thẻ đẩy, hoặc từ Kafka/JDBC — chi tiết ở đọc/ghi nguồn. Quy tắc bronze:
- Không sửa dữ liệu, chỉ thêm metadata: cột
ingest_ts(thời điểm nạp),source_file,batch_date(ngày chạy). - Append theo ngày, phân vùng (partition) theo
batch_dateđể dễ tái chạy đúng một ngày. - Giữ cả bản ghi lỗi — bronze là "nguồn sự thật thô", để còn truy vết được khi số liệu tầng trên bị nghi ngờ.
Bronze đóng vai trò checkpoint an toàn: nếu file nguồn bị xoá, ta vẫn còn bản thô đã nạp. Nó cũng tách việc "kéo dữ liệu về" (I/O, dễ lỗi mạng) khỏi việc "xử lý" (tính toán nặng), giúp tái chạy tầng sau mà không phải kéo lại từ nguồn.
Silver — làm sạch, chuẩn hoá, validate
Tầng silver là nơi dữ liệu trở nên đáng tin và dùng được. Toàn bộ dùng DataFrame API với hàm built-in F.* (tránh Python UDF — lý do ở bài 4). Các phép chính:
- Ép kiểu & chuẩn hoá: số tiền về
decimal, ngày vềdate/timestampđúng múi giờ, mã tiền tệ về ISO (VND,USD). - Chuẩn hoá giá trị: trim khoảng trắng, viết hoa mã, ánh xạ mã kênh về danh mục chuẩn, chuẩn hoá số tài khoản.
- Khử trùng lặp (dedup): một giao dịch có thể đến hai lần do retry ở nguồn; dùng
dropDuplicatestheo khoá nghiệp vụ hoặcrow_number()giữ bản mới nhất. - Validate chất lượng: kiểm tra các luật ở data quality — không null khoá chính, số dư không âm với sản phẩm không cho thấu chi, ngày giao dịch ≤ ngày chạy. Bản ghi vi phạm được tách sang bảng quarantine kèm lý do, thay vì âm thầm bị bỏ.
Silver nên là bảng sạch, đã dedup, đúng schema, ở mức chi tiết gốc (grain giao dịch/tài khoản). Đây là tầng mà phần lớn truy vấn ad-hoc và mô hình sẽ đọc.
Gold — tổng hợp phục vụ BI & ML
Tầng gold biến silver thành các bảng phục vụ đích cụ thể: bảng tổng hợp cho dashboard, và bảng đặc trưng (feature) cho MLlib. Đây là nơi join nhiều bảng và aggregate — cũng là nơi tốn shuffle nhất, nên tối ưu tập trung ở đây.
Ví dụ gold điển hình ở NCB: bảng daily_customer_360 gộp số dư cuối ngày, tổng chi tiêu thẻ, số giao dịch theo kênh, phân loại hành vi — mỗi khách một dòng/ngày. Bảng này vừa nuôi dashboard giám đốc, vừa làm feature cho mô hình churn và điểm tín dụng. Gold thường denormalize (gộp sẵn) để BI đọc nhanh, và ghi vào bảng có phân vùng theo ngày để dashboard chỉ quét phần cần.
Cả ba tầng ghi ra định dạng lakehouse có ACID — Delta Lake hoặc Apache Iceberg — để có transaction, time travel, và cập nhật an toàn (xem thêm Delta trên Databricks).
Cấu trúc code job production
Một job production không phải một notebook dài. Tách rõ ba phần: config (tham số, đường dẫn, ngày chạy), IO (đọc/ghi), transform (logic thuần, nhận DataFrame trả DataFrame). Transform tách riêng vì nó test được không cần cluster — chủ đề ở streaming & testing.
# jobs/daily_etl.py (MINH HOẠ, rút gọn)
import argparse
from pyspark.sql import SparkSession, functions as F
from pyspark.sql.window import Window
# ---- IO ----
def read_bronze(spark, path, batch_date):
return (spark.read.format("delta").load(path)
.where(F.col("batch_date") == batch_date))
def write_delta(df, path, batch_date):
(df.write.format("delta").mode("overwrite")
.option("replaceWhere", f"batch_date = '{batch_date}'") # idempotent
.partitionBy("batch_date").save(path))
# ---- TRANSFORM (logic thuần, test được) ----
def to_silver(tx_raw):
w = Window.partitionBy("tx_id").orderBy(F.col("ingest_ts").desc())
return (tx_raw
.withColumn("amount", F.col("amount").cast("decimal(18,2)"))
.withColumn("currency", F.upper(F.trim("currency")))
.withColumn("channel", F.coalesce(F.col("channel"), F.lit("UNKNOWN")))
.withColumn("rn", F.row_number().over(w)).where("rn = 1").drop("rn")
.where(F.col("tx_id").isNotNull() & (F.col("amount") >= 0)))
def to_gold(tx_silver, cust_dim):
agg = (tx_silver.groupBy("customer_id", "batch_date")
.agg(F.sum("amount").alias("total_spend"),
F.count("*").alias("tx_count"),
F.countDistinct("channel").alias("channels_used")))
# bảng khách hàng nhỏ -> broadcast, tránh shuffle bên lớn
return agg.join(F.broadcast(cust_dim), "customer_id", "left")
# ---- ORCHESTRATION ----
def main(batch_date):
spark = (SparkSession.builder.appName(f"daily_etl_{batch_date}")
.config("spark.sql.adaptive.enabled", "true") # AQE
.config("spark.sql.shuffle.partitions", "400")
.getOrCreate())
tx_raw = read_bronze(spark, "s3://lake/bronze/tx", batch_date)
cust = spark.read.format("delta").load("s3://lake/silver/customer_dim")
silver = to_silver(tx_raw)
write_delta(silver, "s3://lake/silver/tx", batch_date)
gold = to_gold(silver, cust)
write_delta(gold, "s3://lake/gold/customer_daily", batch_date)
if __name__ == "__main__":
p = argparse.ArgumentParser()
p.add_argument("--batch-date", required=True) # tham số ngày chạy
main(p.parse_args().batch_date)
Điểm cần chú ý: ngày chạy là tham số, không hardcode current_date(). Nhờ đó có thể backfill (chạy lại cho ngày quá khứ) và tái chạy đúng một ngày mà không đụng ngày khác — nền tảng của idempotency.
Tối ưu thực chiến
Đây là phần phân biệt job "chạy được" với job "chạy trong cửa sổ đêm mà không đốt tiền cluster". Nền tảng lý thuyết ở shuffle & partitions và tuning; dưới đây là cẩm nang áp dụng.
1. Partition đúng & pruning. Ghi bảng phân vùng theo cột thường lọc (batch_date) để truy vấn chỉ quét phân vùng cần (partition pruning), giảm I/O hàng chục lần. Đừng phân vùng theo cột có cardinality quá cao (như customer_id) — sinh hàng triệu thư mục nhỏ, giết metadata.
2. Broadcast join bảng nhỏ. Khi join bảng lớn (giao dịch) với bảng nhỏ (danh mục khách, danh mục chi nhánh cỡ vài chục MB), dùng F.broadcast() để Spark gửi bảng nhỏ tới mọi executor, loại bỏ shuffle bên lớn. AQE cũng tự broadcast nếu ước lượng dưới spark.sql.autoBroadcastJoinThreshold (mặc định 10MB), nhưng khi biết chắc thì ép tay an toàn hơn.
3. Tránh shuffle thừa. Mỗi groupBy, join, distinct, repartition là một shuffle — ghi dữ liệu ra đĩa và truyền qua mạng, khâu đắt nhất. Gộp các phép cùng khoá, lọc trước khi join/aggregate để giảm dữ liệu, và tránh repartition không cần thiết.
4. Xử lý data skew. Khi một khoá join lệch (ví dụ tài khoản nội bộ tổng hợp gom cực nhiều giao dịch), một task ôm hầu hết dữ liệu, cả job chờ nó. Cách xử lý: bật AQE skew join (spark.sql.adaptive.skewJoin.enabled) để Spark tự tách partition lệch; hoặc salting (thêm hậu tố ngẫu nhiên vào khoá lệch rồi gộp lại); hoặc tách riêng khoá "nóng" xử lý bằng broadcast.
5. Tránh Python UDF. UDF Python bắt dữ liệu vượt biên JVM↔Python từng dòng, chậm gấp bội và phá predicate pushdown. Luôn tìm F.* trước; nếu buộc phải có logic Python, dùng pandas_udf vectorized — chi tiết bài 4.
6. Cache hợp lý. cache()/persist() chỉ đáng khi một DataFrame được dùng lại nhiều lần (ví dụ silver vừa ghi vừa dùng tính gold). Cache bừa làm đầy bộ nhớ, gây spill và GC. Nhớ unpersist() khi xong.
7. AQE — Adaptive Query Execution. Bật spark.sql.adaptive.enabled (mặc định bật từ Spark 3.2). AQE dùng thống kê runtime để: gộp partition shuffle nhỏ lại (coalesce), đổi sort-merge join thành broadcast khi bảng thực nhỏ, và xử lý skew. Đây là "tối ưu miễn phí" nên gần như luôn để bật.
8. Số partition & kích thước file. spark.sql.shuffle.partitions mặc định 200 thường quá ít cho hàng trăm triệu dòng — mỗi partition quá to gây spill; đặt sao cho mỗi partition xử lý cỡ 128–256MB. Ở đầu ra, tránh small files problem: dùng coalesce/AQE để mỗi file gold cỡ vài trăm MB, tránh sinh hàng nghìn file bé làm chậm đọc sau này.
9. Đọc explain(). Trước khi tối ưu mù, gọi df.explain("formatted") để xem physical plan. Tìm Exchange (shuffle) — có bao nhiêu, ở đâu; BroadcastHashJoin vs SortMergeJoin; và có predicate pushdown / partition pruning chưa (PushedFilters, PartitionFilters). Plan nói sự thật; đọc được nó là điều kiện cần để tối ưu đúng chỗ.
Vận hành: điều phối, idempotency, giám sát
Điều phối bằng Airflow. Job không tự chạy — một DAG Airflow gọi spark-submit mỗi đêm, truyền --batch-date {{ ds }} (ngày logic của lần chạy). Airflow lo lịch, retry, cảnh báo, và phụ thuộc (chờ file nguồn tới mới chạy). Mô hình spark-submit từ Airflow (qua SparkSubmitOperator hoặc gọi API cluster) tách rõ điều phối khỏi tính toán.
Idempotency & tái chạy. Nguyên tắc sống còn: chạy lại cùng batch_date phải cho cùng kết quả, không nhân đôi. Cơ chế: ghi Delta với replaceWhere "batch_date = ..." (chỉ ghi đè đúng phân vùng ngày đó), hoặc MERGE theo khoá. Nhờ đó khi job đêm lỗi giữa chừng, sáng chạy lại toàn bộ ngày mà không sợ đếm trùng. Kết hợp checkpoint (đặc biệt với streaming ở bài 7) để phần đã xong không phải làm lại.
Giám sát. Theo dõi Spark UI: tab Stages xem task nào lệch/spill, tab SQL xem plan thực, tab Executors xem GC và memory. Xuất metric (thời lượng job, số dòng vào/ra, số bản ghi quarantine) ra hệ giám sát để phát hiện bất thường. Chi phí cluster là số thực: job chạy lâu = trả tiền lâu; tối ưu shuffle và autoscaling giảm hoá đơn trực tiếp.
Đóng gói & test. Đóng gói code thành package (wheel), test transform bằng unit test trên SparkSession local; chi tiết triển khai production ở Spark production. Test trước khi lên cluster tránh việc phát hiện lỗi logic sau 40 phút chạy và một hoá đơn cluster vô ích.
Use case thực tế
NCB — job ETL cuối ngày trên cluster. (Các con số dưới đây là ước lượng minh hoạ để hình dung quy mô, không phải số đo thực tế production.)
Bối cảnh ước lượng: mỗi đêm pipeline xử lý ~250 triệu giao dịch (thẻ + kênh) và ~8 triệu tài khoản cho ~5 triệu khách. Cluster ước lượng ~20 executor, mỗi executor 4 core / 16GB. Cửa sổ chạy: 01:00–05:00.
Diễn biến một đêm:
- Bronze (nạp thô ~250 triệu dòng từ landing + Kafka): ~15 phút, chủ yếu I/O.
- Silver (ép kiểu, dedup, validate, tách quarantine): ~35 phút; khoảng 0,3% bản ghi vào quarantine (thiếu khoá, sai kiểu tiền).
- Gold (
daily_customer_360: join giao dịch với dim khách + chi nhánh, aggregate về mức khách/ngày): ban đầu ~50 phút.
Điểm nghẽn ước lượng ở tầng gold: bảng dim khách bị đổ vào sort-merge join gây shuffle 250 triệu dòng, cộng skew ở nhóm tài khoản nội bộ. Ba can thiệp:
F.broadcast(customer_dim)(bảng ~200MB) → bỏ shuffle bên lớn.- Bật AQE skew join → tách partition tài khoản nội bộ.
- Đặt
shuffle.partitions=400và coalesce đầu ra → hết spill, file gold ~256MB.
Kết quả ước lượng: gold từ ~50 phút xuống ~18 phút (giảm ~60%), tổng pipeline từ ~1h40 xuống ~1h08, nằm gọn trong cửa sổ đêm và giảm tương ứng chi phí giờ-cluster. Khi một đêm file thẻ tới trễ và job fail giữa gold, nhờ replaceWhere theo batch_date, đội chỉ cần retry đúng ngày đó từ Airflow — số liệu không bị đếm trùng, các ngày khác không đụng tới.
Ghi nhớ
- Medallion là xương sống: bronze (thô, bất biến) → silver (sạch, dedup, validate) → gold (tổng hợp/đặc trưng cho BI & ML). Mỗi tầng là bảng vật lý, tái chạy độc lập.
- Ghi ra lakehouse có ACID (Delta/Iceberg) để có transaction, time travel, ghi đè an toàn.
- Cấu trúc code: tách config / IO / transform; transform là hàm thuần để test không cần cluster. Ngày chạy là tham số, không hardcode.
- Tối ưu: partition + pruning, broadcast bảng nhỏ, giảm shuffle, xử lý skew (AQE/salting), tránh Python UDF, cache có chủ đích, đặt
shuffle.partitionshợp lý, tránh small files, và đọcexplain()trước khi tối ưu. - AQE là tối ưu runtime gần như luôn nên bật.
- Vận hành: Airflow gọi
spark-submitvới--batch-date; idempotency bằngreplaceWhere/MERGEđể tái chạy an toàn; giám sát Spark UI và chi phí cluster; đóng gói + test trước khi lên production. - Điểm nghẽn thật thường ở tầng gold (join + aggregate + skew) — đó là nơi đặt công tối ưu để đưa job về trong cửa sổ đêm.
Nguồn tham khảo
- Spark SQL, DataFrames and Datasets Guide — Apache Spark Documentation
- Performance Tuning — Apache Spark Documentation (AQE, broadcast join, coalesce/skew, shuffle partitions)
- Delta Lake Documentation — ACID transactions, time travel,
replaceWhere,MERGE - Apache Iceberg Documentation — table format cho lakehouse
- Apache Airflow Documentation — DAG, scheduling,
SparkSubmitOperator - Thông tư 11/2021/TT-NHNN — Quy định về phân loại tài sản có, mức trích, phương pháp trích lập dự phòng rủi ro trong hoạt động của tổ chức tín dụng
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ẻ!