FDE PulseViệc làm FDE đang mở 316Mới đăng 7 ngày qua 10Chủ đề nổi bật: Đào tạo kỹ năng FDE tại Đông Nam Á

Tờ báo của nghề Forward Deployed Engineer

Bách khoa

Pipeline Airflow hằng đêm cho trợ lý AI: chạy lại vẫn không hỏng

Chỉ cần một task bị retry lúc 2 giờ sáng là dữ liệu nạp đêm qua có thể bị nhân đôi. Nếu bạn là FDE, sáng hôm sau người phải giải thích chuyện đó với khách hàng có thể chính là bạn.

Đồ hoạPipeline đồng bộ hằng đêm cho trợ lý AI
  1. 1Scheduler tạo Dag RunTham số schedule "0 2 * * *" kích hoạt DAG lúc 2 giờ sáng, mỗi lần chạy là một Dag Run
  2. 2Extract tăng dầnChỉ kéo dữ liệu trong cửa sổ của ngày chạy, không tính theo giờ hiện tại
  3. 3Load thô lên remote storageGhi đè phân vùng dt=ngày trên kho kiểu S3, không ghi ra đĩa cục bộ
  4. 4TransformChuẩn hóa dữ liệu thô sau khi đã nạp, theo kiểu ELT
  5. 5Làm mới nguồn cho trợ lýThay thế dữ liệu của đúng ngày đó trong nguồn tri thức mà trợ lý AI đọc

Mỗi bước nhận ngày của Dag Run làm đầu vào và ghi đè đúng phân vùng ngày đó, nên chạy lại vẫn an toàn.

Đồ hoạ: FDE Times

Tóm tắt nhanh

  • Mỗi task phải cho cùng kết quả khi chạy lại, vì Airflow có thể retry task bị lỗi.
  • Chỉ kéo dữ liệu của đúng một ngày và ghi đè lên đúng phân vùng của ngày đó, đừng bao giờ append.
  • Để file DAG thật nhẹ, đẩy dữ liệu lên remote storage và đặt lịch bằng tham số schedule.
Chia sẻLinkedInFacebookX

Tài liệu best practices của Airflow có một câu đáng in ra dán cạnh màn hình: Airflow có thể retry task khi task lỗi, nên task phải cho cùng một kết quả dù chạy lại bao nhiêu lần. Câu này nghe hiển nhiên, nhưng nó quyết định cách thiết kế toàn bộ pipeline bên dưới.

Thử hình dung thế này: khách hàng có một trợ lý AI trả lời câu hỏi của nhân viên hỗ trợ, dựa trên ticket và tài liệu sản phẩm.

Dữ liệu cần được cập nhật mỗi đêm. Giả sử task nạp dữ liệu bị retry lúc 2 giờ sáng và nhân đôi 1.000 ticket thành 2.000. Sáng hôm sau trợ lý sẽ trích dẫn trùng lặp, đếm sai, và người dùng mất niềm tin vào nó.

Bài này hướng dẫn bạn dựng đúng pipeline đó, theo hướng chạy lại bao nhiêu lần cũng an toàn. Trợ lý AI chỉ trả lời dựa trên dữ liệu nó được nạp, nên nếu bạn muốn làm FDE, đây là phần việc đáng luyện kỹ trước khi ra site khách hàng.

Bạn sẽ dựng cái gì?

Airflow là nền tảng để xây dựng và chạy workflow. Mỗi workflow được biểu diễn bằng một DAG (Directed Acyclic Graph) gồm các Task. DAG quy định quan hệ phụ thuộc giữa các task, từ đó quyết định thứ tự chạy.

DAG trong bài có bốn task nối tiếp nhau. extract kéo ticket của đúng một ngày, load đẩy dữ liệu thô lên remote storage, transform chuẩn hóa dữ liệu, còn refresh_index làm mới nguồn tri thức mà trợ lý đọc. Cách nạp thô trước rồi mới biến đổi chính là ELT: dữ liệu thô vào kho tập trung trước, sau đó mới được chuẩn hóa.

