Skip to content

Consume every queued asset event from concurrent mapped outlets - #71997

Open
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events
Open

Consume every queued asset event from concurrent mapped outlets#71997
Vamsi-klu wants to merge 2 commits into
apache:mainfrom
Vamsi-klu:pr1/54659-mapped-outlet-events

Conversation

@Vamsi-klu

Copy link
Copy Markdown
Contributor

Consume every queued asset event from concurrent mapped outlets

closes: #54659

What I did

Tests only. Mapped producer @task(outlets=[asset]).expand(...) succeeds, then one scheduler tick. Assert all N events land on the consumer run and in triggering_asset_events. No production change. No newsfragment.

Why I did

Three mapped outlets were finishing together and the consumer only saw a subset in triggering_asset_events. Consume-by-event-id is already on main. Existing tests still insert ADRQ by hand, so they never covered this emit path.

How I did

Emit: dag_maker.run_ti(..., map_index=N)register_asset_changes_in_db.
Consume: SchedulerJobRunner._create_dagruns_for_dags.
Context: get_template_context after loading consumed_asset_events.

One tick batches all visible events. That is the default, not a bug. Leftovers stay for the next tick.

What's the impact

None at runtime. If the mapped emit path regresses, these tests fail instead of silently dropping events.

What's the testing

airflow-core/tests/unit/jobs/test_scheduler_job.py

  • test_mapped_outlet_asset_events_consumed_in_one_tick
  • test_mapped_outlet_asset_events_consumed_across_staggered_ticks
  • test_mapped_outlet_asset_events_same_timestamp_are_all_consumed
  • test_mapped_outlet_asset_events_and_condition_waits_for_all_assets
  • test_mapped_outlet_asset_alias_events_are_all_consumed
uv run --project airflow-core pytest airflow-core/tests/unit/jobs/test_scheduler_job.py \
  -k test_mapped_outlet -q

Was generative AI tooling used to co-author this PR?
  • Yes (Grok 4.6)

Generated-by: Grok 4.6 following the guidelines

apache#70972 already consumes queued asset events by id, so concurrent mapped
outlets no longer strand events behind a timestamp watermark. These
tests lock that contract for mapped producers so a later consume-path
change cannot silently drop events from triggering_asset_events.

closes: apache#54659
The same-timestamp and AND cases only pinned consume if events were
already queued by hand. Running mapped TIs through register/queue
keeps those pins on the user path, including alias context lookup.

closes: apache#54659
@boring-cyborg boring-cyborg Bot added the area:Scheduler including HA (high availability) scheduler label Aug 23, 2026
@Vamsi-klu

Copy link
Copy Markdown
Contributor Author

cc @uranusjr @ashb @XD-DENG mapped outlet consume-all coverage on the scheduler path.


Drafted-by: Cursor Grok 4.6; reviewed by @Vamsi-klu before posting

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:Scheduler including HA (high availability) scheduler

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Asset-triggered DAGs miss events from concurrently completing dynamic mapped tasks

1 participant