apache / apache/airflow

Access specific values within an `XCom` value using Taskflow API

Open
#16,618 13 comments 6 reactions 0 assignees View on GitHub
kind:feature
Dominant language
Python
Stars
46.9k
Forks
17.8k
Avg merge
2d 10h
Merged PRs (30d)
483

Description

**Description**
Currently the `output` property of operators doesn't support accessing a specific value within an `XCom` but rather the _entire_ `XCom` value. Ideally the behavior of calling the `XComArg` via the `output` property would function the same as the `task_instance.xcom_pull()` method in which a user has immediate access the `XCom` value and can directly access specific values in that `XCom`.

For example, in the [example DAG](https://github.com/apache/airflow/blob/main/airflow/providers/apache/beam/example_dags/example_beam.py) in the Apache Beam provider, the `jobId` arg in the `DataflowJobStatusSensor` task is a templated value using the `task_instance.xcom_pull()` method and is then accessing the `dataflow_job_id` key within the `XCom` value:
```python
start_python_job_dataflow_runner_async = BeamRunPythonPipelineOperator(
task_id="start_python_job_dataflow_runner_async",
runner="DataflowRunner",
py_file=GCS_PYTHON_DATAFLOW_ASYNC,
pipeline_options={
'tempLocation': GCS_TMP,
'stagingLocation': GCS_STAGING,
'output': GCS_OUTPUT,
},
py_options=[],
py_requirements=['apache-beam[gcp]==2.26.0'],
py_interpreter='python3',
py_system_site_packages=False,
dataflow_config=DataflowConfiguration(
job_name='{{task.task_id}}',
project_id=GCP_PROJECT_ID,
location="us-central1",
wait_until_finished=False,
),
)

wait_for_python_job_dataflow_runner_async_done = DataflowJobStatusSensor(
task_id="wait-for-python-job-async-done",
job_id="{{task_instance.xcom_pull('start_python_job_dataflow_runner_async')['dataflow_job_id']}}",
expected_statuses={DataflowJobStatus.JOB_STATE_DONE},
project_id=GCP_PROJECT_ID,
location='us-central1',
)
```
There is no current, equivalent way to directly access the `dataflow_job_id` value in same manner using the `output` property.

Using `start_python_job_dataflow_runner_async.output["dataflow_job_id"]` yields an equivalent `task_instance.xcom_pull(task_ids='start_python_job_dataflow_runner_async', key='dataflow_job_id'`.

Or even `start_python_job_dataflow_runner_async.output["return_value"]["dataflow_job_id"]` yields the same result: `task_instance.xcom_pull(task_ids='start_python_job_dataflow_runner_async', key='dataflow_job_id'`.

It seems the only way to get the desired behavior currently is to hack around the `__str__` method that's available with `XComArg`:
```python
start_python_job_dataflow_runner_async_output = str(start_python_job_dataflow_runner_async.output).strip("{ }")

wait_for_python_job_dataflow_runner_async_done = DataflowJobStatusSensor(
task_id="wait-for-python-job-async-done",
job_id="{{{{ {start_python_job_dataflow_runner_async_output}['dataflow_job_id'] }}}}",
expected_statuses={DataflowJobStatus.JOB_STATE_DONE},
project_id=GCP_PROJECT_ID,
location='us-central1',
)

```
This approach is not elegant, straightforward, nor user-friendly.

**Use case / motivation**
It's functionally intuitive for users to have direct access to the specific values in an `XCom` related to the `XComArg` via the Taskflow API like the classic `xcom_pull()` method. Ideally using an operator's `.output` property and the `xcom_pull()` method would behave the same way when needing to pass the actual values between operators.

**Are you willing to submit a PR?**
I would love to but I would certainly need some guidance on nuances here.

**Related Issues**
https://github.com/apache/airflow/issues/10285

Contributor guide

Open the contributing guide

Research direction

Start with the Taskflow API's XComArg and operator output behavior, then compare it with task_instance.xcom_pull() in the example DAG at airflow/providers/apache/beam/example_dags/example_beam.py. Determine how direct access to nested XCom values should behave, including the return_value case. Done means the Dataflow example can access dataflow_job_id through output without templating or __str__ workarounds.

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.