PySpark 3 — Đọc/ghi dữ liệu & Schema

14 thg 7, 2026 4 lượt xem
#data-engineering
#parquet
#schema
#pyspark
#io

Sau khi đã quen với DataFrame API ở bài DataFrame API, câu hỏi tiếp theo của mọi pipeline là: dữ liệu vào từ đâu và ra ở đâu? Trong thực tế ngân hàng, một job Spark thường đọc từ nhiều nguồn khác nhau — bảng core qua JDBC, file CSV do đối tác gửi, log JSON, bảng lakehouse — biến đổi, rồi ghi ra kho phân tích. Bài này đi sâu vào lớp đọc/ghi (I/O) và schema — hai thứ quyết định pipeline của bạn nhanh, đúng và bền hay không. Nếu chưa nắm kiến trúc chung, xem lại PySpark tổng quan.

Mô hình đọc/ghi thống nhất

Spark trừu tượng hóa mọi nguồn dữ liệu qua hai đầu vào/ra chung là DataFrameReader (spark.read) và DataFrameWriter (df.write). Cùng một cú pháp áp dụng cho mọi định dạng — chỉ đổi formatoptions:

# Đọc: đủ dạng đầy đủ
df = (spark.read
        .format("parquet")          # hoặc csv, json, orc, avro, jdbc, delta...
        .option("mergeSchema", "true")
        .load("s3a://ncb-lake/raw/txn/"))

# Ghi
(df.write
   .format("parquet")
   .mode("overwrite")
   .save("s3a://ncb-lake/curated/txn/"))

Có các "shortcut" tiện tay: spark.read.parquet(path), df.write.csv(path)… tương đương dạng format(...).load(...). Toàn bộ đều lazy — việc đọc thật sự chỉ xảy ra khi có action.

Các định dạng file

Parquet — lựa chọn mặc định

Parquet là định dạng cột (columnar), nén tốt, có schema nhúng sẵn trong footer của file. Đây là định dạng nên ưu tiên cho mọi dữ liệu trung gian và kho phân tích vì:

  • Lưu theo cột: chỉ đọc đúng cột cần (column pushdown), không quét cả hàng.
  • Nén cao: mặc định snappy (nhanh); có thể chọn zstd/gzip cho tỷ lệ nén cao hơn.
  • Thống kê per row-group (min/max/count) cho phép predicate pushdown — bỏ qua cả khối dữ liệu không thỏa điều kiện lọc.
  • Schema tự mô tả: không cần khai báo lại khi đọc.
df = spark.read.parquet("s3a://ncb-lake/curated/txn/")
(df.write.option("compression", "zstd")
   .mode("overwrite").parquet("s3a://ncb-lake/out/"))

ORC tương tự Parquet (cột, nén, thống kê) và rất phổ biến trong hệ sinh thái Hive; Avro là định dạng theo hàng, hợp cho serialize sự kiện/streaming (schema linh hoạt). Với data lake phân tích, Parquet gần như luôn là mặc định.

CSV và JSON — dữ liệu trao đổi

CSV/JSON hay đến từ đối tác, hệ thống ngoài, export thủ công. Chúng không có schema và không tối ưu cho quét lớn, nên thường chỉ dùng ở tầng raw rồi chuyển ngay sang Parquet.

df_csv = (spark.read.format("csv")
            .option("header", "true")
            .option("delimiter", "|")
            .option("nullValue", "")
            .option("dateFormat", "yyyy-MM-dd")
            .schema(txn_schema)          # nên đưa schema tường minh
            .load("/landing/partner/*.csv"))

df_json = (spark.read.option("multiLine", "true")
             .json("/landing/events/"))

Các option hay dùng của CSV: header, inferSchema, delimiter/sep, quote, escape, nullValue, dateFormat, encoding. Của JSON: multiLine, primitivesAsString. Định dạng text đọc mỗi dòng thành một cột value — hữu ích cho log thô cần tự parse.

Đọc nhiều file, thư mục và partition

Đường dẫn có thể là một thư mục, một danh sách, hay có wildcard — Spark tự gộp tất cả file bên trong thành một DataFrame:

spark.read.parquet("/lake/txn/")                 # cả thư mục
spark.read.parquet("/lake/txn/dt=2026-07-13/")   # một partition
spark.read.parquet("/lake/txn/dt=2026-07-*/")    # wildcard

Nếu dữ liệu được ghi theo partition discovery (cấu trúc thư mục dt=2026-07-13/), Spark tự nhận cột dt như một cột phân vùng và cho phép lọc theo nó mà không đọc thư mục khác — đây là partition pruning, nền tảng của việc đọc nhanh dữ liệu lịch sử.

Schema — trái tim của I/O đúng đắn

