apache / apache/airflow

Add filtering of XCOMs for Dynamic Task Mapping

Open
#35,969 0 comments 0 reactions 0 assignees View on GitHub
area:dynamic-task-mapping kind:feature
Dominant language
Python
Stars
46.9k
Forks
17.8k
Avg merge
2d 7h
Merged PRs (30d)
484

Description

### Description

Allow filtering mappable XCOMs such that inputs to dynamic task mapping run on a subset of the XCOM.

```python
@task
def get_integers():
return [-1, 0, 1]

@task
def print_x(x):
print(x)

print_x.expand(get_integers().filter(lambda x: x>0))
# Prints:
# 1
```

### Use case/motivation

Currently, this can be implemented with a dedicated task just for filtering. This becomes undesired if you intend to filter to partition the XCOM output as inputs of multiple task (i.e., branching/partitioning):

### Dynamic task branching with XCOM filtering
```python
@task
def get_integers():
return [-1, 0, 1]

@task
def print_x(x):
print(x)

integers = get_integers()
task_print_positive = print_x.expand(integers.filter(lambda x: x>0))
task_print_negatives = print_x.expand(integers.filter(lambda x: x<0))
task_print_zeros = print_x.expand(integers.filter(lambda x: x==0))

# The flow logic is effectively:
# integers >> [task_print_positive, task_print_negatives, task_print_zeros]
```

### Dynamic task branching with Taskflow API
```python
@task
def get_integers():
return [-1, 0, 1]

@task
def print_x(x):
print(x)

@task
def filter_positive(x):
return list(filter(x>0, x))

@task
def filter_negative(x):
return list(filter(x<0, x))

@task
def filter_zero(x):
return list(filter(x==0, x))

integers = get_integers()
task_print_positive = print_x.expand(filter_positive(integers))
task_print_negatives = print_x.expand(filter_negatives(integers))
task_print_zeros = print_x.expand(filter_zeros(integers))

# The flow logic is effectively:
# integers >> [filter_positive >> task_print_positive, filter_negative >> task_print_negatives, filter_zeros >> task_print_zeros]
```

### Related issues

_No response_

### Are you willing to submit a PR?

- [ ] Yes I am willing to submit a PR!

### Code of Conduct

- [X] I agree to follow this project's [Code of Conduct](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md)

Contributor guide

Open the contributing guide

Research direction

Start with the Dynamic Task Mapping and Taskflow API examples in the issue, comparing the proposed filter usage with the current dedicated filtering-task workaround. Define the filtering and branching behavior for mappable XCOMs, then add coverage for positive, negative, and zero partitions; no implementation files or tests are named in the issue.

Written by the indexing model from the issue text.

Assessment

Tech stack
python
Domain
data-engineering
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.