Batch 7 — Orchestration, Phụ thuộc & SLA

22 thg 7, 2026 2 lượt xem
#orchestration
#dag
#data-engineering
#sla
#batch
#sensor

Batch 7 — Orchestration, Phụ thuộc & SLA

Một pipeline batch "quy mô lớn" hiếm khi là một job đơn lẻ. Nó là hàng chục — đôi khi hàng trăm — bước: kéo dữ liệu từ core banking, giải nén, validate, join với danh mục, tính toán, ghi vào warehouse, refresh mart, gửi báo cáo. Mỗi bước phụ thuộc vào bước trước, chạy trên lịch khác nhau, có thể lỗi bất cứ lúc nào, và toàn bộ phải xong trước một mốc giờ cam kết. Orchestration chính là lớp lo bài toán đó: cái gì chạy, chạy khi nào, chạy sau cái gì, và làm gì khi hỏng.

Mô hình tinh thần: đừng nghĩ orchestrator là "cron cao cấp". Cron chỉ trả lời câu hỏi "đến giờ chưa?". Orchestrator trả lời thêm ba câu hỏi khó hơn: "các bước trước đã xong đúng chưa?", "dữ liệu nguồn đã sẵn sàng chưa?", và "nếu bước thứ 17 chết thì phải làm gì để không phải chạy lại từ đầu?". Bài này đi qua bốn trụ cột: DAG & phụ thuộc, lập lịch (time-based vs data-aware), xử lý lỗi & phục hồi, và SLA.


1. DAG — đồ thị phụ thuộc task

Nền tảng của mọi orchestrator hiện đại (Airflow, Dagster, Prefect) là DAG — Directed Acyclic Graph: đồ thị có hướng (cạnh A→B nghĩa là B chạy sau A) và không có chu trình (không được A→B→A, vì như vậy không có điểm bắt đầu). Node là task, cạnh là quan hệ phụ thuộc.

Tính "acyclic" không phải chi tiết học thuật — nó là điều kiện để tồn tại một topological order (thứ tự tô-pô): một cách xếp task thành hàng sao cho mọi phụ thuộc đều đứng trước. Không có thứ tự này thì scheduler không biết bắt đầu từ đâu. Đó cũng là lý do orchestrator từ chối DAG có vòng lặp ngay khi parse.

Vài tính chất quan trọng đọc ra từ DAG này:

  • Song song tự nhiên: extract_coreextract_cards không có cạnh nối nhau nên orchestrator chạy chúng đồng thời. build_mart_loanbuild_mart_deposit cũng vậy. Đây là nguồn tăng tốc lớn — bạn khai báo phụ thuộc thật, không thừa, để tối đa hoá song song.
  • Fan-in / fan-out: validate_schema là điểm fan-in (chờ 2 nguồn); load_fact_txn là điểm fan-out (toả ra 2 mart). Điểm fan-in thường là nơi cần trigger rule (chỉ chạy khi tất cả upstream thành công — mặc định all_success).
  • Critical path: đường dài nhất (theo thời gian) quyết định tổng thời lượng pipeline. Muốn về đích SLA sớm hơn, tối ưu task trên critical path, không phải task nhánh.

Task vs Asset — hai cách nghĩ về DAG

Có hai trường phái mô hình hoá:

Cách nghĩNode là gìVí dụ công cụƯu điểm
Task-centricHành động ("chạy script X")Airflow (cổ điển), PrefectTrực quan, hợp với job mệnh lệnh
Asset/data-centricDữ liệu được tạo ra ("bảng fact_txn ngày D")Dagster (Software-Defined Assets), Airflow AssetsLineage rõ, dễ lập lịch data-aware, dễ backfill

Trong mô hình asset, bạn khai báo "bảng này được sinh từ những bảng kia" và orchestrator suy ra DAG. Cách này gắn liền với lập lịch data-aware và làm backfill sạch hơn vì mỗi asset có partition rõ ràng.