Vì sao KHÔNG nên inferSchema với dữ liệu lớn

Với CSV/JSON, inferSchema=true bảo Spark tự đoán kiểu cột. Nghe tiện, nhưng có hai vấn đề nghiêm trọng ở quy mô lớn:

  1. Tốn một lượt quét toàn bộ dữ liệu chỉ để suy luận kiểu — trước cả khi làm việc thật. Với terabyte dữ liệu, đây là chi phí khổng lồ và lặp lại mỗi lần đọc.
  2. Dễ đoán sai và không ổn định. Một cột số tài khoản có số 0 đầu (0012345) có thể bị suy thành số nguyên và mất số 0; cột tiền có ô rỗng có thể bị coi là string; giá trị đặc biệt ở cuối file có thể đổi kiểu cả cột. Kết quả: schema thay đổi giữa các lần chạy → pipeline gãy khó lường.

Quy tắc: luôn khai báo schema tường minh cho dữ liệu sản xuất. inferSchema chỉ nên dùng khi khám phá nhanh một file lạ.

Định nghĩa StructType/StructField

from pyspark.sql.types import (StructType, StructField, StringType,
    LongType, DecimalType, TimestampType, DateType)

txn_schema = StructType([
    StructField("txn_id",     StringType(),      False),
    StructField("account_no", StringType(),      False),  # giữ số 0 đầu
    StructField("amount",     DecimalType(18, 2), True),  # tiền: dùng Decimal
    StructField("kind",       StringType(),      True),
    StructField("txn_ts",     TimestampType(),   True),
    StructField("dt",         DateType(),        True),
])

Vài lưu ý nghiệp vụ: số tài khoản, CIF, số thẻ luôn là String (không phải số) để giữ số 0 đầu và tránh tràn; số tiền dùng DecimalType chứ không DoubleType để tránh sai lệch dấu phẩy động khi cộng dồn. Đối số thứ ba (nullable) khai báo cột có được phép null hay không.

Đọc schema từ Parquet

Vì Parquet nhúng schema sẵn, khi cần tái sử dụng bạn có thể đọc một file mẫu rồi lấy schema đó áp cho nguồn khác — tránh gõ tay lại:

ref_schema = spark.read.parquet("/lake/txn/dt=2026-07-13/").schema
df = spark.read.schema(ref_schema).csv("/landing/txn_backfill/")
print(df.schema.json())   # xuất JSON để version-control schema

Xử lý bản ghi lỗi

Dữ liệu thật luôn bẩn: dòng thiếu cột, sai kiểu, ký tự lạ. Reader có option mode kiểm soát cách xử lý (áp dụng cho CSV/JSON):

modeHành vi
PERMISSIVE (mặc định)Giữ dòng lỗi, đặt các trường không parse được thành null, gom bản ghi thô vào cột columnNameOfCorruptRecord
DROPMALFORMEDBỏ luôn dòng lỗi, không báo
FAILFASTNém lỗi và dừng job ngay khi gặp dòng đầu tiên hỏng
df = (spark.read.format("csv").schema(txn_schema)
        .option("mode", "PERMISSIVE")
        .option("columnNameOfCorruptRecord", "_corrupt")
        .load("/landing/partner/*.csv"))

# tách và đếm dòng lỗi để cảnh báo
bad = df.filter("_corrupt is not null")
print("Số dòng lỗi:", bad.count())

Nguyên tắc vận hành: với dữ liệu tài chính, thà FAILFAST hoặc tách riêng dòng lỗi ra "quarantine" để rà soát, còn hơn âm thầm DROPMALFORMED làm mất giao dịch mà không ai biết. Việc kiểm soát chất lượng này gắn với data quality.

Ghi dữ liệu

mode — cách ứng xử khi đích đã tồn tại

modeÝ nghĩa
appendThêm dữ liệu mới vào đích
overwriteXóa và ghi đè toàn bộ (hoặc partition liên quan nếu bật dynamic overwrite)
error/errorifexists (mặc định)Báo lỗi nếu đích đã có
ignoreKhông làm gì nếu đích đã có

partitionBy — phân vùng theo cột

partitionBy chia dữ liệu ra các thư mục con theo giá trị cột, tạo cấu trúc dt=2026-07-13/. Điều này giúp truy vấn sau chỉ đọc đúng partition cần (partition pruning). Nên phân vùng theo cột lọc thường xuyên và có lực phân giải vừa phải — ngày giao dịch, chi nhánh, loại sản phẩm.

(df.write.mode("overwrite")
   .partitionBy("dt")               # 1 thư mục/ngày
   .parquet("s3a://ncb-lake/curated/txn/"))

