Dynamisk taskmappning

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Problemet med hårdkodade tasks

@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")

Exempel på dynamisk taskmappning

Bygg datapipelines med 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)

Skalas vid körning

Bygg datapipelines med Airflow

Mappade tasks i gränssnittet

 

  • Mappade instanser visas som indexerade instanser under en enda task
  • Klicka på var och en för individuella loggar och status
  • DAG:en förblir läsbar även med hundratals instanser

Mappade tasks i gränssnittet

Bygg datapipelines med 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 och expand

Bygg datapipelines med Airflow

Korsproduktproblemet

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

# Cross product: 3 x 3 = 9 instances (not what we want!) process.expand(path=files, dest=destinations)

Zippa listor

Bygg datapipelines med 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"]

# Paired: 3 instances (a->a, b->b, c->c) process.expand_kwargs(files.zip(destinations))

 

  • zip() matchar element ett-till-ett efter position
  • expand_kwargs packar upp varje par som nyckelord­argument
Bygg datapipelines med Airflow

När ska du använda vad

 

Mönster Användningsfall
expand() Iterera över en enskild lista
partial() Fixa delade parametrar för alla instanser
zip() Para ihop flera listor efter position

 

  • Hjälper dig att bygga skalbara DAG:ar
  • När data förändras justeras antalet tasks automatiskt vid körning
Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...