Pipeline Architecture
"A pipeline is a set of data processing elements connected in series, where the output of one element is the input of the next. Pipelines are the backbone of modern data processing and CI/CD systems." — Michael T. Nygard, Release It! (2007)
Pipeline Architecture (hay Pipes & Filters) là một trong những kiến trúc phần mềm cổ điển nhất, bắt nguồn từ các hệ điều hành Unix (1970s) với khái niệm pipe (|). Trong kiến trúc này, dữ liệu được xử lý qua một chuỗi các bước (stages/filters) kết nối với nhau qua các kênh (pipes/channels). Mỗi bước nhận đầu vào, xử lý, và gửi kết quả sang bước tiếp theo.
Tổng quan
Lịch sử và nguồn gốc
Pipeline architecture có lịch sử lâu đời trong khoa học máy tính:
- 1970s: Unix pipes (
ls | grep | sort) — triết lý "do one thing and do it well" - 1978: Ken Thompson giới thiệu pipeline trong Unix — mỗi lệnh là một filter
- 1990s: Data Transformation Services (DTS), ETL pipelines trong data warehousing
- 2004: MapReduce (Google) — pipeline cho distributed data processing
- 2009: Apache Hadoop — open-source distributed pipeline
- 2014: Apache Spark — in-memory distributed pipeline
- 2015-nay: CI/CD pipelines (Jenkins, GitHub Actions, GitLab CI), ML pipelines (Kubeflow, TFX), stream processing (Kafka Streams, Flink)
Những người tiên phong
| Tên | Đóng góp |
|---|---|
| Ken Thompson & Dennis Ritchie | Unix pipes — origin of pipeline architecture |
| Doug McIlroy | Bell Labs — "pipes" concept và Unix philosophy |
| Jeffrey Dean & Sanjay Ghemawat | MapReduce — distributed pipeline |
| Matei Zaharia | Apache Spark — in-memory pipeline processing |
| Jay Kreps | Apache Kafka — event streaming pipeline |
| Brendan Gregg | Performance analysis pipeline |
Các loại Pipeline
| Loại | Mô tả | Ví dụ |
|---|---|---|
| Data Pipeline | Xử lý dữ liệu qua nhiều stages | ETL, ELT, data transformation |
| CI/CD Pipeline | Build → Test → Deploy | Jenkins, GitHub Actions |
| Processing Pipeline | Xử lý media/file pipeline | Image processing, video transcoding |
| Streaming Pipeline | Real-time data processing | Kafka Streams, Flink |
| ML Pipeline | ML model training pipeline | Kubeflow, TFX |
| Event Pipeline | Event processing pipeline | Event sourcing, CQRS |
| Log Pipeline | Log aggregation pipeline | ELK Stack (Elasticsearch, Logstash, Kibana) |
Bài toán
Hệ thống Xử lý Đơn hàng Thương mại Điện tử Quy mô Lớn
Giả sử bạn đang xây dựng TikiNgon — một nền tảng giao đồ ăn và hàng tạp hóa với hàng triệu đơn hàng mỗi ngày tại Việt Nam. Mỗi đơn hàng trải qua nhiều bước xử lý:
- Đặt hàng: Customer tạo đơn
- Xác thực: Kiểm tra thông tin, số dư ví
- Thanh toán: Xử lý payment (Visa, MoMo, COD)
- Kiểm tra kho: Check tồn kho, reserve item
- Xác nhận người bán: Gửi cho merchant, chờ xác nhận
- Đóng gói: Gán shipper, in hóa đơn
- Vận chuyển: Logistic tracking
- Hoàn thành: Giao hàng thành công
- Đánh giá: Gửi review request
Khó khăn với kiến trúc monolithic/traditional
Vấn đề 1 — Mỗi bước xử lý có yêu cầu khác nhau về tài nguyên:
Monolith: [Đặt hàng → Thanh toán → Kho → Gửi → Vận chuyển]
⬆ Tất cả trong 1 server
- Bước validate: Cần CPU (string parsing, regex)
- Bước payment: Cần I/O network (gọi API ngân hàng)
- Bước kho: Cần memory (cache inventory)
- Bước vận chuyển: Cần I/O disk (logistics routes)
→ Tài nguyên hỗn độn, không thể tối ưu từng bước
Vấn đề 2 — Một bước chậm kéo theo toàn bộ pipeline:
Nếu bước xác thực thanh toán mất 10 giây (do timeout ngân hàng), tất cả đơn hàng phía sau bị block. Hàng nghìn đơn hàng bị delay chỉ vì một bước.
# Synchronous pipeline — blocking
def process_order(order):
validate(order) # Nhanh: 100ms
process_payment(order) # Chậm: 10s (API ngân hàng timeout)
check_inventory(order) # Bị block vì chờ payment xong!
confirm_merchant(order) # Bị block tiếp!
assign_shipper(order) # Không thể chạy!
Vấn đề 3 — Khó thêm/xóa bước xử lý:
Khi business yêu cầu thêm bước "AI fraud detection" giữa payment và inventory, bạn phải:
- Sửa code ở vị trí chính xác
- Deploy lại toàn bộ ứng dụng
- Nguy cơ ảnh hưởng đến các bước khác
# Thêm bước mới phải sửa code hiện tại
def process_order(order):
validate(order)
process_payment(order)
# Phải chèn ở đây:
fraud_check(order) # New step — sửa function!
check_inventory(order)
# ...
Vấn đề 4 — Không thể retry từng bước riêng:
Khi bước payment fail, bạn phải chạy lại toàn bộ pipeline từ đầu:
- Đã validate → lại validate (lãng phí)
- Đã check kho → lại check (nguy cơ duplicate)
- Không thể resume từ bước fail
Vấn đề 5 — Khó scale từng bước:
Bước validate cần 2 servers, nhưng bước vận chuyển cần 50 servers. Trong monolith, bạn phải scale toàn bộ app.
Pipeline Architecture giải quyết vấn đề
- Decoupled stages: Mỗi stage độc lập, giao tiếp qua message queue
- Independent scaling: Mỗi stage scale riêng (validate: 2 pods, shipping: 50 pods)
- Fault isolation: Một stage fail không ảnh hưởng stage khác
- Retry per stage: Retry stage fail, không cần restart pipeline
- Dynamic pipeline: Thêm/xóa stage không ảnh hưởng code hiện tại
- Resource optimization: Mỗi stage dùng resource phù hợp
- Monitoring per stage: Biết chính xác stage nào chậm
Nguyên lý thiết kế
1. Single Responsibility per Stage
Mỗi stage làm đúng MỘT việc:
# GOOD: mỗi stage một responsibility
class ValidateOrder: ...
class ProcessPayment: ...
class CheckInventory: ...
# BAD: stage làm nhiều việc
class ValidateAndPaymentAndInventory: ...
2. Standardized Interface (Pipe Interface)
Tất cả stages đều có cùng interface:
@runtime_checkable
class Stage(Protocol):
async def process(self, context: PipelineContext) -> PipelineContext: ...
Input/output qua PipelineContext — một dict chứa tất cả dữ liệu pipeline.
3. Immutability
Context không nên bị mutate trực tiếp. Mỗi stage nên trả về context mới (hoặc copy-on-write).
4. Idempotency
Mỗi stage có thể chạy lại nhiều lần mà không gây side effects:
# GOOD: idempotent
def process_payment(context):
if context.get("payment_processed"):
return context # Already done, skip
# Process payment
return context | {"payment_processed": True}
# BAD: not idempotent
def process_payment(context):
charge_credit_card(context["amount"]) # Sẽ charge 2 lần nếu retry!
5. Error Handling per Stage
Mỗi stage tự xử lý lỗi và quyết định:
- Retry: Tạm thời (network timeout)
- Skip: Có thể bỏ qua (optional step)
- Abort: Dừng pipeline (critical error)
6. Pipeline Configuration
Pipeline được cấu hình động, không hardcode:
pipeline = Pipeline([
ValidateStage(),
PaymentStage(config.PAYMENT_GATEWAY),
InventoryStage(config.WAREHOUSE_API),
...
])
7. Observability
Mỗi stage phải emit metrics:
- Duration
- Success/failure count
- Input/output size
- Retry count
8. Backpressure
Khi stage sau chậm hơn stage trước, cần cơ chế backpressure để không làm overflow queue.
Cấu trúc chi tiết
Thành phần cốt lõi
┌─────────────────────────────────────────────────────────────────────────┐
│ PIPELINE ARCHITECTURE │
│ │
│ INPUT ──→ [Pipe 1] ──→ [Stage 1] ──→ [Pipe 2] ──→ [Stage 2] ──→ ... │
│ │ │
│ ▼ │
│ PipelineContext │
│ { │
│ "order_id": "123", │
│ "amount": 500000, │
│ "status": "pending", │
│ "errors": [], │
│ "results": {} │
│ } │
│ │
│ ┌─────────────────────────────────────────────────────────────────┐ │
│ │ PIPELINE MANAGER │ │
│ │ │ │
│ │ ┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │ │
│ │ │ Stage Runner │ │ Retry Logic │ │ Circuit │ │ │
│ │ │ │ │ (exponential │ │ Breaker │ │ │