# 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.

Bản gốc: https://fdetimes.net/vi/bach-khoa/thuc-hanh-pipeline-du-lieu-hang-dem-bang-airflow/

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.

```python
# 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.

```python
# 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.

```python
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.

**Điểm mấu chốt:** Một task tốt nhận ngày chạy làm đầu vào và ghi đè lên đúng một phân vùng của ngày đó, nên chạy lại chỉ cho ra đúng kết quả cũ.

**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.

```python
# 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?"

**Thử ngay tuần này:**

- Viết hàm raw_path(run_date) và một test gọi hàm load hai lần liên tiếp, rồi xác nhận số dòng không đổi
- Rà một file DAG có sẵn (của bạn hoặc của team) và gạch hết mọi lệnh truy cập DB hay gọi mạng đang nằm ở top-level
- Thêm vào CV một dòng mô tả pipeline của bạn theo công thức: lịch chạy, cách tách phần tăng dần (incremental), cách chạy lại khi lỗi

## Nguồn

- [Airflow Overview (Apache Airflow docs)](https://airflow.apache.org/docs/apache-airflow/stable/core-concepts/overview.html)

- [Cron and Time Intervals (Apache Airflow docs)](https://airflow.apache.org/docs/apache-airflow/stable/authoring-and-scheduling/cron.html)

- [Best Practices (Apache Airflow docs)](https://airflow.apache.org/docs/apache-airflow/stable/best-practices.html)

- [DAG writing best practices (Astronomer Docs)](https://www.astronomer.io/docs/learn/dag-best-practices)

- [What is data engineering? (IBM Think)](https://www.ibm.com/think/topics/data-engineering)
