I have two tasks inside a TaskGroup that need to pull xcom values to supply the job_flow_id and step_id. Here's the code:
with TaskGroup('execute_my_steps') as execute_my_steps:
config = {some dictionary}
dependencies = {another dictionary}
task_id = 'execute_spark_job_step'
task_name = 'spark_job'
add_step = EmrAddStepsOperator(
task_id=task_id,
job_flow_id="{{ task_instance.xcom_pull(dag_id='my_dag', task_ids='emr', key='return_value') }}",
steps=create_emr_step(args=config, d=dependencies),
aws_conn_id='aws_default',
retries=3,
dag=dag
)
wait_for_step = EmrStepSensor(
task_id='wait_for_' + task_name + '_step',
job_flow_id="{{ task_instance.xcom_pull(dag_id='my_dag', task_ids='emr', key='return_value') }}",
step_id="{{ task_instance.xcom_pull(dag_id='my_dag', task_ids='" + task_id + "', key='return_value') }}",
retries=3,
dag=dag,
mode='reschedule'
)
add_step >> wait_for_step
The problem is the step_id does not render correctly. The wait_for_step value in the UI rendered template shows as 'None', however, the xcom return_value for execute_spark_job_step is there (this is the emr step_id).
wait_for_step rendered template:

When I remove the TaskGroup, it renders fine and the step waits until the job enters the completed state.
I need this to be in a task group because I will be looping through a larger config file and creating multiple steps.
Why doesn't this work? Do I need a nested TaskGroup? I tried using a TaskGroup without the context manager and still no luck.
