PySpark 7 — Structured Streaming & Kiểm thử/Đóng gói
Hai kỹ năng biến code PySpark thành sản phẩm chạy được
Cho tới đây trong series, ta đã viết các phép biến đổi trên DataFrame tĩnh (batch). Bài này ghép nốt hai mảnh còn thiếu để một job PySpark thực sự chạy trên production: (A) xử lý dữ liệu chảy liên tục bằng Structured Streaming, và (B) kiểm thử + đóng gói job đúng cách để không đẩy lỗi ra hệ thống thật.
Phần lý thuyết nền của streaming — semantics, cơ chế state store, so sánh với Flink — đã được trình bày ở Spark Structured Streaming. Ở đây ta tập trung vào cách viết bằng PySpark: API cụ thể, cấu trúc code, cạm bẫy Python. Cả hai kỹ năng này đều là thứ phân biệt "chạy được trong notebook" với "vận hành được lúc 2 giờ sáng".
PHẦN A — Structured Streaming bằng PySpark
Ý tưởng cốt lõi: bảng vô hạn
Structured Streaming coi luồng dữ liệu như một bảng không ngừng được nối thêm dòng (unbounded table). Bạn viết truy vấn trên bảng đó y hệt như trên DataFrame tĩnh; Spark chịu trách nhiệm chạy lại truy vấn một cách tăng dần (incremental) mỗi khi có dữ liệu mới. Nhờ vậy code streaming và code batch dùng chung một API DataFrame — đây là điểm mạnh lớn nhất so với các engine đời cũ.
Vòng đời một job gồm ba khối: spark.readStream (nguồn) → các phép biến đổi DataFrame → writeStream (đích).
readStream: các nguồn
Khai báo nguồn gần giống spark.read nhưng thay bằng readStream:
- Kafka (
format("kafka")): nguồn production phổ biến nhất; cầnkafka.bootstrap.serversvàsubscribe. Chi tiết Kafka xem lưu trữ & độ tin cậy Kafka. Giá trị message nằm ở cộtvaluedạng binary, phảiCASTvà parse. - File (
format("json"/"parquet"/...)): Spark theo dõi thư mục, phát hiện file mới. Tốt cho landing zone nhưng phải quản lý số lượng file. - Socket (
format("socket")): CHỈ để thử nghiệm/demo, không dùng production (không tái lập được). - Rate: nguồn sinh dữ liệu giả để test hiệu năng.
Biến đổi giống hệt batch
Sau khi có streaming DataFrame, bạn dùng đúng các API đã học ở DataFrame API: select, filter, withColumn, groupBy, join... Điểm khác biệt Spark tự lo phần lớn. Chỉ có vài phép không hỗ trợ trên stream (ví dụ sort toàn cục ngoài aggregation, hay một số kiểu multi-aggregation), và một số phép cần watermark mới thực hiện được (join stream-stream, dedup theo thời gian).
Bốn khái niệm phải nắm đúng
1. Micro-batch. Mặc định, Spark xử lý stream bằng cách chia thành các lô nhỏ liên tiếp: cứ theo chu kỳ, gom dữ liệu mới đến, chạy truy vấn trên lô đó rồi ghi kết quả. Độ trễ cỡ giây (không phải mili-giây thật sự). Có chế độ Continuous processing cho độ trễ ~1ms nhưng còn thử nghiệm và hạn chế; thực tế gần như luôn dùng micro-batch.
2. Output mode — quyết định cái gì được ghi ra sink mỗi micro-batch:
| Mode | Ý nghĩa | Dùng khi |
|---|---|---|
append | Chỉ ghi dòng mới, không đổi nữa | Không aggregation, hoặc aggregation có watermark đã "chốt" cửa sổ |
update | Ghi các dòng kết quả vừa thay đổi | Aggregation cần cập nhật liên tục (dashboard) |
complete | Ghi toàn bộ bảng kết quả mỗi lần | Aggregation nhỏ, cần snapshot đầy đủ |
Chọn sai mode là lỗi thường gặp: aggregation không watermark mà xin append sẽ báo lỗi ngay lúc khởi động.
3. Trigger — quyết định khi nào chạy micro-batch tiếp theo:
- Mặc định (không set): chạy lô mới ngay khi lô trước xong.
processingTime="1 minute": chạy theo nhịp cố định.availableNow=True: xử lý hết dữ liệu đang có rồi dừng — cực hữu ích để chạy stream như một batch định kỳ (thay choonce, vốn đã deprecated), tận dụng checkpoint để chỉ đọc phần mới.
4. Checkpoint — BẮT BUỘC cho mọi job production. checkpointLocation là một thư mục bền vững (HDFS/S3/ABFS) nơi Spark lưu offset đã đọc và trạng thái aggregation. Nhờ nó, khi job crash và khởi động lại, Spark biết đọc tiếp từ đâu và khôi phục state — đây là nền tảng của bảo đảm exactly-once (với sink hỗ trợ, như Delta/Kafka). Mỗi query phải có checkpoint riêng; không được dùng chung thư mục cho hai query khác nhau.
Watermark & windowed aggregation
Dữ liệu thực đến trễ và lộn xộn: một giao dịch lúc 10:00 có thể tới hệ thống lúc 10:03. Ta thường muốn tổng hợp theo event-time (thời điểm sự việc xảy ra) chứ không phải thời điểm nhận. Windowed aggregation nhóm dữ liệu theo cửa sổ thời gian:
from pyspark.sql import functions as F
agg = (events
.withWatermark("event_time", "10 minutes")
.groupBy(F.window("event_time", "5 minutes"), "branch_id")
.agg(F.sum("amount").alias("total"),
F.count("*").alias("txn_count")))
withWatermark("event_time", "10 minutes") nói với Spark: "chấp nhận dữ liệu trễ tối đa 10 phút; sau ngưỡng đó coi như cửa sổ đã đóng và xóa state của nó". Không có watermark, Spark phải giữ trạng thái của mọi cửa sổ mãi mãi → state phình vô hạn → job chết. Watermark là cơ chế đánh đổi giữa độ chính xác (bắt được dữ liệu trễ) và chi phí bộ nhớ.
writeStream & foreachBatch
writeStream khai báo sink. Với Delta/Parquet bạn ghi thẳng bằng format("delta"). Nhưng khi cần logic ghi tùy ý — MERGE/upsert vào Delta, ghi đồng thời nhiều đích, gọi API ngoài — dùng foreachBatch: hàm của bạn nhận (batch_df, batch_id) và batch_df là một DataFrame tĩnh bình thường, nên bạn được dùng đầy đủ mọi API batch (kể cả MERGE). Đây cũng là chỗ đảm bảo idempotency: dùng batch_id để bỏ qua lô đã ghi, hoặc dùng MERGE theo khóa để ghi lại an toàn khi retry.
Ví dụ: Kafka → cửa sổ → Delta
# MINH HOẠ (Python) — không phải SQL sandbox
from pyspark.sql import SparkSession, functions as F
from pyspark.sql.types import StructType, StringType, DoubleType, TimestampType
spark = SparkSession.builder.appName("txn-stream").getOrCreate()
schema = (StructType()
.add("txn_id", StringType()).add("account_no", StringType())
.add("amount", DoubleType()).add("event_time", TimestampType())
.add("branch_id", StringType()))
raw = (spark.readStream.format("kafka")
.option("kafka.bootstrap.servers", "broker:9092")
.option("subscribe", "transactions")
.option("startingOffsets", "latest")
.load())
parsed = (raw
.select(F.from_json(F.col("value").cast("string"), schema).alias("d"))
.select("d.*"))
windowed = (parsed
.withWatermark("event_time", "10 minutes")
.groupBy(F.window("event_time", "5 minutes"), "branch_id")
.agg(F.sum("amount").alias("total"),
F.count("*").alias("txn_count")))
query = (windowed.writeStream
.format("delta")
.outputMode("append") # cửa sổ có watermark → append hợp lệ
.option("checkpointLocation", "/mnt/chk/txn_window")
.trigger(processingTime="1 minute")
.start("/mnt/delta/txn_window"))
query.awaitTermination()
Nền tảng Delta/lakehouse xem Spark Delta Lakehouse và Databricks Delta Lake.
PHẦN B — Kiểm thử & đóng gói job PySpark
Streaming chạy được rồi, nhưng làm sao biết logic đúng trước khi lên production? PySpark có đặc thù riêng khi test và đóng gói vì phải khởi động cả một cụm Spark (dù local).
Nguyên tắc số một: tách logic transform khỏi I/O
Job production trộn ba thứ: đọc nguồn (I/O), biến đổi (logic), ghi đích (I/O). Chỉ có logic biến đổi là thứ chứa nghiệp vụ dễ sai và đáng test nhất. Vì vậy hãy viết các phép biến đổi thành hàm thuần: nhận DataFrame vào, trả DataFrame ra, không tự đọc/ghi, không tạo SparkSession bên trong.
# transforms.py — hàm thuần, dễ test
from pyspark.sql import DataFrame, functions as F
def flag_large_txn(df: DataFrame, threshold: float = 500_000_000) -> DataFrame:
return df.withColumn(
"is_large",
F.when(F.col("amount") >= threshold, F.lit(True)).otherwise(F.lit(False)),
)
Hàm này chạy y hệt cho batch lẫn streaming, và test được bằng vài dòng dữ liệu mẫu — không cần Kafka.
Kim tự tháp test cho PySpark
Đáy tháp — test hàm transform — chiếm số lượng lớn nhất vì rẻ và nhanh. Càng lên cao càng ít và đắt.
SparkSession fixture với pytest
Khởi động Spark tốn vài giây, nên tạo một session dùng chung cho cả phiên test qua fixture scope="session":
# conftest.py
import pytest
from pyspark.sql import SparkSession
@pytest.fixture(scope="session")
def spark():
s = (SparkSession.builder
.master("local[*]")
.appName("tests")
.config("spark.sql.shuffle.partitions", "1") # nhỏ → test nhanh
.config("spark.ui.enabled", "false")
.getOrCreate())
yield s
s.stop()
local[*] chạy Spark ngay trong tiến trình test, không cần cụm. Đặt shuffle.partitions=1 để tránh sinh 200 partition rỗng làm test chậm.
So sánh DataFrame bằng chispa
So sánh hai DataFrame "bằng tay" rất cực vì thứ tự dòng/cột và kiểu dữ liệu. Thư viện chispa cung cấp assert chuyên dụng, báo lỗi trực quan (highlight ô sai):
# test_transforms.py
from chispa import assert_df_equality
from transforms import flag_large_txn
def test_flag_large_txn(spark):
cols = ["txn_id", "amount"]
inp = spark.createDataFrame(
[("t1", 600_000_000.0), ("t2", 100_000.0)], cols)
got = flag_large_txn(inp)
expected = spark.createDataFrame(
[("t1", 600_000_000.0, True), ("t2", 100_000.0, False)],
["txn_id", "amount", "is_large"])
assert_df_equality(got, expected, ignore_row_order=True)
ignore_row_order=True vì Spark không đảm bảo thứ tự dòng. Còn ignore_nullable=True hữu ích khi schema chỉ khác ở cờ nullable. Với so sánh gần đúng số thực, chispa có assert_approx_df_equality(precision=...).
Test code khác test dữ liệu
Một điểm hay bị lẫn: pytest test logic code của bạn (given input X thì transform ra Y). Nó không kiểm tra dữ liệu production có sạch không. Kiểm thử chất lượng dữ liệu thật (null bất thường, phá vỡ ràng buộc, drift) là một khâu riêng khi job chạy — xem Data Quality. Đừng kỳ vọng bộ pytest bắt được dữ liệu bẩn từ nguồn; và ngược lại, đừng để rule chất lượng dữ liệu thay cho unit test logic.
Đóng gói project
Cấu trúc gợi ý cho một job đóng gói được:
txn_job/
├── src/txn_job/
│ ├── __init__.py
│ ├── transforms.py # hàm thuần
│ ├── io.py # đọc/ghi
│ └── main.py # entrypoint: parse args → gọi transforms
├── tests/
│ ├── conftest.py
│ └── test_transforms.py
├── pyproject.toml
└── requirements.txt
main.py là entrypoint: nhận tham số (đường dẫn nguồn/đích, ngày chạy, ngưỡng) qua argparse, tạo SparkSession, gọi io + transforms, ghi kết quả. Tham số hóa thay vì hard-code giúp cùng một job chạy cho nhiều ngày/môi trường.
spark-submit & quản lý dependency
Chạy job bằng spark-submit. Ba cờ quan trọng nhất:
| Cờ | Công dụng |
|---|---|
--py-files | Gửi code Python phụ (file .zip/.egg/.py) tới executor |
--packages | Kéo dependency Maven (JVM), ví dụ connector Kafka, Delta |
--files | Gửi file cấu hình/tài nguyên (yaml, cert...) tới executor |
--conf | Đặt cấu hình Spark lúc chạy (bộ nhớ, shuffle, extension) |
Điểm bẫy lớn của PySpark: dependency có hai loại. Thư viện JVM (connector Kafka, Delta) đi qua --packages. Thư viện Python (pandas, scikit-learn của bạn) phải có mặt trên mọi executor. Vài cách phổ biến:
- venv + zip / conda-pack: đóng gói toàn bộ môi trường Python thành archive, gửi qua
--archives, trỏPYSPARK_PYTHONvào đó → executor dùng đúng phiên bản thư viện. - Docker image: build image chứa sẵn Python deps, dùng trên K8s/YARN — tái lập tốt nhất, hợp CI/CD.
--py-files package.zip: đủ cho code của chính bạn nếu deps đã có trên cluster.
# MINH HOẠ (shell)
spark-submit \
--master yarn --deploy-mode cluster \
--packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.0,io.delta:delta-spark_2.12:3.1.0 \
--py-files dist/txn_job.zip \
--files conf/prod.yaml \
--conf spark.sql.shuffle.partitions=200 \
--conf spark.sql.streaming.stateStore.providerClass=...RocksDB... \
src/txn_job/main.py \
--source transactions --sink /mnt/delta/txn --run-date 2026-07-14
Lưu ý: phiên bản package (Scala 2.12, Spark 3.5.0) phải khớp cluster, nếu không job chết vì lỗi lớp không tương thích.
Logging, CI & idempotency
- Logging: dùng
loggingchuẩn ở driver; tránh in ở trong UDF/executor vì log rải rác nhiều máy. Log tham số đầu vào và số dòng ở mỗi bước để dễ truy vết. - CI: chạy
pytestvớilocal[*]ngay trong pipeline CI (GitHub Actions/GitLab). Cần cài Java + đúng PySpark; cân nhắc chạy trong container để đồng nhất. Vì test dùng dữ liệu mẫu nhỏ, toàn bộ suite thường xong trong 1–2 phút — đủ nhanh để chặn merge khi đỏ. - Idempotency & tái lập: job phải chạy lại cho ra cùng kết quả. Ghi theo partition ngày và overwrite đúng partition (hoặc MERGE theo khóa) thay vì append mù; luôn truyền
--run-datethay vì đọcnow()bên trong; với streaming thì checkpoint +foreachBatchMERGE lo phần này.
Nền tảng chung về Spark xem PySpark tổng quan; triển khai production sâu hơn ở Spark production.
Use case thực tế
Bối cảnh NCB — giám sát giao dịch gần thời gian thực. Đội rủi ro cần một bảng theo dõi cập nhật mỗi phút: mỗi chi nhánh, trong cửa sổ 5 phút, tổng giá trị giao dịch và số giao dịch lớn (≥ 500 triệu), để phát hiện bất thường sớm thay vì chờ báo cáo cuối ngày.
Kiến trúc. Core banking đẩy sự kiện giao dịch vào Kafka topic transactions (~1.500 giao dịch/giây giờ cao điểm, ~80 triệu bản ghi/ngày). Job PySpark Structured Streaming:
readStreamtừ Kafka, parse JSON theo schema cố định.withWatermark("event_time", "10 minutes")(giao dịch từ chi nhánh xa có thể trễ vài phút),groupBy(window 5 phút, branch_id)tínhsum(amount),countvà số giao dịch lớn.writeStreamforeachBatchMERGE vào bảng Deltarisk.txn_5mintheo khóa(window, branch_id)→ idempotent khi retry.- Trigger
processingTime="1 minute", checkpoint trên ABFS; job tự khôi phục sau khi node bị kill nhờ offset + state đã lưu.
Chặn lỗi trước production. Logic phân loại giao dịch lớn và quy đổi ngoại tệ nằm trong transforms.py (hàm thuần). Bộ pytest có ~40 case dùng chispa: ngưỡng biên (đúng 500 triệu), số âm (giao dịch điều chỉnh), null event_time, đa tiền tệ. Một lần refactor đổi toán tử > thành >= bị test biên bắt ngay trong CI (đỏ pipeline), không lọt ra bảng rủi ро. Suite chạy 90 giây trên runner container có sẵn Java 11 + PySpark 3.5.
Kết quả ước lượng (minh họa). Độ trễ dashboard từ "cuối ngày" xuống ~60–90 giây; đội rủi ro rút ngắn thời gian phát hiện chuỗi giao dịch bất thường từ hàng giờ xuống vài phút. Nhờ tách transform + chispa, số bug logic lọt production trong quý giảm rõ; nhờ checkpoint, một sự cố node executor ban đêm được job tự phục hồi mà không mất/nhân đôi dữ liệu (MERGE idempotent).
Ghi nhớ
- Structured Streaming = cùng API DataFrame với batch:
readStream→ biến đổi →writeStream. Học batch trước là học được 80% streaming. - Checkpoint là bắt buộc cho mọi job production: lưu offset + state, nền tảng khôi phục và exactly-once. Mỗi query một thư mục riêng.
- Output mode (append/update/complete) và trigger phải chọn khớp với việc có aggregation/watermark hay không — chọn sai là lỗi lúc khởi động.
- Watermark giới hạn dữ liệu trễ để state không phình vô hạn; cần cho windowed aggregation, dedup, join stream-stream.
- foreachBatch cho ghi tùy ý và idempotent (MERGE theo khóa, dùng
batch_id); batch_df là DataFrame tĩnh bình thường. - Test PySpark: tách hàm transform thuần khỏi I/O, dùng
SparkSessionfixturelocal[*](scope="session",shuffle.partitionsnhỏ), so DataFrame bằng chispa vớiignore_row_order. - Test code ≠ test dữ liệu: pytest kiểm logic; chất lượng dữ liệu production là khâu riêng (Data Quality).
- Đóng gói: entrypoint tham số hóa;
spark-submitvới--py-files(code Python),--packages(JVM/connector),--files,--conf; deps Python qua conda-pack/venv-zip/Docker để tái lập. - Chạy pytest trong CI để chặn merge khi đỏ; đảm bảo idempotency bằng overwrite theo partition/MERGE và truyền
--run-datethay chonow().
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.
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ẻ!