在任务间共享数据

使用 Airflow 构建数据流水线

Volker Janz

Senior Developer Advocate at Astronomer

什么是 XCom?

XCom push pull flow

 

  • XCom = 任务间跨通信
  • 值存储在元数据数据库
  • 经典用法:显式调用 xcom_push()xcom_pull()
使用 Airflow 构建数据流水线

XComArg:TaskFlow 方式

 

  • @task 的返回值会变成一个 XComArg
  • XComArg 是惰性引用,不是真实数据
  • 传给其他函数即可建立依赖
  • 运行时,Airflow 通过从 XCom 拉取来解析

XComArgs in TaskFlow

使用 Airflow 构建数据流水线

代码中的 XComArg

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

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


config = get_config() # XComArg
run_pipeline(config) # 依赖 + 数据
使用 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)
使用 Airflow 构建数据流水线

XCom 的限制

XCom size limits by backend

 

  • 最大值取决于您的数据库后端
  • 必须可JSON 序列化(dict、list、string、number)
  • 每个值都会给数据库增加负载
使用 Airflow 构建数据流水线

自定义 XCom 后端

Custom XCom backend routing

  • 将值路由到对象存储
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD:仅当达到阈值时才把 XCom 存到对象存储
使用 Airflow 构建数据流水线

让我们一起练习吧!

使用 Airflow 构建数据流水线

Preparing Video For Download...