Chia sẻ dữ liệu giữa các task

Xây dựng Data Pipeline với Airflow

Volker Janz

Senior Developer Advocate at Astronomer

XCom là gì?

Luồng push/pull XCom

 

  • XCom = trao đổi chéo giữa các task
  • Giá trị được lưu trong cơ sở dữ liệu metadata
  • Cách cổ điển: dùng rõ xcom_push()xcom_pull()
Xây dựng Data Pipeline với Airflow

XComArg: cách của TaskFlow

 

  • Giá trị trả về của @task trở thành một XComArg
  • XComArg là một tham chiếu lười, không phải dữ liệu thật
  • Truyền nó vào hàm khác sẽ tạo phụ thuộc
  • Khi chạy, Airflow sẽ resolve bằng cách pull từ XCom

XComArg trong TaskFlow

Xây dựng Data Pipeline với Airflow

XComArg trong mã

@task
def get_config():
    return {"name": "daily_etl"}

@task
def run_pipeline(config):
    print(config["name"])


config = get_config() # XComArg
run_pipeline(config) # dependency + data
Xây dựng Data Pipeline với Airflow

Kết hợp kiểu cổ điển và TaskFlow

get_timestamp = BashOperator(
    task_id="get_timestamp",
    bash_command="date +%Y-%m-%d")

@task
def log_timestamp(ts_value):
    print(f"Timestamp: {ts_value}")


log_timestamp(get_timestamp.output)
Xây dựng Data Pipeline với Airflow

Giới hạn của XCom

Giới hạn kích thước XCom theo backend

 

  • Kích thước tối đa phụ thuộc vào database backend của bạn
  • Phải tuần tự hóa JSON được (dict, list, string, number)
  • Mỗi giá trị đều tăng tải lên cơ sở dữ liệu
Xây dựng Data Pipeline với Airflow

Backend XCom tùy chỉnh

Định tuyến backend XCom tùy chỉnh

  • Điều hướng giá trị sang object storage
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: chỉ lưu XCom vào object storage khi đạt ngưỡng
Xây dựng Data Pipeline với Airflow

Cùng luyện tập nào!

Xây dựng Data Pipeline với Airflow

Preparing Video For Download...