การแชร์ข้อมูลระหว่าง Task

การสร้าง Data Pipeline ด้วย Airflow

Volker Janz

Senior Developer Advocate at Astronomer

XCom คืออะไร?

XCom push pull flow

 

  • XCom = การสื่อสารข้ามงาน (cross-communication)
  • ค่าต่าง ๆ จัดเก็บไว้ใน metadata database
  • แนวทางดั้งเดิม: ใช้ xcom_push() และ xcom_pull() โดยตรง
การสร้าง Data Pipeline ด้วย Airflow

XComArgs: แนวทาง TaskFlow

 

  • ค่าที่ return จาก @task จะกลายเป็น XComArg
  • XComArg คือ การอ้างอิงแบบ lazy ไม่ใช่ข้อมูลจริง
  • การส่งไปยังฟังก์ชันอื่นจะสร้าง dependency ขึ้น
  • เมื่อ runtime Airflow จะ resolve โดยดึงค่าจาก XCom

XComArgs ใน TaskFlow

การสร้าง Data Pipeline ด้วย Airflow

XComArgs ในโค้ด

@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
การสร้าง Data Pipeline ด้วย Airflow

ผสมแนวทางดั้งเดิมกับ 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)
การสร้าง Data Pipeline ด้วย Airflow

ข้อจำกัดของ XCom

ขีดจำกัดขนาด XCom ตาม backend

 

  • ขนาดสูงสุดขึ้นอยู่กับ database backend ที่ใช้
  • ต้องสามารถ JSON serialize ได้ (dict, list, string, number)
  • ทุกค่าที่เก็บจะเพิ่ม ภาระ ให้กับฐานข้อมูล
การสร้าง Data Pipeline ด้วย Airflow

Custom XCom backends

การกำหนดเส้นทาง XCom backend แบบกำหนดเอง

  • กำหนดเส้นทางค่าไปยัง object storage
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: จัดเก็บ XCom ใน object storage เมื่อถึง threshold ที่กำหนด
การสร้าง Data Pipeline ด้วย Airflow

มาฝึกกันเถอะ!

การสร้าง Data Pipeline ด้วย Airflow

Preparing Video For Download...