2. Lập lịch — time-based vs data-aware

Một DAG cần biết khi nào khởi động. Có hai họ cơ chế, và pipeline thật thường trộn cả hai.

2.1. Time-based (cron / interval)

Cách kinh điển: chạy theo lịch giờ. Biểu thức cron năm trường phút giờ ngày tháng thứ là chuẩn de-facto.

# Chạy 2:00 sáng mỗi ngày (EOD batch)
0 2 * * *

# 15 phút một lần trong giờ hành chính T2-T6
*/15 8-17 * * 1-5

# 6:00 sáng ngày làm việc đầu mỗi tháng (báo cáo tháng)
0 6 1 * *

Hai cạm bẫy cần khắc ghi:

  • Timezone & DST: 0 2 * * * theo giờ nào? Nếu orchestrator chạy UTC còn nghiệp vụ theo giờ VN (UTC+7), "2h sáng" bị lệch 7 tiếng. Với ngân hàng, luôn khai báo timezone tường minh (Asia/Ho_Chi_Minh). DST không ảnh hưởng VN nhưng ảnh hưởng chi nhánh quốc tế — 2h sáng có thể xảy ra hai lần hoặc không lần nào trong ngày đổi giờ.
  • Data interval, không phải "thời điểm chạy": các orchestrator hiện đại phân biệt logical date / data interval (khoảng dữ liệu mà lần chạy này xử lý) với wall-clock time (giờ thực khi chạy). Run của ngày 2026-07-20 thường khởi động lúc cuối interval (rạng sáng 21) nhưng xử lý dữ liệu của ngày 20. Hiểu sai điểm này là nguồn lỗi off-by-one-day kinh điển trong report.

2.2. Data-aware (sensor / asset)

Cron giả định "cứ đến giờ là dữ liệu có sẵn". Thực tế thường không: file sao kê từ đối tác đến lúc 2:10 hôm nay nhưng 2:40 hôm khác; job EOD của core banking xong sớm/muộn tuỳ khối lượng. Nếu cứng nhắc chạy lúc 2:30, thỉnh thoảng bạn xử lý dữ liệu thiếu.

Data-aware scheduling đảo ngược: chỉ chạy khi dữ liệu nguồn thật sự sẵn sàng. Hai kỹ thuật:

  • Sensor (task-centric): một task đặc biệt chờ một điều kiện — file xuất hiện trên object storage, partition có trong Hive metastore, row đếm > 0, message đến topic. Sensor poll định kỳ đến khi điều kiện đúng thì "nhả" cho downstream.
  • Asset/Dataset trigger (data-centric): DAG B khai báo "tôi phụ thuộc asset X"; khi DAG A cập nhật X xong, orchestrator tự kích DAG B. Không cần poll — đây là mô hình event-driven giữa các DAG.

Mẹo vận hành: sensor kiểu poll cổ điển giữ một worker slot suốt thời gian chờ — 100 sensor chờ đồng thời có thể "ăn" cạn pool và làm nghẽn cả cluster. Airflow giải bằng chế độ reschedule (nhả slot giữa các lần poll) và deferrable operators / triggerer (chờ bất đồng bộ, gần như không tốn worker). Luôn ưu tiên deferrable cho sensor chờ lâu.

Điểm mấu chốt trong sơ đồ: sensor phải có timeout. Chờ vô hạn là phản pattern — nếu nguồn không bao giờ đến, bạn muốn fail nhanh và cảnh báo, chứ không muốn một task treo âm thầm che mất sự cố thượng nguồn.


3. Retry, backoff & timeout — chống lỗi thoáng qua

