Python hiện đại 8 — Ghép pipeline dữ liệu ngân hàng thực chiến
Suốt series Python hiện đại cho Data, chúng ta học từng công cụ riêng lẻ: Polars xử lý DataFrame tốc độ Rust, DuckDB truy vấn SQL trên Parquet, Pydantic validate dữ liệu, async/httpx ingest song song, uv/ruff quản lý môi trường, pytest kiểm thử. Bài cuối này làm điều quan trọng nhất — ghép tất cả lại thành một pipeline dữ liệu ngân hàng chạy được thật, hằng ngày, có thể bàn giao cho vận hành.
Bài toán: đóng sổ dữ liệu cuối ngày
Mỗi cuối ngày làm việc, một ngân hàng như NCB phải gom dữ liệu giao dịch từ nhiều nguồn: file sao kê CSV từ core banking, file Parquet từ hệ thống thẻ, feed từ đối tác trung gian thanh toán (ví điện tử, cổng QR) qua API. Số lượng file dao động vài chục đến vài trăm, tổng cỡ vài GB đến vài chục GB. Yêu cầu nghiệp vụ:
- Ingest: kéo hết dữ liệu từ các nguồn về, càng nhanh càng tốt vì cửa sổ đóng sổ hẹp.
- Validate: bắt bản ghi hỏng (thiếu trường, số tiền âm bất thường, mã tiền tệ sai) và tách riêng để không làm ô nhiễm báo cáo — nhưng không được làm rơi âm thầm.
- Transform & tổng hợp: chuẩn hóa, khử trùng lặp, tính số dư/doanh số theo chi nhánh, loại tiền, kênh.
- Xuất & nạp: ghi Parquet vào vùng lưu trữ (data lake) và tạo báo cáo cho nghiệp vụ.
- Đảm bảo chất lượng & tái chạy được: nếu lỗi giữa chừng, chạy lại không tạo dữ liệu nhân đôi (idempotent).
Đây chính là chỗ mọi mảnh ghép của series khớp vào nhau.
Kiến trúc end-to-end
Từng bước dưới đây gọi lại đúng bài đã học trong series.
1. Môi trường tái lập và chất lượng code
Nền móng của pipeline production là tái lập được (reproducible). Dùng uv để khóa phụ thuộc và ruff để chuẩn hóa code: một file pyproject.toml khai báo phụ thuộc, uv.lock ghim phiên bản chính xác. Máy dev, máy CI, và server chạy pipeline đều dựng môi trường từ cùng một lock — không còn cảnh "chạy được ở máy tôi".
# pyproject.toml (rút gọn)
[project]
name = "eod-pipeline"
requires-python = ">=3.11"
dependencies = [
"polars>=1.0", "duckdb>=1.0", "pydantic>=2.7",
"pydantic-settings>=2.3", "httpx>=0.27",
]
[dependency-groups]
dev = ["pytest>=8", "ruff>=0.5"]
[tool.ruff]
line-length = 100
[tool.ruff.lint]
select = ["E", "F", "I", "B", "UP"]
ruff check và ruff format chạy trước mỗi commit (qua pre-commit) và trong CI. Đây là lớp bảo vệ rẻ nhất: bắt lỗi cú pháp, import thừa, biến chưa dùng trước khi chúng kịp gây sự cố lúc 23h đóng sổ.
2. Cấu trúc project: tách I/O, validate, transform
Sai lầm phổ biến là nhét mọi thứ vào một script 800 dòng. Pipeline bền vững tách bạch tác vụ có tác dụng phụ (đọc/ghi, gọi mạng) khỏi logic thuần (biến đổi, tổng hợp) — vì logic thuần dễ test, còn I/O thì cần mock.
eod_pipeline/
├── pyproject.toml
├── src/eod/
│ ├── config.py # pydantic-settings
│ ├── models.py # Pydantic: Transaction, ...
│ ├── ingest.py # async httpx + đọc file (I/O)
│ ├── validate.py # tách hợp lệ / lỗi
│ ├── transform.py # Polars lazy (thuần)
│ ├── warehouse.py # DuckDB + ghi Parquet (I/O)
│ └── run.py # orchestrate
└── tests/
├── test_validate.py
└── test_transform.py
Cấu hình gom về một chỗ bằng pydantic-settings — đọc từ biến môi trường, có validate và giá trị mặc định, không rải os.getenv khắp nơi:
# config.py
from pydantic_settings import BaseSettings, SettingsConfigDict
class Settings(BaseSettings):
model_config = SettingsConfigDict(env_prefix="EOD_")
partner_api_base: str
partner_api_key: str
landing_dir: str = "/data/landing"
warehouse_dir: str = "/data/warehouse"
max_concurrency: int = 8
settings = Settings() # tự đọc EOD_PARTNER_API_KEY, ...
3. Ingest song song bằng async/httpx
Với vài chục nguồn, kéo tuần tự thì phần lớn thời gian là ngồi chờ mạng. Async và httpx cho phép mở nhiều request cùng lúc, giới hạn bằng Semaphore để không quá tải đối tác. File cục bộ (CSV/Parquet) thì đọc song song bằng thread pool vì đó là I/O đĩa.
# ingest.py
import asyncio, httpx
async def fetch_partner(client, sem, branch: str) -> list[dict]:
async with sem:
r = await client.get(f"/eod/{branch}", timeout=30)
r.raise_for_status()
return r.json()["records"]
async def ingest_partners(branches: list[str]) -> list[dict]:
sem = asyncio.Semaphore(settings.max_concurrency)
headers = {"Authorization": f"Bearer {settings.partner_api_key}"}
async with httpx.AsyncClient(base_url=settings.partner_api_base,
headers=headers) as client:
tasks = [fetch_partner(client, sem, b) for b in branches]
results = await asyncio.gather(*tasks, return_exceptions=True)
rows = []
for b, res in zip(branches, results):
if isinstance(res, Exception):
log.error("ingest thất bại chi nhánh %s: %s", b, res)
continue # ghi nhận, không làm sập cả mẻ
rows.extend(res)
return rows
Lưu ý return_exceptions=True: một chi nhánh lỗi không được kéo sập toàn bộ. Ta ghi log và tiếp tục — cuối mẻ báo cáo rõ nguồn nào thiếu để vận hành xử lý.
4. Validate bằng Pydantic v2 — tách dòng lỗi
Dữ liệu từ nhiều nguồn không bao giờ sạch. Pydantic v2 định nghĩa hợp đồng dữ liệu (data contract) một lần, rồi ép mọi bản ghi phải tuân thủ. Điểm mấu chốt nghiệp vụ: không được để bản ghi lỗi rơi âm thầm — phải đưa vào vùng cách ly (quarantine / dead-letter) để đối soát.
# models.py
from datetime import datetime
from decimal import Decimal
from pydantic import BaseModel, field_validator
class Transaction(BaseModel):
txn_id: str
account_no: str
amount: Decimal
currency: str
kind: str # debit / credit
booked_at: datetime
branch: str
@field_validator("currency")
@classmethod
def valid_ccy(cls, v: str) -> str:
if v not in {"VND", "USD", "EUR"}:
raise ValueError(f"currency không hợp lệ: {v}")
return v
# validate.py
def split_valid(rows: list[dict]) -> tuple[list[Transaction], list[dict]]:
good, bad = [], []
for r in rows:
try:
good.append(Transaction.model_validate(r))
except Exception as e:
bad.append({"raw": r, "error": str(e)})
return good, bad
Dùng Decimal cho số tiền (không dùng float) là quy tắc bắt buộc trong tài chính để tránh sai số dấu phẩy động. Danh sách bad được ghi ra Parquet riêng — vừa để đối soát, vừa là số liệu chất lượng dữ liệu theo dõi hằng ngày.
5. Transform và tổng hợp bằng Polars lazy
Sau khi có tập bản ghi sạch, dựng nó thành Polars DataFrame rồi dùng lazy API để mô tả toàn bộ phép biến đổi trước khi chạy. Lazy cho phép Polars tối ưu kế hoạch (đẩy bộ lọc xuống sớm, chỉ đọc cột cần) và xử lý khối lượng lớn hiệu quả, thậm chí streaming khi vượt RAM.
# transform.py
import polars as pl
def build_daily(txns: list[Transaction], run_date: str) -> pl.DataFrame:
df = pl.DataFrame([t.model_dump() for t in txns])
return (
df.lazy()
# khử trùng lặp theo txn_id — nền tảng của idempotency
.unique(subset=["txn_id"], keep="last")
.filter(pl.col("amount") > 0)
.with_columns(
pl.col("booked_at").dt.date().alias("book_date"),
)
.group_by(["branch", "currency", "kind"])
.agg(
pl.len().alias("txn_count"),
pl.col("amount").sum().alias("total_amount"),
)
.with_columns(pl.lit(run_date).alias("run_date"))
.collect() # đến đây mới thực thi
)
Bước .unique(subset=["txn_id"]) rất quan trọng: nếu một nguồn gửi trùng, hoặc pipeline chạy lại, ta vẫn ra một kết quả duy nhất theo txn_id.
6. Truy vấn kiểm tra bằng DuckDB trên Parquet
DuckDB đọc trực tiếp Parquet bằng SQL, không cần nạp vào server database. Rất hợp cho bước kiểm tra chất lượng và tạo báo cáo tổng: viết SQL quen thuộc, chạy trên file kết quả. DuckDB và Polars chia sẻ dữ liệu qua Arrow gần như không tốn chi phí copy.
# warehouse.py
import duckdb
def quality_check(df: pl.DataFrame) -> dict:
con = duckdb.connect()
con.register("daily", df.to_arrow())
row = con.execute("""
SELECT
SUM(txn_count) AS total_txn,
COUNT(*) FILTER (WHERE total_amount <= 0) AS bad_agg,
COUNT(DISTINCT branch) AS n_branch
FROM daily
""").fetchone()
return {"total_txn": row[0], "bad_agg": row[1], "n_branch": row[2]}
7. Ghi Parquet, nạp kho, và idempotency
Kết quả ghi ra Parquet phân vùng theo ngày. Chiến lược ghi đè theo phân vùng (overwrite partition) là cách đơn giản để đạt idempotency: chạy lại ngày 2026-07-13 sẽ thay thế toàn bộ phân vùng của ngày đó, không cộng dồn.
def write_partition(df: pl.DataFrame, run_date: str):
path = f"{settings.warehouse_dir}/eod/run_date={run_date}"
# ghi tạm rồi đổi tên (atomic) — tránh trạng thái nửa vời
tmp = path + ".tmp"
df.write_parquet(f"{tmp}/data.parquet")
os.replace(tmp, path)
Với kho hiện đại, phân vùng Parquet này có thể là bảng trong lakehouse dạng Iceberg — khi đó việc ghi lại một ngày tận dụng cơ chế snapshot/ACID của table format thay vì đổi tên thư mục thủ công. Nguyên tắc idempotency vẫn giữ nguyên: một khóa nghiệp vụ (ngày + txn_id) → một dòng kết quả.
Ghép orchestration lại:
# run.py
async def main(run_date: str, branches: list[str]):
raw = await ingest_partners(branches)
raw += read_local_files(settings.landing_dir) # CSV/Parquet cục bộ
good, bad = split_valid(raw)
write_quarantine(bad, run_date) # không bỏ sót lỗi
daily = build_daily(good, run_date)
stats = quality_check(daily)
if stats["bad_agg"] > 0:
raise RuntimeError(f"phát hiện {stats['bad_agg']} nhóm tổng bất thường")
write_partition(daily, run_date)
log.info("hoàn tất %s: %s", run_date, stats)
8. Test và CI làm gác cổng
Logic đã tách thuần nên test bằng pytest rất gọn. Ưu tiên test những chỗ dễ sai: validate biên, khử trùng lặp, tổng hợp.
# tests/test_transform.py
def test_dedup_va_tong():
txns = [
Transaction(txn_id="A", account_no="1", amount=Decimal("100"),
currency="VND", kind="credit",
booked_at=datetime(2026,7,13,9), branch="HN"),
Transaction(txn_id="A", account_no="1", amount=Decimal("100"),
currency="VND", kind="credit",
booked_at=datetime(2026,7,13,9), branch="HN"), # trùng
]
out = build_daily(txns, "2026-07-13")
assert out["txn_count"].sum() == 1 # trùng đã bị khử
assert out["total_amount"].sum() == Decimal("100")
CI (GitHub Actions/GitLab CI) chạy uv sync, ruff check, rồi pytest. Nếu bất kỳ bước nào đỏ, không cho merge. Đây là gác cổng: pipeline chạy production chỉ từ code đã qua kiểm thử.
9. Vận hành: lịch chạy, quan sát, chất lượng dữ liệu
Chạy theo lịch: đóng gói run.py thành một job, kích hoạt bằng cron cho môi trường đơn giản, hoặc Airflow khi cần phụ thuộc/retry/backfill. Airflow bổ sung theo dõi lần chạy, cảnh báo khi trễ, và chạy lại một ngày lịch sử (backfill) — hưởng lợi trực tiếp từ tính idempotent đã xây.
Quan sát (observability): log có cấu trúc mỗi bước (số bản ghi vào/ra, số dòng lỗi, thời gian). Xuất các con số này thành chỉ số chất lượng dữ liệu: tỉ lệ dòng lỗi, số chi nhánh thiếu, độ lệch tổng doanh số so với ngày trước. Khi tỉ lệ lỗi vượt ngưỡng, cảnh báo sớm thay vì để nghiệp vụ phát hiện qua báo cáo sai.
Tái chạy: nhờ ghi đè phân vùng và khử trùng lặp theo txn_id, chỉ cần chạy lại với cùng run_date là dữ liệu trở về đúng — không sợ nhân đôi.
Use case thực tế
Bối cảnh NCB, pipeline báo cáo cuối ngày chạy trên một server khỏe (16 vCPU, 64 GB RAM), không cần cụm phân tán. Các số liệu dưới đây là ước lượng minh họa cho quy mô tầm trung, không phải đo lường chính thức:
- Đầu vào: ~35 nguồn (25 file CSV/Parquet cục bộ + 10 feed API đối tác), tổng khoảng 40 triệu bản ghi giao dịch/ngày, dữ liệu thô ~6 GB.
- Ingest async: kéo 10 API song song với
Semaphore(8)— ước tính ~40–60 giây thay vì ~5–7 phút nếu tuần tự. Đọc file cục bộ song song thêm ~30 giây. - Validate Pydantic: bắt được ~0,3–0,8% bản ghi lỗi (thiếu trường, currency lạ, số tiền âm) đưa vào quarantine — ước tính cỡ 120.000–320.000 dòng, đối soát riêng mỗi sáng.
- Transform Polars lazy + DuckDB: khử trùng lặp, group_by theo chi nhánh/tiền tệ/kênh trên ~40 triệu dòng ước tính ~20–40 giây nhờ đa lõi.
- Tổng thời gian một mẻ: ước tính ~3–5 phút, thoải mái trong cửa sổ đóng sổ.
- Tái chạy: khi một feed đối tác về muộn, chạy lại
run_datehôm đó ghi đè đúng phân vùng, không nhân đôi số liệu.
Điểm mấu chốt nghiệp vụ: 0,3–0,8% dòng lỗi không bị làm rơi âm thầm mà được cách ly và đếm — đó là khác biệt giữa một script tạm bợ và một pipeline kiểm soát được chất lượng.
Ghi nhớ
- Pipeline production = ghép các mảnh của cả series: uv/ruff (môi trường + chất lượng) → async/httpx (ingest song song) → Pydantic v2 (validate, tách dòng lỗi) → Polars lazy + DuckDB (transform, tổng hợp, kiểm tra) → Parquet/kho → pytest/CI (gác cổng).
- Tách bạch I/O (đọc/ghi, mạng) khỏi logic thuần (transform) — logic thuần dễ test, I/O thì mock.
- Không để bản ghi lỗi rơi âm thầm: đưa vào quarantine/dead-letter và đếm như chỉ số chất lượng dữ liệu.
- Dùng Decimal cho số tiền, không dùng float.
- Idempotency là yêu cầu bắt buộc: khử trùng lặp theo khóa nghiệp vụ (
txn_id) + ghi đè theo phân vùng ⇒ chạy lại không nhân đôi. - Cấu hình gom về pydantic-settings, đọc từ biến môi trường; đừng rải
os.getenv. - CI chạy
ruff check+pytestlàm điều kiện merge — chỉ code đã kiểm thử mới lên production. - Vận hành cần quan sát được: log có cấu trúc, chỉ số tỉ lệ lỗi/độ lệch tổng, cảnh báo sớm; lịch chạy bằng cron cho đơn giản hoặc Airflow khi cần retry/backfill.
- Trên một máy khỏe, khối lượng vài chục triệu dòng/ngày hoàn toàn xử lý được bằng Python hiện đại — chưa cần cụm phân tán.
Nguồn tham khảo
- Polars — User Guide (đặc biệt mục Lazy API và Streaming)
- DuckDB Documentation — mục Reading & Writing Parquet
- Pydantic v2 Documentation và pydantic-settings
- Apache Parquet — Documentation
- HTTPX Documentation — mục Async Support
- pytest Documentation
- uv Documentation và Ruff Documentation
- Python
decimal— Decimal fixed-point and floating-point arithmetic
Bài viết liên quan
Vì sao Python là ngôn ngữ số một của data engineer: vai trò trong pipeline (ingest/transform/orchestrate), hệ sinh thái thư viện (pandas/polars/pyarrow/sqlalchemy), quản lý môi trường (venv/uv/poetry), và khi nào dùng Python vs SQL/Spark.
Học cách tổ chức code Python: định nghĩa hàm với tham số vị trí/từ khoá/mặc định, *args/**kwargs, lambda và hàm bậc cao, closure, decorator, generator với yield. Đóng gói code thành module và package, cô lập thư viện bằng môi trường ảo venv, quản lý phụ thuộc với pip và requirements.txt để dự án tái lập được trên mọi máy.
Biến script thành pipeline đáng tin cậy: cấu trúc project & packaging (uv/poetry), type hints & pydantic, kiểm thử với pytest, logging & cấu hình, đóng gói Docker, và tích hợp CI cho code dữ liệu.
Hướng dẫn OOP trong Python từ class/instance, kế thừa và super(), đa hình & duck typing, encapsulation tới dunder methods, @property, classmethod/staticmethod, dataclass và type hints (mypy). Kèm nguyên tắc clean code: đặt tên rõ nghĩa, hàm nhỏ, DRY, SOLID cùng chuẩn PEP8 với công cụ ruff/black.
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ẻ!