Để làm theo, bạn cần một môi trường Airflow chạy được trên máy, Python cơ bản, và một nơi lưu trữ từ xa kiểu S3. Code bên dưới là bản giản lược để minh họa thiết kế. Đường dẫn import, tên operator và các tham số ngoài schedule thay đổi theo phiên bản, nên bạn cần đối chiếu với tài liệu của bản Airflow đang cài.

Bước 1: Vì sao file DAG phải thật nhẹ?

Scheduler của Airflow vừa kích hoạt workflow theo lịch, vừa giao task cho executor chạy, và để làm việc đó nó phải đọc (parse) các file DAG của bạn. Tài liệu Airflow nói rõ: ở top-level của file DAG không được truy cập database, không tính toán nặng và không gọi mạng.

# SAI: câu query chạy mỗi lần scheduler parse file, không phải lúc task chạy
rows = source_db.query("SELECT * FROM tickets")

# ĐÚNG: mọi I/O nằm trong hàm, chỉ chạy khi task được thực thi
def extract(run_date: str):
    ...

Kiểm tra: mở file DAG, đọc từ trên xuống. Ngoài import, khai báo hằng số và định nghĩa hàm, không được có dòng nào chạm vào DB hay API.

Bước 2: Chỉ kéo dữ liệu của một ngày

Tài liệu của Astronomer khuyên chia pipeline thành các lần extract và load tăng dần (incremental) ở mọi chỗ có thể. Nhờ vậy, khi có lỗi, bạn chỉ phải xử lý lại phần dữ liệu bị ảnh hưởng. Với đồng bộ hằng đêm, đơn vị tự nhiên là một ngày.

# sync_utils.py: logic thuần Python, test được mà không cần Airflow
def extract_window(run_date: str) -> tuple[str, str]:
    return f"{run_date} 00:00:00", f"{run_date} 23:59:59"

def raw_path(run_date: str) -> str:
    # bucket "acme-raw" chỉ là ví dụ
    return f"s3://acme-raw/tickets/dt={run_date}/part.jsonl"

Mấu chốt là cửa sổ dữ liệu được tính từ ngày của lần chạy, không phải từ “bây giờ”. Hôm nay bạn chạy lại job của ngày 07/10 thì nó vẫn kéo đúng ticket của ngày 07/10, không lẫn sang dữ liệu ngày 08/10.

Kiểm tra: gọi extract_window("2026-10-07") hai lần ở hai thời điểm khác nhau, kết quả phải giống hệt nhau.

Bước 3: Ghi đè, đừng append

Quay lại con số 1.000 thành 2.000 ở đầu bài. Nếu load ghi thêm (append) vào một bảng chung, mỗi lần retry lại cộng thêm một bản sao. Còn nếu load ghi đè lên đúng phân vùng dt=2026-10-07, thì chạy một lần hay ba lần, phân vùng đó vẫn chỉ có 1.000 dòng.

def load(run_date: str, rows: list[dict]):
    # storage là client giả định cho remote storage (S3 hoặc tương đương)
    storage.put(raw_path(run_date), rows, overwrite=True)

Dữ liệu phải lên remote storage thay vì nằm trên đĩa cục bộ, vì trong Airflow mỗi task có thể chạy trên một máy khác nhau. File mà extract ghi ra /tmp trên máy A sẽ không tồn tại khi transform chạy trên máy B. Tài liệu Airflow khuyên không lưu file hay config ở filesystem cục bộ, còn dữ liệu lớn thì nên đặt ở nơi như S3.

Kiểm tra: viết một test gọi load hai lần với cùng run_date, sau đó đếm số dòng. Nếu số dòng tăng gấp đôi thì task chưa idempotent.

Bước 4: Transform và làm mới nguồn cho trợ lý

transform đọc từ raw_path(run_date), làm sạch dữ liệu (bỏ trường rỗng, chuẩn hóa định dạng), rồi ghi kết quả sang một phân vùng “sạch”, cũng chia theo ngày. Sau đó refresh_index cập nhật nguồn tri thức mà trợ lý AI truy vấn theo cùng nguyên tắc: thay thế phần dữ liệu của ngày đó chứ không chèn thêm bản trùng.