Cảnh báo over-partition: đừng phân vùng theo cột có quá nhiều giá trị (như account_no, timestamp giây). Hàng triệu thư mục, mỗi thư mục vài KB, sẽ tạo ra "small files problem" — driver quá tải liệt kê file, đọc chậm thảm hại. Kinh nghiệm: mỗi partition nên đủ lớn (hàng trăm MB tới vài GB).

bucketBy — chia xô

bucketBy băm dữ liệu vào số "xô" cố định theo cột khóa, giúp các join/aggregation sau tránh shuffle lại. Chỉ dùng được khi ghi ra bảng managed (saveAsTable), không dùng với save(path) đơn thuần:

(df.write.mode("overwrite")
   .bucketBy(32, "account_no").sortBy("account_no")
   .saveAsTable("curated.txn_bucketed"))

Kiểm soát số file: repartition vs coalesce

Số file output = số partition của DataFrame khi ghi. Quá nhiều file nhỏ hại cả ghi lẫn đọc; quá ít file lớn thì mất song song. Hai công cụ:

  • repartition(n) / repartition("dt"): xáo trộn lại (shuffle) để đạt đúng n partition hoặc phân bố đều theo cột — dùng khi cần tăng số partition hoặc cân bằng dữ liệu lệch.
  • coalesce(n): gộp partition mà không shuffle — rẻ hơn, nhưng chỉ để giảm số partition.
# gộp còn ~1 file/partition ngày trước khi ghi
(df.repartition("dt").write.partitionBy("dt")
   .parquet("s3a://ncb-lake/curated/txn/"))

Chi tiết cơ chế shuffle và cách chọn số partition tối ưu xem shuffle & partitions.

Kết nối nguồn ngoài

JDBC — đọc từ database quan hệ

Spark đọc trực tiếp từ Postgres/Oracle/SQL Server qua JDBC. Điểm mấu chốt: mặc định Spark dùng một kết nối duy nhất, kéo cả bảng qua một task — cực chậm và dễ nghẽn. Để song song hóa, khai báo partitionColumn, lowerBound, upperBound, numPartitions: Spark tự sinh nhiều truy vấn WHERE partitionColumn BETWEEN ... chạy song song.

jdbc = (spark.read.format("jdbc")
   .option("url", "jdbc:postgresql://core-db:5432/corebank")
   .option("dbtable", "(SELECT * FROM txn WHERE dt = '2026-07-13') t")
   .option("user", user).option("password", pwd)
   .option("driver", "org.postgresql.Driver")
   # song song hóa
   .option("partitionColumn", "txn_id_num")
   .option("lowerBound", 1).option("upperBound", 50_000_000)
   .option("numPartitions", 16)
   .option("fetchsize", 10_000)
   .load())

Vài lưu ý sống còn khi đọc từ core banking:

  • Cảnh báo tải nguồn: mỗi partition là một truy vấn thật lên DB sản xuất. numPartitions=16 nghĩa là 16 truy vấn đồng thời — có thể đè sập core. Phối hợp với DBA, chạy ngoài giờ cao điểm, hoặc đọc từ read-replica.
  • dbtable có thể là một subquery để đẩy lọc/chọn cột xuống DB (chỉ kéo dữ liệu cần).
  • partitionColumn nên là cột số/khóa phân bố đều để tránh partition lệch.

Object storage — S3, HDFS

Với data lake, đường dẫn dùng scheme s3a:// (AWS/MinIO) hay hdfs://. Cấu hình credential/endpoint qua spark.hadoop.*. Đọc/ghi Parquet trên object storage là mô hình chuẩn của lakehouse hiện đại.

Bảng lakehouse — Delta Lake & Iceberg

Parquet thuần không có transaction, không time-travel, ghi đè một phần dễ hỏng. Các table format giải quyết điều đó bằng lớp metadata trên Parquet, cho ACID, cập nhật/xóa, và du hành thời gian:

# Delta Lake
df.write.format("delta").mode("append").save("/lake/txn_delta")
spark.read.format("delta").option("versionAsOf", 12).load("/lake/txn_delta")

# Iceberg (qua catalog đã cấu hình)
df.writeTo("catalog.curated.txn").append()
spark.read.table("catalog.curated.txn")

Xem sâu ở Delta Lake, Delta lakehouse với Spark, Iceberg tổng quan và kiến trúc chung ở lakehouse.

Predicate & column pushdown với Parquet

Đây là lý do Parquet nhanh vượt trội. Khi bạn viết:

(spark.read.parquet("/lake/txn/")
   .select("txn_id", "amount")          # column pushdown
   .filter("dt = '2026-07-13' and amount > 100000000"))  # predicate pushdown

