PySpark 11 — Cấu hình, spark-submit & Cluster
Từ laptop lên cluster
Suốt series, ta chạy PySpark bằng spark.master = "local[*]" — một JVM duy nhất giả lập cluster, tiện để học và test. Nhưng job thật của NCB xử lý hàng trăm triệu bản ghi mỗi đêm; nó phải chạy trên hàng chục máy với hàng trăm core. Bài này trả lời câu hỏi vận hành: làm sao đóng gói code Python + thư viện, gửi nó lên cluster, xin đúng lượng tài nguyên, và để nó chạy an toàn trong môi trường ngân hàng.
Đây là góc nhìn Python bổ sung cho Spark 8 — Production. Chỗ khó riêng của PySpark là: mã của bạn là Python, nhưng cluster nói JVM. Mỗi executor phải khởi động một tiến trình Python worker, và tiến trình đó cần đúng bộ thư viện bạn dùng trong code (pandas, numpy, thư viện nội bộ). Quản lý dependency Python trên cluster là vấn đề đặc thù mà người dùng Scala không gặp.
spark-submit — cửa ngõ vào cluster
spark-submit là script chuẩn để nộp một ứng dụng Spark. Cấu trúc lệnh:
spark-submit \
--master <cluster-url> \
--deploy-mode <client|cluster> \
--conf <key>=<value> \
--py-files <file.zip,...> \
--packages <group:artifact:version,...> \
main_job.py --arg1 val1 --arg2 val2
File cuối cùng (main_job.py) là entry point — script Python chứa SparkSession.builder.... Mọi tham số sau nó được truyền vào sys.argv của chính script, không phải cho Spark. Các cờ quan trọng:
--master — chỉ định cluster manager và địa chỉ:
| Giá trị | Ý nghĩa |
|---|---|
local[*] | Chạy 1 máy, dùng mọi core (chỉ để test) |
yarn | Nộp lên YARN (Hadoop) — phổ biến on-prem |
k8s://https://<api>:6443 | Spark on Kubernetes |
spark://host:7077 | Standalone cluster của Spark |
--deploy-mode — quyết định driver chạy ở đâu, và đây là điểm nhiều người nhầm:
client(mặc định): driver chạy ngay trên máy bạn gõ lệnh (edge node / máy CI). Executor vẫn chạy trên cluster. Log driver in thẳng ra terminal — tiện debug và cho notebook/REPL. Nhược điểm: máy đó phải giữ kết nối suốt job; nếu tắt SSH là job chết. Không hợp cho job đêm tự động.cluster: driver được đóng gói và chạy bên trong cluster như một container/process do manager cấp phát. Máy submit chỉ gửi rồi thoát. Bền vững, hợp cho production và scheduler (Airflow, Oozie). Nhược điểm: log driver nằm trên cluster, phải xem qua UI/yarn logs.
--py-files — đính kèm code Python phụ (module nội bộ) dưới dạng .zip, .egg hoặc .py. Spark phân phát chúng tới mọi executor và thêm vào PYTHONPATH. Đây là cách gửi package mylib/ của bạn lên cluster.
--archives — gửi file nén (.zip, .tar.gz) và giải nén trên executor. Dùng để đóng gói cả một virtualenv/conda env — chìa khoá giải bài toán dependency (mục dưới).
--packages — kéo thư viện Java/Scala từ Maven theo toạ độ group:artifact:version. Spark tự tải và phân phát. Ví dụ hay dùng ở ngân hàng:
--packages org.postgresql:postgresql:42.7.3,\
io.delta:delta-spark_2.12:3.2.0,\
org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.1
lần lượt là JDBC connector PostgreSQL, Delta Lake, và connector Kafka. --jars làm việc tương tự nhưng bạn tự trỏ tới file .jar cục bộ (khi cluster không ra được Internet — thường gặp ở ngân hàng, phải mirror Maven nội bộ hoặc dùng --jars).
--files — gửi file cấu hình/tài nguyên (vd log4j2.properties, file cert, file YAML tham số) tới thư mục làm việc của executor.
--conf — đặt bất kỳ thuộc tính spark.* nào; lặp lại cờ này nhiều lần cho nhiều tham số.
Cấu hình tài nguyên
Bạn khai báo cho cluster manager mỗi executor "to" bao nhiêu và cần mấy cái:
| Cờ | Ý nghĩa |
|---|---|
--num-executors N | Số executor (YARN). K8s dùng dynamic hoặc spark.executor.instances |
--executor-cores C | Số core mỗi executor = số task chạy song song trong nó |
--executor-memory M | Heap JVM mỗi executor (vd 8g) |
--driver-memory M | Heap của driver |
spark.executor.memoryOverhead | Bộ nhớ ngoài heap cho mỗi executor |
Điểm PySpark hay cháy: memoryOverhead. Vùng này (mặc định max(384MB, 0.1 × executor-memory)) chứa buffer shuffle, netty, và tiến trình Python worker. Vì PySpark chạy pandas/numpy/UDF trong Python nằm ngoài heap JVM, nếu code dùng nhiều bộ nhớ Python mà overhead nhỏ, YARN sẽ giết container với lỗi Container killed... exceeds physical memory. Với job PySpark nặng UDF/pandas, hãy nâng overhead lên 15–25% executor-memory.
Cân đối cores/memory. Đừng đặt executor-cores quá cao (>5): nhiều task cùng dùng chung heap gây tranh chấp GC và HDFS throughput giảm. Công thức thực dụng phổ biến: chọn executor-cores = 4–5, rồi tính số executor từ tổng core cluster (chừa 1 core + ~1GB cho daemon YARN mỗi node). Ví dụ node 16 core / 64GB, chừa 1 core → 15 core dùng được → 3 executor × 5 core, mỗi executor ~ (64−1)/3 ≈ 19GB, chia ra ~16GB heap + ~3GB overhead.
Dynamic allocation — thay vì cố định --num-executors, để Spark tự co giãn số executor theo tải:
--conf spark.dynamicAllocation.enabled=true \
--conf spark.dynamicAllocation.minExecutors=2 \
--conf spark.dynamicAllocation.maxExecutors=40 \
--conf spark.shuffle.service.enabled=true
Rất hợp cho cluster dùng chung: job xin thêm executor khi có nhiều task chờ, trả lại khi rảnh. Cần bật external shuffle service để executor bị thu hồi không làm mất dữ liệu shuffle.
Cluster manager: chọn cái nào
- YARN — quản lý tài nguyên của Hadoop, phổ biến nhất on-prem (nhiều ngân hàng VN có sẵn cụm Hadoop/Cloudera). Tích hợp Kerberos, quản lý queue theo phòng ban.
- Kubernetes — Spark chạy driver/executor thành pod. Ưu điểm: đóng gói dependency bằng Docker image (giải sạch bài toán Python env), hạ tầng thống nhất với microservice, scale linh hoạt. Nộp bằng
--master k8s://...với--conf spark.kubernetes.container.image=.... Xu hướng đang lên nhưng vận hành phức tạp hơn YARN. - Standalone — scheduler đơn giản đi kèm Spark, ít tính năng đa nhiệm; hợp cluster nhỏ chuyên dụng.
- Managed (đám mây) — Databricks, AWS EMR, GCP Dataproc: nhà cung cấp lo cluster, autoscale, tối ưu sẵn (Photon của Databricks). Bạn tập trung vào code, đổi lại chi phí cao hơn và phụ thuộc vendor. Databricks còn có Delta, Unity Catalog, notebook cộng tác — sẽ nói ở mục workflow.
Quản lý dependency Python trên executor
Đây là phần khiến PySpark khác biệt. Khi executor cần import pandas, thư viện đó phải có mặt trên node executor, không chỉ trên máy driver. Ba cách chính:
1. venv-pack / conda-pack — đóng gói môi trường ảo thành một archive, gửi qua --archives, giải nén trên executor:
# Trên máy build: tạo env rồi nén
conda create -y -n pysparkenv python=3.10 pandas numpy
conda activate pysparkenv
conda pack -o pysparkenv.tar.gz
# Khi submit
spark-submit \
--archives pysparkenv.tar.gz#environment \
--conf spark.pyspark.python=./environment/bin/python \
main_job.py
#environment là tên thư mục sau giải nén; spark.pyspark.python (hoặc biến PYSPARK_PYTHON) trỏ mọi executor dùng đúng Python trong archive đó. Cách này đảm bảo mọi node có y hệt bộ thư viện.
2. --py-files cho code nội bộ — đóng gói package .zip của bạn (không phải third-party nặng) và gửi kèm. Thường kết hợp: --archives cho env, --py-files cho mã dự án.
3. Docker image — trên Kubernetes, đưa toàn bộ thư viện vào image, mọi pod dùng chung. Sạch nhất về mặt tái lập, không cần đóng gói env lúc runtime.
Nguyên tắc cốt lõi: phiên bản Python và thư viện của driver phải khớp với executor, nếu không sẽ gặp lỗi pickling/serialization khó hiểu. Biến PYSPARK_PYTHON (executor) và PYSPARK_DRIVER_PYTHON (driver) cho phép trỏ tường minh.
Thứ tự ưu tiên cấu hình
Cùng một spark.* có thể đặt ở ba nơi, ưu tiên từ cao đến thấp:
- Trong code (
.config("spark.sql.shuffle.partitions", "400")) thắng tất cả — nhưng lưu ý: vài config chỉ có tác dụng nếu đặt trước khi SparkSession khởi tạo (config runtime của cluster như số executor thì phải đặt qua submit). --confhợp cho tham số phụ thuộc lần chạy, tách khỏi code.spark-defaults.confdo admin đặt mặc định chung cho toàn cluster.
Quy tắc thực dụng: cấu hình logic ứng dụng (shuffle partitions, AQE) đặt trong code cho version-control; cấu hình tài nguyên & môi trường (executor, master, packages) đặt ở submit để linh hoạt.
Tham số hoá job. Đừng hardcode ngày chạy hay đường dẫn. Nhận qua sys.argv/argparse rồi truyền lúc submit — để Airflow đưa {{ ds }} vào:
import argparse
p = argparse.ArgumentParser()
p.add_argument("--run-date", required=True) # vd 2026-07-13
p.add_argument("--input-path", required=True)
args = p.parse_args()
df = spark.read.parquet(f"{args.input_path}/dt={args.run_date}")
Logging. Spark dùng log4j2; gửi log4j2.properties qua --files rồi trỏ --conf spark.driver.extraJavaOptions=-Dlog4j.configurationFile=log4j2.properties. Trong Python, dùng logging chuẩn cho log nghiệp vụ; tránh print vì trên cluster nó biến mất trong stdout của executor.
Notebook & Databricks workflow
Trên notebook (Jupyter, Zeppelin, Databricks), bạn không gõ spark-submit. Môi trường đã tạo sẵn biến spark (một SparkSession client-mode gắn với cluster). Bạn viết code ô-này-ô-kia, mỗi action chạy ngay — hợp khám phá và phát triển.
Trên Databricks, có hai chế độ chạy:
- Interactive cluster + notebook: như trên, tương tác trực tiếp, cluster luôn bật.
- Job cluster (Workflows): Databricks tự spin up cluster riêng, chạy notebook/script như một job theo lịch, rồi tắt — tương đương
spark-submit --deploy-mode clusternhưng khai báo bằng UI/JSON thay vì lệnh. Dependency khai báo qua cluster libraries (PyPI/Maven) hoặc%pip installtrong notebook, không cần conda-pack.
Khác biệt cốt lõi: spark-submit cho bạn kiểm soát tuyệt đối và chạy được ở mọi nơi; notebook/Databricks đổi kiểm soát lấy tiện lợi và tích hợp. Nhiều đội phát triển trên notebook, đóng gói lại thành script + spark-submit/Job để chạy production nghiêm túc, có version control và test (xem testing PySpark).
Bảo mật trong môi trường ngân hàng
Job đọc/ghi dữ liệu khách hàng, credential không được lộ:
- Không hardcode user/password JDBC trong code hay trong
--conf(config hiện trong Spark UI và log!). Lấy từ secret manager: Databricks Secrets (dbutils.secrets.get(...)), Hashicorp Vault, hoặc biến môi trường của scheduler được mã hoá. - Kerberos/keytab trên cluster YARN doanh nghiệp: submit kèm
--principalvà--keytabđể job tự gia hạn ticket, không cần người ngồi nhập mật khẩu. - Phân quyền: job chạy dưới service account giới hạn quyền, chỉ truy cập đúng queue YARN và đúng schema/bảng của nghiệp vụ đó.
- Môi trường kiểm soát: submit qua edge node hoặc pipeline CI/CD có kiểm duyệt, không cho engineer gõ lệnh tuỳ ý lên production; log truy cập được audit.
Ví dụ: lệnh submit đầy đủ + cấu trúc project
Cấu trúc project chuẩn để đóng gói (minh hoạ):
ncb_etl/
├── main_job.py # entry point: parse args, build SparkSession
├── ncb_lib/ # package nội bộ
│ ├── __init__.py
│ ├── cleaning.py
│ └── aggregations.py
├── conf/
│ └── log4j2.properties
├── env/pysparkenv.tar.gz # env đóng gói bằng conda-pack
└── build.sh # zip ncb_lib/ -> ncb_lib.zip
Lệnh submit hoàn chỉnh (minh hoạ shell):
spark-submit \
--master yarn \
--deploy-mode cluster \
--name ncb_daily_txn_etl \
--num-executors 20 \
--executor-cores 4 \
--executor-memory 12g \
--driver-memory 8g \
--conf spark.executor.memoryOverhead=3g \
--conf spark.sql.shuffle.partitions=400 \
--conf spark.sql.adaptive.enabled=true \
--packages org.postgresql:postgresql:42.7.3,io.delta:delta-spark_2.12:3.2.0 \
--archives env/pysparkenv.tar.gz#environment \
--conf spark.pyspark.python=./environment/bin/python \
--py-files ncb_lib.zip \
--files conf/log4j2.properties \
--principal [email protected] \
--keytab /etc/security/svc_dp_etl.keytab \
main_job.py \
--run-date 2026-07-13 \
--input-path hdfs:///raw/cards
Trong main_job.py (minh hoạ Python), credential JDBC lấy từ biến môi trường do scheduler nạp, không hardcode:
import os
jdbc_props = {
"user": os.environ["JDBC_USER"],
"password": os.environ["JDBC_PASSWORD"], # từ Vault, không nằm trong code
"driver": "org.postgresql.Driver",
}
dim = spark.read.jdbc(os.environ["JDBC_URL"], "dwh.dim_customer", properties=jdbc_props)
Use case thực tế
Bài toán. Đội Data Platform NCB đưa job ETL giao dịch thẻ cuối ngày (đã xây trong ETL banking) từ môi trường dev-notebook lên chạy tự động trên cluster YARN Cloudera dùng chung với các phòng ban khác. Khối lượng ~180 triệu bản ghi/ngày, cửa sổ chạy 01:00–03:00.
Đóng gói. Code nghiệp vụ trong package ncb_lib được build.sh nén thành ncb_lib.zip (gửi bằng --py-files). Môi trường Python (pandas, numpy, thư viện chuẩn hoá mã KH nội bộ) đóng bằng conda-pack thành pysparkenv.tar.gz (~450MB), gửi bằng --archives; nhờ đó cả 20 executor dùng đúng Python 3.10 + đúng thư viện, hết cảnh "chạy được ở laptop, chết trên cluster". Connector JDBC PostgreSQL và Delta kéo qua --jars từ mirror Maven nội bộ (cluster không ra Internet).
Tài nguyên. Cấu hình như lệnh ví dụ: 20 executor × 4 core = 80 task song song, mỗi executor 12GB heap + 3GB overhead (nâng overhead vì có bước pandas_udf). Tổng chiếm ~20 × 15GB ≈ 300GB RAM và 80 core — nằm trong quota queue dp_etl. Với 180 triệu bản ghi và 400 shuffle partition (~450K bản ghi/partition), job hoàn tất trong ~35–45 phút, gọn trong cửa sổ đêm.
Vận hành & bảo mật. Airflow trên edge node gọi spark-submit --deploy-mode cluster, truyền --run-date {{ ds }}; driver chạy trong cluster nên task Airflow không giữ tài nguyên. Xác thực bằng keytab của service account svc_dp_etl, chỉ có quyền trên queue và schema được cấp; mật khẩu JDBC lấy từ Vault qua biến môi trường, không lộ trong Spark UI. Khi lỗi, Airflow retry idempotent (ghi đè partition dt= của ngày đó), không cần can thiệp tay.
Kết quả. Chuyển từ "chạy tay trên notebook mỗi sáng, hay quên thư viện" sang job đêm tự động, tái lập được, có audit — thời gian chờ số liệu sáng của đội BI giảm từ ~10:00 xuống trước 06:00.
Ghi nhớ
spark-submitlà cửa ngõ:--masterchọn cluster manager, entry point.pyở cuối, tham số sau nó vàosys.argv.--deploy-mode:client→ driver ở máy submit (debug/notebook);cluster→ driver trong cluster (production/scheduler). Job đêm luôn dùngcluster.- Đóng gói:
--py-filescho code nội bộ (.zip),--archivescho cả venv/conda env,--packages/--jarscho connector JDBC/Kafka/Delta,--filescho file cấu hình. - Bài toán Python đặc thù: mọi executor phải có đúng phiên bản Python + thư viện như driver — dùng conda-pack/venv-pack hoặc Docker image; trỏ
PYSPARK_PYTHON/spark.pyspark.python. - Tài nguyên:
executor-coresgiữ 4–5, nângmemoryOverhead(15–25%) cho job nặng pandas/UDF vì Python nằm ngoài heap; cân nhắc dynamic allocation trên cluster dùng chung. - Thứ tự ưu tiên config:
builder.config()>--conf>spark-defaults.conf. Tham số hoá ngày/đường dẫn qua argparse, đừng hardcode. - Bảo mật: không hardcode credential JDBC (config lộ trong UI/log) — dùng secret manager/Vault, keytab Kerberos, service account phân quyền, submit qua môi trường kiểm soát.
- Notebook/Databricks đổi kiểm soát lấy tiện lợi; production nghiêm túc vẫn nên đóng gói script + submit/Job có version control.
Xem thêm nền tảng ở PySpark tổng quan và vận hành sâu ở Spark 8 — Production.
Nguồn tham khảo
- Apache Spark Documentation — Submitting Applications (spark-submit)
- Apache Spark Documentation — Configuration
- Apache Spark Documentation — Running Spark on YARN
- Apache Spark Documentation — Running Spark on Kubernetes
- Apache Spark Documentation — Job Scheduling: Dynamic Resource Allocation
- PySpark Documentation — Python Package Management
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ẻ!