Theo Astronomer, DAG và task idempotent giúp rút ngắn thời gian phục hồi sau sự cố và tránh mất dữ liệu. Ở khách hàng, lợi ích này rất cụ thể. Khi bên nguồn dữ liệu báo vừa sửa lỗi cho ngày 05/10, bạn chỉ cần chạy lại đúng ngày đó, không phải xóa sạch rồi đồng bộ lại từ đầu.

Bước 5: Đặt lịch chạy lúc 2 giờ sáng

Lịch chạy được đặt bằng tham số schedule của DAG, có thể là biểu thức cron, một timedelta hoặc một preset. Mỗi lần đến lịch, Airflow tạo một Dag Run tương ứng và lưu vào database backend.

# nightly_sync.py: khung giản lược, đối chiếu import với phiên bản của bạn
with DAG(
    dag_id="nightly_assistant_sync",
    schedule="0 2 * * *",   # cron: 2 giờ sáng mỗi ngày
    start_date=datetime(2026, 10, 1),
) as dag:
    extract_task >> load_task >> transform_task >> refresh_index_task

Các hàm ở bước 2 và 3 nên nhận run_date từ Dag Run đang chạy, đừng lấy từ đồng hồ hệ thống. Cách đọc ngày của Dag Run bên trong task khác nhau giữa các phiên bản Airflow, nên bạn cần tra tài liệu của bản đang cài.

Toàn bộ tính idempotent dựa vào lựa chọn này: chạy lại Dag Run nào thì xử lý đúng ngày của Dag Run đó.

Kiểm tra: sau đêm đầu tiên, mở giao diện Airflow và xác nhận đã có một Dag Run mới, cả bốn task đều thành công. Sau đó kích hoạt chạy lại thủ công đúng Dag Run đó và xác nhận số dòng không đổi.

Ba lỗi bạn sẽ gặp ở khách hàng

Lỗi phổ biến nhất là dùng “giờ hiện tại” để tính cửa sổ dữ liệu, kiểu WHERE updated_at > now() - 1 day. Chạy đúng giờ thì không sao, nhưng nếu retry lúc 9 giờ sáng thì cửa sổ đã trượt đi, và bạn sẽ mất hoặc trùng dữ liệu mà không có thông báo lỗi nào.

Lỗi thứ hai là đặt logic nặng ở top-level, thường do ai đó muốn “lấy danh sách bảng từ DB rồi sinh task động”. Lỗi thứ ba là ghi file trung gian ra đĩa cục bộ: chạy trên laptop thì ổn, nhưng lên môi trường nhiều worker của khách hàng là hỏng ngay.

Từ laptop đến site khách hàng

Pipeline trong bài là một lát cắt nhỏ của data engineering, mảng mà IBM định nghĩa là thiết kế và xây dựng hệ thống thu thập, lưu trữ và phân tích dữ liệu ở quy mô lớn.

Khi một mô tả công việc FDE nhắc đến “data pipeline”, “Airflow” hay “integrate with customer systems”, thực chất họ đang hỏi bạn có làm được lát cắt ấy trên hệ thống thật của khách hàng hay không.

Vì thế, trong CV đừng chỉ ghi “dùng Airflow”. Hãy viết rõ bạn đã thiết kế pipeline incremental theo ngày, có thể chạy lại một ngày bất kỳ mà không làm trùng dữ liệu. Câu đó cho thấy bạn hiểu điều khách hàng thật sự lo: một đêm lỗi là sáng hôm sau người dùng thấy ngay trợ lý trả lời sai.

Ngày đầu tiên tại khách hàng, đừng vội hỏi “dùng model nào”. Hãy hỏi trước: “Nếu job đêm nay lỗi, sáng mai mình chạy lại thế nào?”

5 nguồn
Đọc tiếp trên lộ trình · Chặng 2: Kỹ thuật rộngLần đầu vào database của khách hàng: đọc schema, truy vấn trên bản sao và khoá phiên chỉ đọcTrước khi gõ truy vấn đầu tiên ở khách hàng, hãy tìm một bản sao để đọc; nếu buộc phải dùng primary, hãy khoá phiên ở chế độ chỉ đọc và đặt thời hạn cho mọi câu lệnh.