动态任务映射

使用 Airflow 构建数据流水线

Volker Janz

Senior Developer Advocate at Astronomer

硬编码任务的问题

@task
def process_file_1():
    transform("/data/file_1.csv")
@task
def process_file_2():
    transform("/data/file_2.csv")
@task
def process_file_3():
    transform("/data/file_3.csv")

动态任务映射示例

使用 Airflow 构建数据流水线

.expand()

@task
def fetch_files():
    return ["/data/file_1.csv", "/data/file_2.csv", "/data/file_3.csv"]

@task(max_active_tis_per_dagrun=2) # control max parallelism def process_file(path: str): print(f"processing file: {path}")
files = fetch_files() process_file.expand(path=files)

运行时扩展

使用 Airflow 构建数据流水线

UI 中的映射任务

 

  • 映射实例作为单个任务下的索引实例出现
  • 点进每个实例可查看独立日志与状态
  • 即使有数百个实例,DAG 仍清晰易读

UI 中的映射任务

使用 Airflow 构建数据流水线

.partial()

@task
def process_file(path: str, output_dir: str):
    transform(path, output_dir)

files = get_files() process_file.partial(output_dir="/out").expand(path=files)

partial 与 expand

使用 Airflow 构建数据流水线

笛卡尔积问题

files = get_files()           # ["/data/a.csv", "/data/b.csv", "/data/c.csv"]
destinations = get_targets()  # ["s3://out/a", "s3://out/b", "s3://out/c"]

# 笛卡尔积:3 x 3 = 9 个实例(并非所需!) process.expand(path=files, dest=destinations)

拉链配对列表

使用 Airflow 构建数据流水线

.zip()

files = get_files()           # ["/data/a.csv", "/data/b.csv", "/data/c.csv"]
destinations = get_targets()  # ["s3://out/a", "s3://out/b", "s3://out/c"]

# 成对:3 个实例(a->a,b->b,c->c) process.expand_kwargs(files.zip(destinations))

 

  • zip() 按位置进行一一对应
  • expand_kwargs 将每对元素解包为关键字参数
使用 Airflow 构建数据流水线

何时用哪种方式

 

模式 适用场景
expand() 映射到单个列表
partial() 固定各实例的共享参数
zip() 按位置配对多个列表

 

  • 帮助构建可扩展的 DAGs
  • 数据变化时,任务数量会在运行时自动调整
使用 Airflow 构建数据流水线

让我们一起练习吧!

使用 Airflow 构建数据流水线

Preparing Video For Download...