Xây dựng Data Pipeline ‘bất tử’ với Prefect: Đừng để script Python chết trong im lặng

Python tutorial - IT technology blog
Python tutorial - IT technology blog

Tạm biệt nỗi lo script Python “ngỏm” lúc nửa đêm

Tưởng tượng bạn vừa hoàn thành script cào giá crypto và treo nó lên VPS bằng Cron job lúc 11 giờ đêm. Bạn yên tâm đi ngủ, hy vọng sáng ra sẽ có một file CSV đầy ắp dữ liệu. Nhưng thực tế phũ phàng: Một lỗi kết nối API xảy ra lúc 11:05, script dừng chạy, và bạn mất trắng 8 tiếng dữ liệu quý giá.

Trước đây, mình thường phải tốn hàng giờ để bới log file nặng vài trăm MB hoặc viết mỏi tay các khối try...except chỉ để bắt lỗi mạng. Mọi thứ thay đổi khi mình chuyển sang Prefect. Công cụ này giúp biến những đoạn code rời rạc thành một hệ thống pipeline bền bỉ, có khả năng tự phục hồi và giám sát trực quan.

Bắt đầu nhanh: Chạy Pipeline đầu tiên sau 2 phút

Prefect không bắt bạn phải học lại từ đầu. Nó hoạt động dựa trên chính code Python thuần túy của bạn. Đầu tiên, hãy cài đặt thư viện:

pip install -U prefect

Hãy xem cách mình biến một script lấy dữ liệu thời tiết bình thường thành một flow chuyên nghiệp bằng các decorator @task@flow:

from prefect import task, flow
import random

@task(retries=3, retry_delay_seconds=10)
def get_data():
    # Giả lập lỗi API với tỉ lệ 30%
    if random.random() > 0.7:
        raise ValueError("API không phản hồi!")
    return [25, 28, 30, 22]

@task
def transform_data(data):
    return [x * 1.8 + 32 for x in data] # Chuyển C sang F

@flow(name="Weather Pipeline")
def weather_flow():
    raw = get_data()
    processed = transform_data(raw)
    print(f"Nhiệt độ (F): {processed}")

if __name__ == "__main__":
    weather_flow()

Điểm “ăn tiền” nằm ở retries=3. Nếu API gặp sự cố, Prefect sẽ tự động thử lại sau 10 giây. Bạn không cần viết thêm bất kỳ dòng code logic retry phức tạp nào. Mọi tiến trình chạy đều được ghi log chi tiết ngay tại terminal.

Tại sao Prefect lại vượt trội hơn Cron job hay Airflow?

Nhiều bạn sẽ thắc mắc: “Dùng Cron cho nhẹ máy, hoặc Airflow cho đúng chuẩn enterprise chứ dùng Prefect làm gì?”. Từ kinh nghiệm thực chiến của mình, đây là 3 lý do cốt lõi:

1. Dashboard quan sát (Observability) cực xịn

Với Cron job, bạn hoàn toàn mù tịt về trạng thái script. Với Prefect, chỉ cần gõ prefect server start, bạn sẽ có ngay một Dashboard tại localhost:4200. Tại đây, bạn có thể xem biểu đồ chạy theo thời gian, task nào tốn nhiều tài nguyên nhất và nguyên nhân chính xác khi một task bị fail.

2. Xử lý lỗi thông minh với Exponential Backoff

Khi xử lý khoảng 500.000 bản ghi khách hàng, việc database bị quá tải là chuyện thường ngày. Prefect hỗ trợ exponential backoff, giúp thời gian chờ giữa các lần retry tăng dần. Điều này giúp hệ thống của bạn không “dội bom” yêu cầu vào server khi nó đang gặp sự cố.

3. Code-first: Viết code như đang viết Python thường

Airflow yêu cầu bạn cấu trúc code theo DAG khá gò bó và cồng kềnh. Prefect thì khác. Bạn chỉ cần giữ nguyên logic cũ và bọc thêm decorator. Nó tôn trọng cách bạn viết code, giúp việc chuyển đổi từ script cũ sang pipeline mới chỉ mất vài phút.

Nâng cao: Tiết kiệm tài nguyên với Caching

Giả sử bạn cần tải một file báo cáo 2GB từ Google Drive. Bạn chắc chắn không muốn tải lại file này nếu các bước xử lý sau đó bị lỗi. Prefect xử lý việc này bằng cache_key_fn:

from prefect.tasks import task_input_hash
from datetime import timedelta

@task(cache_key_fn=task_input_hash, cache_expiration=timedelta(hours=2))
def download_heavy_file(file_id):
    # Task này sẽ không chạy lại nếu file_id không đổi trong 2 giờ
    print("Đang tải file cực nặng...")
    return "Data content"

Triển khai (Deployment) lên Server trong một nốt nhạc

Để pipeline tự chạy hàng ngày lúc 8 giờ sáng, bạn không cần đụng vào crontab -e của Linux. Hãy dùng lệnh deploy:

prefect deploy weather_script.py:weather_flow -n "Daily-Check" --cron "0 8 * * *"

Lịch trình này có thể bật/tắt hoặc chỉnh sửa ngay trên giao diện web. Cực kỳ tiện lợi khi bạn cần thay đổi giờ chạy mà không muốn SSH vào server.

Kinh nghiệm thực tế từ các dự án lớn

Sau hơn một năm dùng Prefect quản lý hệ thống ETL, mình rút ra 4 bài học xương máu:

  • Chia nhỏ task tối đa: Đừng gộp lấy dữ liệu, xử lý và lưu DB vào một task. Nếu bước lưu DB lỗi, bạn sẽ phải chạy lại từ đầu cả quá trình lấy dữ liệu tốn kém.
  • Dùng Blocks để bảo mật: Thay vì để API Key trong file .env dễ lộ, hãy lưu chúng vào Prefect Blocks. Code của bạn sẽ sạch hơn và an toàn hơn.
  • Gắn Tag để quản lý: Khi có trên 20 pipeline, hãy dùng tag như production, crawling để lọc nhanh trên Dashboard.
  • Đừng lạm dụng Retry: Nếu lỗi do logic code (ví dụ chia cho 0), retry 100 lần cũng không giải quyết được gì. Chỉ dùng retry cho lỗi ngoại cảnh như Network hoặc Timeout.

Xây dựng Data Pipeline không khó, nhưng làm sao để nó chạy ổn định mới là thử thách thực sự. Prefect giúp bạn gánh vác phần hạ tầng phức tạp để bạn tập trung hoàn toàn vào logic xử lý dữ liệu. Nếu bạn đang mệt mỏi với việc quản lý script chạy ngầm, hãy thử chuyển sang Prefect ngay hôm nay.

Share: