PySpark 9 — Dữ liệu lồng & phức tạp (struct, array, map, JSON)
Vì sao dữ liệu ngân hàng hầu như luôn lồng
Ở tổng quan PySpark ta hình dung Spark làm việc với bảng phẳng: hàng và cột đơn giản. Nhưng thực tế dữ liệu chảy vào một ngân hàng như NCB hiếm khi phẳng. Một message giao dịch thẻ từ Kafka là một JSON lồng nhiều tầng: thông tin thẻ nằm trong một object con, danh sách mục hàng (line items) là một mảng các object, thông tin thiết bị và định vị lại là object khác. Một payload API Open Banking, một dòng log ứng dụng, một event từ core banking — tất cả đều có cấu trúc cây, không phải bảng.
Spark xử lý được điều này nhờ hệ thống kiểu phức hợp (complex types) ngay trong DataFrame — không cần "đập phẳng" dữ liệu trước khi nạp. Đây là điểm mạnh lớn: bạn giữ được cấu trúc gốc, parse có kiểm soát, rồi mới quyết định làm phẳng bao nhiêu. Bài này nối tiếp DataFrame API, đi sâu vào ba kiểu lồng cốt lõi và cách parse JSON — kỹ năng gần như bắt buộc với bất kỳ ai làm data engineer trong ngành tài chính.
Ba kiểu phức hợp: struct, array, map
Spark có ba kiểu phức hợp, ánh xạ gần đúng với JSON:
| Kiểu | Ý nghĩa | Tương ứng JSON | Ví dụ ngân hàng |
|---|---|---|---|
| StructType | Bản ghi lồng, các field có tên và kiểu cố định | object {...} | Thông tin thẻ: {card_no, brand, issuer} |
| ArrayType | Mảng các phần tử cùng kiểu | array [...] | Danh sách mục hàng của giao dịch |
| MapType | Cặp khoá–giá trị, khoá cùng kiểu, giá trị cùng kiểu | object có khoá động | Bộ tham số/metadata linh hoạt |
Khác biệt then chốt giữa struct và map: struct có schema cố định (biết trước tên field, kiểm tra lúc compile/analyze), còn map có khoá động (không biết trước, chỉ biết kiểu khoá). Khi số field cố định và đã biết — dùng struct; khi khoá thay đổi tuỳ bản ghi (ví dụ bag of attributes) — dùng map.
Truy cập cột lồng
Ba cách truy cập, tuỳ kiểu:
from pyspark.sql import functions as F
# Struct: dấu chấm hoặc getField
df.select("txn.card.card_no") # dot notation
df.select(F.col("txn.card").getField("brand"))
# Array: getItem theo chỉ số
df.select(F.col("items").getItem(0)) # phần tử đầu
df.select(F.col("items")[0]["sku"]) # lồng array→struct
# Map: getItem theo khoá
df.select(F.col("meta").getItem("channel"))
df.select(F.col("meta")["device_id"])
Dấu chấm còn hoạt động xuyên qua mảng struct: df.select("items.sku") trên một mảng struct trả về một mảng các sku — Spark tự "chiếu" field qua từng phần tử. Đây là mẹo rất tiện nhưng cũng dễ gây nhầm khi bạn tưởng mình lấy được một giá trị đơn.
JSON: từ chuỗi thô đến DataFrame có schema
Đọc file JSON trực tiếp
Khi dữ liệu đã là file JSON (một object mỗi dòng — JSON Lines), spark.read.json tự suy luận schema:
df = spark.read.json("s3://raw/txn/2026-07-13/*.json")
# JSON nhiều dòng (một object trải nhiều dòng, hoặc mảng ở gốc):
df = spark.read.option("multiLine", True).json("path")
Suy luận schema tiện lúc khám phá nhưng không nên dùng trong production: Spark phải quét dữ liệu một lượt để đoán kiểu, chậm và không ổn định. Cách đúng là cấp schema tường minh qua .schema(...) — vừa nhanh vừa bắt lỗi sớm. Xem thêm đọc/ghi nguồn dữ liệu.
from_json — vũ khí chính cho message Kafka
Tình huống phổ biến nhất ở ngân hàng: bạn đọc từ Kafka, cột value là một chuỗi JSON (string), không phải cấu trúc. Bạn cần from_json để parse chuỗi đó thành struct theo schema mình định nghĩa:
from pyspark.sql.types import (StructType, StructField, StringType,
DoubleType, ArrayType, TimestampType, MapType)
txn_schema = StructType([
StructField("txn_id", StringType()),
StructField("amount", DoubleType()),
StructField("ccy", StringType()),
StructField("ts", TimestampType()),
StructField("card", StructType([
StructField("card_no", StringType()),
StructField("brand", StringType()),
])),
StructField("items", ArrayType(StructType([
StructField("sku", StringType()),
StructField("qty", DoubleType()),
StructField("price", DoubleType()),
]))),
StructField("meta", MapType(StringType(), StringType())),
])
parsed = df.select(F.from_json("value", txn_schema).alias("txn"))
from_json với schema tường minh là cách chuẩn để xử lý message Kafka JSON — cực hay dùng khi đọc topic giao dịch (xem Kafka storage & reliability). Nếu chuỗi JSON hỏng, các field sẽ là null thay vì làm crash job — hành vi này quan trọng cho pipeline chịu lỗi.
Các hàm JSON liên quan:
to_json: chiều ngược lại — serialize struct/array/map thành chuỗi JSON (để ghi ra Kafka, hoặc gộp nhiều cột thành một payload).schema_of_json: suy luận schema từ một chuỗi JSON mẫu — tiện để lấy nhanh khung schema rồi chỉnh tay, không nên dùng runtime.get_json_object: rút một giá trị theo đường dẫn JSONPath ($.card.brand) mà không cần parse cả struct — nhẹ khi chỉ cần vài field.json_tuple: rút nhiều field top-level cùng lúc từ chuỗi JSON, nhanh hơn nhiều lầnget_json_object.
Quy tắc thực dụng: cần nhiều field / cấu trúc lồng → from_json với schema; chỉ cần một hai field nông → get_json_object/json_tuple.
Mảng (Array): explode và higher-order functions
explode — mảng thành nhiều dòng
Thao tác kinh điển nhất với mảng là explode: biến một dòng có mảng N phần tử thành N dòng, mỗi dòng một phần tử. Đây là bước "chuẩn hoá" đưa mảng mục hàng về dạng bảng chi tiết.
# Mỗi giao dịch có mảng items → mỗi item một dòng
detail = parsed.select("txn.txn_id",
F.explode("txn.items").alias("item"))
# explode_outer: giữ dòng cả khi mảng rỗng/null (item = null)
detail = parsed.select("txn.txn_id",
F.explode_outer("txn.items").alias("item"))
# posexplode: kèm chỉ số vị trí (0,1,2...) — hữu ích khi thứ tự có nghĩa
parsed.select("txn.txn_id",
F.posexplode("txn.items").alias("pos", "item"))
Khác biệt sống còn: explode loại bỏ dòng có mảng rỗng hoặc null (như inner join ngầm), còn explode_outer giữ lại dòng đó với giá trị null. Nếu bạn cần đếm cả giao dịch không có item nào, dùng explode_outer — nhiều bug số liệu ở ngân hàng đến từ chỗ này.
Higher-order functions — xử lý mảng KHÔNG cần explode
Từ Spark 2.4+, ta có các hàm bậc cao thao tác trực tiếp trên mảng mà không phải explode rồi gom lại — nhanh hơn và giữ nguyên cấu trúc dòng:
# transform: áp một biểu thức lên từng phần tử
F.transform("prices", lambda p: p * 1.1) # tăng 10%
# filter: lọc phần tử theo điều kiện
F.filter("items", lambda x: x["qty"] > 0)
# aggregate: gấp mảng thành một giá trị (khởi tạo, hàm gộp)
F.aggregate("prices", F.lit(0.0), lambda acc, p: acc + p) # tổng
# exists / forall: kiểm tra tồn tại / toàn bộ
F.exists("items", lambda x: x["price"] > 1000)
Ngoài ra: array (tạo mảng từ nhiều cột), array_contains(col, value), size (độ dài mảng, trả -1 nếu null), array_distinct, sort_array, slice, array_union/array_intersect.
Gom ngược lại: collect_list / collect_set
Chiều ngược của explode là gom nhiều dòng thành một mảng, thường trong groupBy:
parsed.groupBy("txn.card.card_no").agg(
F.collect_list("txn.txn_id").alias("all_txns"), # giữ trùng, giữ thứ tự
F.collect_set("txn.ccy").alias("currencies"), # loại trùng
)
Lưu ý collect_list/collect_set không đảm bảo thứ tự (trừ khi kèm sắp xếp) và có thể ngốn bộ nhớ nếu nhóm quá lớn — cẩn thận với khách hàng có hàng triệu giao dịch.
Struct: tạo, làm phẳng, sửa field lồng
Tạo và làm phẳng
# Tạo struct từ nhiều cột (gộp lại)
df.select(F.struct("card_no", "brand").alias("card"))
# Làm phẳng: struct.* bung mọi field con thành cột top-level
flat = parsed.select("txn.*") # txn_id, amount, ccy, ts, card, items, meta
# card vẫn còn là struct → làm phẳng tiếp:
flat = parsed.select("txn.txn_id", "txn.amount",
F.col("txn.card.card_no").alias("card_no"),
F.col("txn.card.brand").alias("brand"))
struct.* (dấu sao) là cách nhanh để bung một tầng struct thành các cột. Với struct lồng nhiều tầng, phải bung từng tầng, hoặc viết hàm đệ quy duyệt schema để tự sinh danh sách cột phẳng — nhiều team có sẵn util flatten_df kiểu này.
Sửa field lồng bằng withField
Trước đây, đổi một field bên trong struct rất phiền (phải dựng lại cả struct). Từ Spark 3.1+, Column.withField và dropFields giải quyết gọn:
# Thêm/ghi đè field brand bên trong struct card
parsed.withColumn("txn",
F.col("txn").withField("card",
F.col("txn.card").withField("brand", F.lit("VISA"))))
# Bỏ field nhạy cảm khỏi struct
parsed.withColumn("txn", F.col("txn").dropFields("card.card_no"))
Map: khoá–giá trị động
F.map_keys("meta") # mảng các khoá
F.map_values("meta") # mảng các giá trị
F.explode("meta") # map → 2 cột key, value (mỗi cặp một dòng)
F.create_map(F.lit("k1"), F.col("v1"), F.lit("k2"), F.col("v2"))
explode trên map trả về hai cột key, value (khác explode array chỉ trả một cột) — tiện khi cần "unpivot" một bag metadata thành bảng dài. Dùng map khi khoá thật sự động; nếu khoá cố định, struct luôn tốt hơn vì được kiểm tra schema.
Schema phức tạp và tiến hoá (nested schema evolution)
Đây là nơi đau đầu nhất trong thực tế. Payload từ hệ thống nguồn thay đổi theo thời gian: đội thẻ thêm field merchant_country vào object card, đổi kiểu amount từ số sang chuỗi, hoặc đổi mảng items thêm một field lồng. Vài cạm bẫy:
- Field mới trong struct: nếu bạn parse bằng schema tường minh, field mới sẽ bị bỏ qua âm thầm — không lỗi, nhưng mất dữ liệu. Ngược lại nếu để Spark tự suy schema, hai file với schema khác nhau có thể merge thành union hoặc gây lỗi.
- Đổi kiểu field lồng:
amountlúc100.0(double) lúc"100"(string) → suy luận thành string, mọi phép tính số ở dưới hỏng. - Field lồng bị đổi tên hoặc di chuyển tầng: dot-path cũ (
txn.card.brand) trả null hàng loạt mà không báo lỗi. - Thứ tự field: với một số format, struct so khớp theo tên; với format khác (như một số reader theo vị trí) lại theo thứ tự — trộn lẫn gây lệch dữ liệu.
Chiến lược phòng thủ: (1) luôn version schema và review khi nguồn đổi; (2) với Delta/Iceberg bật schema evolution có kiểm soát thay vì merge mù; (3) đặt kiểm tra chất lượng ngay sau parse — đếm tỉ lệ null bất thường, kiểm field bắt buộc, so số dòng trước/sau explode (xem data quality). Một field lồng "im lặng thành null" thường là dấu hiệu schema nguồn vừa đổi.
Chiến lược: khi nào giữ lồng, khi nào làm phẳng
Không phải cứ nhận JSON là phải đập phẳng ngay. Cân nhắc:
| Yếu tố | Nghiêng về giữ lồng | Nghiêng về làm phẳng |
|---|---|---|
| Người dùng cuối | Data engineer, đọc bằng Spark | Analyst dùng SQL/BI |
| Truy vấn | Chủ yếu lấy cả bản ghi | Lọc/join/aggregate theo field |
| Quan hệ 1–n | Muốn giữ item gắn với giao dịch | Cần phân tích item độc lập |
| Công cụ downstream | Spark, engine hỗ trợ nested | Warehouse phẳng, dashboard |
Về hiệu năng: giữ lồng tiết kiệm storage và tránh nhân bản dữ liệu (không lặp lại thông tin giao dịch cho mỗi item như sau explode). Làm phẳng lại cho phép column pruning và predicate pushdown hiệu quả hơn, join đơn giản hơn, analyst không cần biết explode. explode một mảng lớn sẽ nhân số dòng — một giao dịch 50 item thành 50 dòng — gây phình dữ liệu và shuffle nặng; hãy explode càng muộn càng tốt và chỉ giữ cột cần. Thực hành phổ biến kiểu medallion: bronze/silver giữ lồng (trung thực với nguồn), gold làm phẳng theo nhu cầu phân tích cụ thể.
Sơ đồ luồng xử lý
Ví dụ đầy đủ: parse payload giao dịch thẻ (minh hoạ)
Ghép mọi thứ lại — parse cột JSON → explode mảng item → làm phẳng struct → chọn cột phẳng. Đây là minh hoạ, tên field và schema mang tính ví dụ:
from pyspark.sql import functions as F
from pyspark.sql.types import (StructType, StructField, StringType,
DoubleType, ArrayType, TimestampType, MapType)
# 1) Schema tường minh cho payload giao dịch thẻ
schema = StructType([
StructField("txn_id", StringType()),
StructField("amount", DoubleType()),
StructField("ccy", StringType()),
StructField("ts", TimestampType()),
StructField("card", StructType([
StructField("card_no", StringType()),
StructField("brand", StringType()),
StructField("issuer", StringType()),
])),
StructField("items", ArrayType(StructType([
StructField("sku", StringType()),
StructField("name", StringType()),
StructField("qty", DoubleType()),
StructField("price", DoubleType()),
]))),
StructField("meta", MapType(StringType(), StringType())),
])
# 2) Đọc từ Kafka: cột value là chuỗi JSON → parse
raw = (spark.read.format("kafka")
.option("subscribe", "card-txn")
.load()
.selectExpr("CAST(value AS STRING) AS value"))
parsed = raw.select(F.from_json("value", schema).alias("t"))
# 3) explode mảng item + làm phẳng struct card → bảng phẳng chi tiết
line_items = (parsed
.select(
F.col("t.txn_id").alias("txn_id"),
F.col("t.ts").alias("ts"),
F.col("t.ccy").alias("ccy"),
F.col("t.card.card_no").alias("card_no"),
F.col("t.card.brand").alias("card_brand"),
F.col("t.meta").getItem("channel").alias("channel"), # lấy 1 khoá map
F.posexplode_outer("t.items").alias("item_pos", "item"),
)
.select("txn_id", "ts", "ccy", "card_no", "card_brand", "channel",
"item_pos",
F.col("item.sku").alias("sku"),
F.col("item.qty").alias("qty"),
F.col("item.price").alias("price"),
(F.col("item.qty") * F.col("item.price")).alias("line_amount"))
)
# 4) kiểm chất lượng nhanh: tỉ lệ item null (dấu hiệu schema đổi)
line_items.select(
F.avg(F.col("sku").isNull().cast("int")).alias("null_sku_ratio")
).show()
Kết quả line_items là bảng phẳng, mỗi dòng một mục hàng, gắn với thông tin giao dịch và thẻ — sẵn sàng cho aggregate, join với danh mục sản phẩm, hoặc feed vào mô hình gian lận.
Use case thực tế
Bối cảnh — NCB, parse payload giao dịch thẻ từ Kafka. Hệ thống thẻ đẩy mỗi giao dịch POS/e-commerce lên topic Kafka dưới dạng JSON lồng nhiều tầng: object card (số thẻ token hoá, brand, issuer), object merchant (mã, tên, MCC, quốc gia), object device/geo, và mảng items (giỏ hàng, mỗi item có sku/qty/price). Đội rủi ro cần một bảng phẳng chi tiết mục hàng cập nhật gần thời gian thực để chấm điểm gian lận và phân tích hành vi chi tiêu.
Vấn đề. Nếu để nguyên JSON, analyst rủi ro không truy vấn được bằng SQL; nếu chỉ lấy tổng tiền giao dịch thì mất chi tiết giỏ hàng — vốn là tín hiệu gian lận quan trọng (ví dụ nhiều item giá trị cao trong một lần quẹt, hoặc pattern sku bất thường).
Cách làm. Job Structured Streaming đọc topic, from_json với schema tường minh (version hoá trong git), explode_outer mảng items để giữ cả giao dịch không có item, làm phẳng struct card/merchant, lấy vài khoá cần từ map meta, rồi ghi Delta theo tầng: silver giữ struct card/merchant (trung thực nguồn), gold là bảng phẳng mục hàng cho phân tích.
Số liệu ước lượng (minh hoạ). Khoảng 8 triệu giao dịch/ngày, trung bình 2,3 item/giao dịch → sau explode ra ~18 triệu dòng mục hàng/ngày. Chuyển từ suy luận schema sang schema tường minh giảm thời gian parse của micro-batch khoảng 35–40% (không còn quét đoán kiểu). Kiểm null_sku_ratio bắt được hai lần đội thẻ đổi tên field lồng trong sáu tháng — mỗi lần chỉ số vọt từ dưới 1% lên trên 90%, cảnh báo kích hoạt trước khi số liệu gian lận bị sai. Tỉ lệ message JSON hỏng ổn định quanh 0,05%, được đẩy sang bảng dead-letter thay vì làm dừng job.
Ghi nhớ
- Ba kiểu phức hợp: StructType (field cố định, có schema), ArrayType (mảng cùng kiểu), MapType (khoá động). Field cố định → struct; khoá động → map.
- Truy cập lồng: dấu chấm
col.a.bhoặcgetFieldcho struct,getItem/[i]cho array,getItem("key")/["key"]cho map. Dot-path chiếu xuyên qua mảng struct. from_json+ schema tường minh là cách chuẩn parse chuỗi JSON của message Kafka — nhanh, chịu lỗi (JSON hỏng thành null). Tránh suy luận schema trong production.- Chỉ cần vài field nông →
get_json_object/json_tuple; cần cấu trúc lồng →from_json.to_jsonđể serialize ngược. explodebỏ dòng mảng rỗng/null;explode_outergiữ lại;posexplodekèm chỉ số. Explode nhân số dòng — explode càng muộn càng tốt.- Higher-order functions (
transform/filter/aggregate/exists) xử lý mảng không cần explode;collect_list/collect_setgom ngược. - Làm phẳng struct bằng
struct.*(một tầng);withField/dropFieldssửa field lồng (Spark 3.1+). - Schema evolution là cạm bẫy lớn: field lồng đổi tên/kiểu/tầng thường "im lặng thành null". Version hoá schema và kiểm tỉ lệ null/số dòng sau parse.
- Chiến lược medallion: bronze/silver giữ lồng (trung thực nguồn), gold làm phẳng theo nhu cầu phân tích — cân bằng storage, hiệu năng và trải nghiệm analyst.
Nguồn tham khảo
- Apache Spark — SQL Data Types (StructType, ArrayType, MapType)
- Apache Spark — Built-in Functions (explode, from_json, to_json, get_json_object, transform, filter, aggregate...)
- Apache Spark — JSON Files Data Source Guide (đọc/ghi JSON, multiLine, schema)
- PySpark API Reference — pyspark.sql.functions
- Databricks — Higher-order functions on Spark SQL (transform, filter, exists, aggregate)
- Sách: Jules Damji et al., Learning Spark, 2nd Edition (O'Reilly) — chương Spark SQL & DataFrame và complex/nested types.
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ẻ!