Batch quy mô lớn sống trong môi trường không đáng tin cậy tạm thời: mạng chớp tắt, database bận, credential hết hạn 3 giây, node bị evict. Phần lớn các lỗi này tự khỏi nếu thử lại. Đó là lý do mọi task nên khai báo chính sách retry.

  • retries: số lần thử lại (thường 2–3 cho task I/O mạng).
  • retry backoff: đừng thử lại ngay lập tức. Dùng exponential backoff (chờ 1', 2', 4', 8'...) — nếu database đang quá tải, dồn retry vào cùng lúc chỉ làm nó chết thêm. Nhiều orchestrator hỗ trợ exponential_backoff=True kèm max_retry_delay. Thêm jitter (ngẫu nhiên hoá) để tránh "thundering herd" khi nhiều task cùng retry đồng bộ.
  • timeout (execution timeout): trần thời gian một task được phép chạy. Không có timeout, một query treo có thể ngốn cả cửa sổ batch và kéo sập SLA mà không ai biết. Đặt timeout hơi lớn hơn p99 thời lượng bình thường.
# Pseudocode trung lập công cụ — chính sách retry cho một task
default_task_policy = {
    "retries": 3,
    "retry_delay": minutes(2),
    "retry_exponential_backoff": True,   # 2', 4', 8'
    "max_retry_delay": minutes(30),
    "execution_timeout": minutes(60),    # cắt task nếu treo > 60'
    "on_failure_callback": alert_oncall, # gọi khi cạn retry
}

# Task gọi API bên ngoài (dễ lỗi thoáng qua) → retry nhiều
fetch_fx_rate = Task(
    fn=call_fx_api,
    retries=5,
    retry_exponential_backoff=True,
)

# Task ghi vào bảng fact (tốn kém, phải idempotent) → retry ít, cẩn trọng
load_fact = Task(
    fn=load_fact_txn,   # dùng MERGE/overwrite-partition, xem batch-04
    retries=2,
)

Cảnh báo sống còn: retry chỉ an toàn khi task idempotent. Nếu load_factINSERT thuần, retry sau khi nó đã ghi một nửa sẽ nhân đôi dữ liệu. Đây là lý do orchestration và idempotency (batch-04) là hai mặt của một đồng xu: retry tự động của orchestrator giả định mỗi task chạy lại cho ra cùng kết quả. Hãy dùng overwrite-by-partition, MERGE, hoặc natural key + dedup — đừng bao giờ để orchestrator retry một INSERT append thô.

Phân loại lỗi: transient vs permanent

Không phải lỗi nào cũng nên retry. Lỗi permanent (schema sai, chia cho 0, file corrupt, quyền bị thu hồi) sẽ thất bại y hệt ở mọi lần thử — retry chỉ trì hoãn cảnh báo và đốt tài nguyên. Task tốt nên phân biệt: ném lỗi fail-fast cho lỗi logic (không retry) và chỉ retry cho lỗi hạ tầng/mạng. Một số framework có ngoại lệ riêng để "fail ngay, bỏ qua retry".


4. Xử lý lỗi & phục hồi — checkpoint, chạy lại từ task lỗi

Retry lo lỗi một task. Nhưng khi cả một run 40 task chết ở task 17, bạn cần một câu trả lời khác: không chạy lại từ đầu.

4.1. Chạy lại từ điểm lỗi (partial re-run)

Vì DAG lưu trạng thái từng task (success/failed/skipped), orchestrator cho phép clear + rerun chỉ những task faileddownstream của chúng. 16 task trước đó vẫn success, không đụng lại. Điều này chỉ đúng nếu output của các task đã thành công còn nguyên vẹn (thường được ghi ra storage bền, không chỉ nằm trong RAM của run).

4.2. Checkpoint bên trong task