Spark không đọc cả file rồi mới lọc. Nhờ schema cột và thống kê min/max trong footer:

  • Column pushdown: chỉ đọc hai cột txn_id, amount; các cột khác không chạm đĩa.
  • Predicate pushdown: bỏ qua cả row-group mà min/max cho biết không có dòng nào thỏa amount > 100tr.
  • Partition pruning: nếu phân vùng theo dt, chỉ mở đúng thư mục ngày đó.

Bài học thực hành: luôn select đúng cột cần và lọc sớm; tránh select("*") khi bảng rộng. Với JDBC, cơ chế tương đương là đẩy WHERE/SELECT vào subquery dbtable.

Use case thực tế

Bối cảnh NCB: Đội Data cần nạp dữ liệu giao dịch hằng ngày từ hệ thống core (Postgres/Oracle) sang data lake để phục vụ báo cáo và mô hình rủi ro. Giả sử một ngày phát sinh ~30 triệu giao dịch (số liệu ước lượng minh họa), bảng txn ở core có ~40 cột nhưng phân tích chỉ cần 12 cột.

Luồng xử lý (chạy 02:00 mỗi đêm, ngoài giờ cao điểm):

  1. Đọc JDBC song song từ core: dùng subquery dbtable lọc đúng dt = hôm qua và chỉ chọn 12 cột cần (column pushdown xuống DB). Đặt partitionColumn = txn_seq (khóa số tự tăng), numPartitions = 12 — tương đương 12 truy vấn đồng thời, đã thống nhất giới hạn với DBA để không đè core. Với ~30 triệu dòng, mỗi task đọc ~2,5 triệu dòng.
  2. Áp schema tường minh: account_no là String (giữ số 0 đầu), amountDecimalType(18,2), txn_ts là Timestamp, thêm cột dt (DateType). Không dùng inferSchema.
  3. Biến đổi nhẹ: chuẩn hóa mã loại giao dịch, lọc bản ghi test, thêm cột dt để phân vùng.
  4. Ghi Parquet phân vùng theo ngày: repartition("dt") rồi write.partitionBy("dt") mode overwrite ở phạm vi partition ngày đó — tạo /lake/curated/txn/dt=2026-07-13/ với vài file ~256 MB thay vì hàng nghìn file nhỏ. Nén zstd.

Kết quả ước lượng: 30 triệu dòng × 40 cột CSV ~ vài chục GB, sau khi chọn 12 cột và nén Parquet chỉ còn vài GB — tiết kiệm dung lượng rõ rệt. Báo cáo cuối ngày truy vấn "tổng chi tiêu theo chi nhánh ngày hôm qua" chỉ quét đúng một partition dt, đọc 3 cột, hoàn tất trong vài giây thay vì quét toàn bảng. Về sau có thể nâng cấp đích sang Delta/Iceberg để có ACID và time-travel mà không đổi logic đọc/ghi.

Câu SQL phân tích tương ứng (minh họa nghiệp vụ, chạy trên kho — không đánh dấu chạy được):

SELECT dt, branch_id, COUNT(*) AS so_gd, SUM(amount) AS tong_tien
FROM   curated.txn
WHERE  dt = DATE '2026-07-13'
GROUP  BY dt, branch_id;

Ghi nhớ

  • Một mô hình I/O thống nhất: spark.read.format(...).load()df.write.format(...).save() cho mọi định dạng; tất cả đều lazy.
  • Parquet là mặc định: định dạng cột, nén, schema tự mô tả, hỗ trợ predicate/column pushdown. CSV/JSON chỉ nên ở tầng raw rồi chuyển sang Parquet.
  • Luôn khai báo schema tường minh (StructType/StructField) cho dữ liệu sản xuất; tránh inferSchema vì tốn một lượt quét và dễ đoán sai. Số tài khoản/CIF là String, tiền là DecimalType.
  • Xử lý bản ghi lỗi có chủ đích: chọn PERMISSIVE + columnNameOfCorruptRecord để tách quarantine, hoặc FAILFAST cho dữ liệu tài chính — đừng âm thầm DROPMALFORMED.
  • Khi ghi: chọn mode đúng; partitionBy theo cột lọc thường xuyên nhưng tránh over-partition; dùng repartition/coalesce kiểm soát số file, tránh small files.
  • JDBC phải song song hóa bằng partitionColumn/numPartitions, nhưng luôn cân nhắc tải lên DB nguồn — phối hợp DBA, chạy ngoài giờ, ưu tiên read-replica.
  • Bảng lakehouse (Delta/Iceberg) thêm ACID và time-travel trên nền Parquet mà giữ nguyên cú pháp đọc/ghi quen thuộc.
  • Lọc và chọn cột sớm để tận dụng pushdown và partition pruning — nhanh nhất là dữ liệu bạn không cần đọc.

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ẻ!