Với task dài (xử lý 500 triệu dòng), lỗi ở phút thứ 55/60 mà phải chạy lại toàn bộ là lãng phí. Hai lớp checkpoint:

  • Checkpoint cấp task (orchestrator): cắt job to thành nhiều task nhỏ theo partition/ngày. Mỗi partition là một task → lỗi partition nào, rerun partition đó. Đây là "dynamic task mapping" / partitioned assets.
  • Checkpoint cấp engine (bên trong task): các engine như Spark Structured Streaming hay job có checkpointLocation tự lưu tiến độ; batch Spark ghi theo partition để rerun ghi đè an toàn. Orchestrator không thay được lớp này — nó thuộc về xử lý idempotent của chính task.

4.3. Bù trạng thái & dọn dẹp

Task lỗi giữa chừng có thể để lại rác: file tạm, bảng staging nửa vời, lock chưa nhả. Thiết kế tốt tách staging → publish: task ghi vào bảng/thư mục tạm, chỉ khi thành công mới swap/rename atomic sang vùng chính (write-audit-publish). Nếu lỗi, vùng chính chưa bị đụng → chạy lại sạch sẽ.


5. SLA — cam kết thời gian & xử lý trễ

SLA (Service Level Agreement) trong batch trả lời: "báo cáo/bảng này phải sẵn sàng trước mấy giờ?". Với ngân hàng đây là ràng buộc thật: sổ phụ phải xong trước giờ mở cửa quầy, báo cáo phân loại nợ phải nộp NHNN trước hạn.

Phân biệt ba khái niệm hay bị gộp:

Khái niệmNghĩaVí dụ
LatencyThời gian pipeline chạy hết45 phút
SLA / deadlineMốc giờ tuyệt đối phải xong"trước 05:00"
SLA missXong sau deadlinexong 05:20 → miss 20'

Điểm tinh tế: SLA gắn với deadline tuyệt đối, không phải "chạy trong bao lâu". Một job 10 phút vẫn miss SLA nếu nó khởi động lúc 04:55 mà deadline là 05:00 — vì nguồn đến muộn. Vì thế theo dõi SLA phải tính cả thời gian chờ nguồn, không chỉ thời gian tính toán.

Cơ chế SLA của orchestrator

  • SLA callback / deadline alert: khai báo mỗi task/DAG phải xong trong X thời gian (hoặc trước mốc Y); quá hạn thì orchestrator kích callback cảnh báo — nhưng không tự dừng task. SLA thường là tín hiệu quan sát, khác với execution_timeout (thực sự giết task).
  • Cảnh báo phân tầng: cảnh báo dự đoán miss (khi critical path đã trễ so với mốc trung gian) quý hơn cảnh báo sau khi đã miss. Kỹ thuật: đặt deadline trung gian ("bảng fact phải xong trước 04:00 để mart kịp 05:00") và alert ngay ở mốc trung gian.
  • Alerting đa kênh: Slack/Teams cho cảnh báo thường, PagerDuty/gọi điện cho SLA nghiệp vụ trọng yếu. Gắn runbook vào alert: cảnh báo mà không có hướng dẫn xử lý chỉ tạo mệt mỏi cho on-call.

Chi tiết SLA/downtime nâng cao (SLI, error budget cho pipeline) nằm ở SLA & downtime trong DataOps. Ở đây ta chỉ dừng ở góc orchestration.


6. Công cụ — trung lập, nhưng cùng một mô hình

Airflow, Dagster, Prefect khác cú pháp nhưng cùng khung khái niệm đã trình bày: DAG, phụ thuộc, lập lịch, retry, sensor/asset, SLA.

Khía cạnhAirflowDagsterPrefect
Đơn vị cốt lõiTask trong DAGSoftware-Defined AssetTask/Flow
Data-awareAssets (Datasets)Asset-native (mặc định)Automations/events
Chờ nguồnSensor / deferrableSensor, asset checksTrigger, wait
Backfilltheo data intervalpartition-nativedeployment
Điểm mạnhhệ sinh thái lớn, chínlineage & test dữ liệuPythonic, dynamic

Nếu bạn đi sâu Airflow, xem series Airflow chuyên sâu — đặc biệt DAG & TaskLập lịch, Timetables & Assets cho phần data-aware. Việc "orchestration đo chất lượng trước khi publish" (quality gate) được nối sang Batch 8 — Data Quality: task validate là một node trong DAG, và nếu nó fail thì downstream load/publish bị chặn (fail-closed).


Use case thực tế

Pipeline EOD tại NCB (số liệu minh hoạ). Batch cuối ngày gồm ~60 task: kéo dữ liệu core banking T24, giao dịch thẻ, dữ liệu khoản vay; validate; tính lãi dồn tích; phân loại nợ (nhóm 1–5); dựng mart; refresh dashboard rủi ro. Deadline nghiệp vụ: bảng phân loại nợ sẵn sàng trước 06:00 để khối QLRR chốt số.

  • Data-aware thay cron cứng: trước đây job chạy cứng 02:30; ~2 lần/tháng file EOD của T24 đến sau 02:30 → xử lý thiếu, phải chạy tay. Chuyển sang sensor deferrable chờ file eod_YYYY-MM-DD.parquet với timeout=04:00. Kết quả (minh hoạ): sự cố dữ liệu thiếu về ~0; slot worker tiết kiệm vì không poll đồng bộ.
  • Partial re-run: một đêm task build_mart_loan OOM ở phút 50. Nhờ WAP + partitioned task, on-call chỉ clear + rerun nhánh mart (12 task) với memory cao hơn, không đụng 45 task extract/transform đã xong → phục hồi trong ~25' thay vì chạy lại ~2h.
  • SLA phân tầng: đặt mốc trung gian fact_txn xong trước 04:30. Một đêm nguồn đến muộn (03:50), critical path trễ; alert dự đoán bắn lúc 04:35 cho on-call trước khi miss deadline 06:00 → kịp cấp thêm compute, về đích 05:40. Không có mốc trung gian, đội chỉ biết khi đã miss.

Ghi nhớ

  • DAG = phụ thuộc, không phải thứ tự tay viết: khai báo cạnh thật để orchestrator tối đa song song; acyclic là điều kiện để có topological order.
  • Cron trả lời "đến giờ chưa", data-aware trả lời "dữ liệu sẵn sàng chưa" — pipeline thật thường cần cả hai; ưu tiên sensor/asset cho nguồn đến giờ bất định.
  • Sensor phải có timeout + chạy deferrable/reschedule: chờ vô hạn che sự cố thượng nguồn và ngốn worker slot.
  • Retry + exponential backoff + jitter cho lỗi thoáng qua; fail-fast cho lỗi permanent; đặt execution_timeout để job treo không nuốt cửa sổ batch.
  • Retry chỉ an toàn khi task idempotent — orchestration và idempotency là hai mặt một đồng xu.
  • Phục hồi = chạy lại từ task lỗi, không từ đầu; thiết kế WAP (staging → publish atomic) và partition hoá để rerun sạch.
  • SLA gắn deadline tuyệt đối, không phải latency; theo dõi cả thời gian chờ nguồn; alert dự đoán miss quý hơn alert sau khi miss, và luôn kèm runbook.

Nguồn tham khảo

  • Reis & Housley, Fundamentals of Data Engineering (O'Reilly) — chương Orchestration & undercurrents.
  • Kleppmann, Designing Data-Intensive Applications (O'Reilly) — batch processing, dataflow & fault tolerance.
  • Apache Airflow — Documentation: DAGs, Scheduling & Timetables, Sensors/Deferrable Operators, SLAs & callbacks (https://airflow.apache.org/docs/).
  • Dagster — Documentation: Software-Defined Assets, Schedules & Sensors, Partitions & Backfills (https://docs.dagster.io/).
  • Prefect — Documentation: Flows, Tasks, Retries & Automations (https://docs.prefect.io/).
  • Google — "The Site Reliability Engineering Book", chương về SLO/SLA & alerting (https://sre.google/books/).

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