caoterry commented on PR #73983: URL: https://github.com/apache/airflow/pull/73983#issuecomment-5934758697
Sure. Setup: Postgres 16, a fresh database migrated to `90e4d18ccadf` (the revision just before this PR), seeded with 100,000 `asset_partition_dag_run` rows for one partitioned Dag (99,500 already fired, each pointing at its own `dag_run`; the newest 500 pending), then `airflow db migrate` to this PR's head and `ANALYZE`. The SQL is compiled from the same ORM expressions as main (linked below); `EXPLAIN (ANALYZE, BUFFERS)`, warm cache. | Query | Before | After | |---|---|---| | Lookup in `AssetManager._get_or_create_apdr`, once per emitted partition key ([manager.py#L743-L751](https://github.com/apache/airflow/blob/8f3e8466c67c0d4bf5db4bdbb99fdce50f436e4f/airflow-core/src/airflow/assets/manager.py#L743-L751)) | Seq Scan + Sort, 99,999 rows filtered, 935 buffers, 6.96 ms | Index Scan Backward on `idx_apdr_target_dag_id_partition_key_id`, 4 buffers, 0.008 ms | | Pending scan in `_create_dagruns_for_partitioned_asset_dags`, every scheduler loop ([scheduler_job_runner.py#L2378-L2395](https://github.com/apache/airflow/blob/8f3e8466c67c0d4bf5db4bdbb99fdce50f436e4f/airflow-core/src/airflow/jobs/scheduler_job_runner.py#L2378-L2395)) | Seq Scan, 99,500 rows filtered, 935 buffers for the scan, 3.72 ms | Index Scan on `idx_apdr_created_dag_run_id_created_at_id` (`created_dag_run_id IS NULL`), 11 buffers for the scan, 0.43 ms | | Deleting 100 `dag_run` rows: the `ON DELETE CASCADE` lookup through `apdr_created_dag_run_id_fkey` ([models/asset.py#L983-L988](https://github.com/apache/airflow/blob/8f3e8466c67c0d4bf5db4bdbb99fdce50f436e4f/airflow-core/src/airflow/models/asset.py#L983-L988)), which `airflow db clean` and Dag deletion hit | FK trigger 360 ms for 100 calls (one sequential scan per deleted run) | FK trigger 0.42 ms | The ~1,000 remaining buffers in the "after" pending scan are `LockRows` for the 500 rows (`FOR UPDATE SKIP LOCKED`), the same in both plans. The lookup runs once per key inside the task-success request, and the cascade runs once per deleted `dag_run`, so both scale with table size times call count. On MySQL, InnoDB already keeps an implicit index on the foreign-key column, which covers the cascade and the `IS NULL` filter there; the lookup gain applies to every backend. Write cost, for completeness: inserting 10,000 rows takes about 55 ms without the two indexes and about 136 ms with them (medians of five runs), roughly 8 µs per row, against roughly 5 ms per key for the registration request itself. The same predicates are also used by the UI's pending-partition count and next-pending lookup ([routes/ui/assets.py#L183-L190](https://github.com/apache/airflow/blob/8f3e8466c67c0d4bf5db4bdbb99fdce50f436e4f/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/assets.py#L183-L190), [#L262-L268](https://github.com/apache/airflow/blob/8f3e8466c67c0d4bf5db4bdbb99fdce50f436e4f/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/assets.py#L262-L268)) and the partitioned Dag run detail ([routes/ui/partitioned_dag_runs.py#L388-L392](https://github.com/apache/airflow/blob/8f3e8466c67c0d4bf5db4bdbb99fdce50f436e4f/airflow-core/src/airflow/api_fastapi/core_api/routes/ui/partitioned_dag_runs.py#L388-L392)). Script, seed and raw output to reproduce: https://github.com/caoterry/airflow-100k-partitions/tree/main/bench/results/pr73983_explain <details><summary>Plans: before (revision 90e4d18ccadf)</summary> ``` -- before: lookup (_get_or_create_apdr) Limit (cost=2435.01..2435.02 rows=1 width=83) (actual time=6.949..6.949 rows=1 loops=1) Buffers: shared hit=935 -> Sort (cost=2435.01..2435.02 rows=1 width=83) (actual time=6.948..6.948 rows=1 loops=1) Sort Key: id DESC Sort Method: quicksort Memory: 25kB Buffers: shared hit=935 -> Seq Scan on asset_partition_dag_run (cost=0.00..2435.00 rows=1 width=83) (actual time=3.368..6.945 rows=1 loops=1) Filter: (((partition_key)::text = 'acct_50000'::text) AND ((target_dag_id)::text = 'consumer'::text)) Rows Removed by Filter: 99999 Buffers: shared hit=935 Planning Time: 0.045 ms Execution Time: 6.957 ms -- before: pending scan (scheduler loop) Limit (cost=1965.32..1971.57 rows=500 width=95) (actual time=3.459..3.687 rows=500 loops=1) Buffers: shared hit=1936 -> LockRows (cost=1965.32..1971.70 rows=510 width=95) (actual time=3.459..3.652 rows=500 loops=1) Buffers: shared hit=1936 -> Sort (cost=1965.32..1966.60 rows=510 width=95) (actual time=3.451..3.470 rows=500 loops=1) Sort Key: asset_partition_dag_run.created_at, asset_partition_dag_run.id Sort Method: quicksort Memory: 67kB Buffers: shared hit=936 -> Nested Loop (cost=0.00..1942.38 rows=510 width=95) (actual time=3.295..3.388 rows=500 loops=1) Join Filter: ((dag.dag_id)::text = (asset_partition_dag_run.target_dag_id)::text) Buffers: shared hit=936 -> Seq Scan on dag (cost=0.00..1.01 rows=1 width=15) (actual time=0.002..0.003 rows=1 loops=1) Filter: ((is_paused IS FALSE) AND (is_draining IS FALSE) AND (is_stale IS FALSE)) Buffers: shared hit=1 -> Seq Scan on asset_partition_dag_run (cost=0.00..1935.00 rows=510 width=89) (actual time=3.291..3.337 rows=500 loops=1) Filter: (created_dag_run_id IS NULL) Rows Removed by Filter: 99500 Buffers: shared hit=935 Planning Time: 0.042 ms Execution Time: 3.720 ms -- before: dag_run delete, 100 runs (ON DELETE CASCADE into APDR) Delete on dag_run (cost=6.11..765.83 rows=0 width=0) (actual time=0.164..0.165 rows=0 loops=1) Buffers: shared hit=504 -> Nested Loop (cost=6.11..765.83 rows=100 width=34) (actual time=0.045..0.110 rows=100 loops=1) Buffers: shared hit=304 -> HashAggregate (cost=5.82..6.82 rows=100 width=32) (actual time=0.043..0.049 rows=100 loops=1) Group Key: "ANY_subquery".id Batches: 1 Memory Usage: 24kB Buffers: shared hit=4 -> Subquery Scan on "ANY_subquery" (cost=0.29..5.57 rows=100 width=32) (actual time=0.006..0.031 rows=100 loops=1) Buffers: shared hit=4 -> Limit (cost=0.29..4.57 rows=100 width=4) (actual time=0.004..0.020 rows=100 loops=1) Buffers: shared hit=4 -> Index Scan using dag_run_pkey on dag_run dag_run_1 (cost=0.29..4252.54 rows=99500 width=4) (actual time=0.003..0.014 rows=100 loops=1) Filter: ((dag_id)::text = 'consumer'::text) Buffers: shared hit=4 -> Index Scan using dag_run_pkey on dag_run (cost=0.29..7.59 rows=1 width=10) (actual time=0.000..0.000 rows=1 loops=100) Index Cond: (id = "ANY_subquery".id) Buffers: shared hit=300 Planning: Buffers: shared hit=7 Planning Time: 0.112 ms Trigger for constraint dag_run_note_dr_fkey: time=0.369 calls=100 Trigger for constraint task_instance_dag_run_fkey: time=0.408 calls=100 Trigger for constraint bdr_dag_run_fkey: time=0.300 calls=100 Trigger for constraint dagrun_asset_event_dag_run_id_fkey: time=0.569 calls=100 Trigger for constraint deadline_dagrun_id_fkey: time=0.274 calls=100 Trigger for constraint apdr_created_dag_run_id_fkey: time=360.463 calls=100 Trigger for constraint task_state_store_dag_run_fkey: time=0.894 calls=100 Execution Time: 363.531 ms ``` </details> <details><summary>Plans: after (this PR)</summary> ``` -- after: lookup (_get_or_create_apdr) Limit (cost=0.42..8.44 rows=1 width=83) (actual time=0.005..0.005 rows=1 loops=1) Buffers: shared hit=4 -> Index Scan Backward using idx_apdr_target_dag_id_partition_key_id on asset_partition_dag_run (cost=0.42..8.44 rows=1 width=83) (actual time=0.005..0.005 rows=1 loops=1) Index Cond: (((target_dag_id)::text = 'consumer'::text) AND ((partition_key)::text = 'acct_50000'::text)) Buffers: shared hit=4 Planning Time: 0.020 ms Execution Time: 0.008 ms -- after: pending scan (scheduler loop) Limit (cost=690.25..695.91 rows=453 width=95) (actual time=0.175..0.403 rows=500 loops=1) Buffers: shared hit=1012 -> LockRows (cost=690.25..695.91 rows=453 width=95) (actual time=0.174..0.375 rows=500 loops=1) Buffers: shared hit=1012 -> Sort (cost=690.25..691.38 rows=453 width=95) (actual time=0.168..0.186 rows=500 loops=1) Sort Key: asset_partition_dag_run.created_at, asset_partition_dag_run.id Sort Method: quicksort Memory: 67kB Buffers: shared hit=12 -> Nested Loop (cost=0.42..670.27 rows=453 width=95) (actual time=0.006..0.107 rows=500 loops=1) Join Filter: ((asset_partition_dag_run.target_dag_id)::text = (dag.dag_id)::text) Buffers: shared hit=12 -> Seq Scan on dag (cost=0.00..1.01 rows=1 width=15) (actual time=0.002..0.002 rows=1 loops=1) Filter: ((is_paused IS FALSE) AND (is_draining IS FALSE) AND (is_stale IS FALSE)) Buffers: shared hit=1 -> Index Scan using idx_apdr_created_dag_run_id_created_at_id on asset_partition_dag_run (cost=0.42..663.60 rows=453 width=89) (actual time=0.003..0.059 rows=500 loops=1) Index Cond: (created_dag_run_id IS NULL) Buffers: shared hit=11 Planning Time: 0.037 ms Execution Time: 0.434 ms -- after: dag_run delete, 100 runs (ON DELETE CASCADE into APDR) Delete on dag_run (cost=6.11..765.83 rows=0 width=0) (actual time=0.218..0.218 rows=0 loops=1) Buffers: shared hit=504 -> Nested Loop (cost=6.11..765.83 rows=100 width=34) (actual time=0.052..0.138 rows=100 loops=1) Buffers: shared hit=304 -> HashAggregate (cost=5.82..6.82 rows=100 width=32) (actual time=0.050..0.059 rows=100 loops=1) Group Key: "ANY_subquery".id Batches: 1 Memory Usage: 24kB Buffers: shared hit=4 -> Subquery Scan on "ANY_subquery" (cost=0.29..5.57 rows=100 width=32) (actual time=0.005..0.036 rows=100 loops=1) Buffers: shared hit=4 -> Limit (cost=0.29..4.57 rows=100 width=4) (actual time=0.004..0.024 rows=100 loops=1) Buffers: shared hit=4 -> Index Scan using dag_run_pkey on dag_run dag_run_1 (cost=0.29..4252.54 rows=99500 width=4) (actual time=0.004..0.017 rows=100 loops=1) Filter: ((dag_id)::text = 'consumer'::text) Buffers: shared hit=4 -> Index Scan using dag_run_pkey on dag_run (cost=0.29..7.59 rows=1 width=10) (actual time=0.001..0.001 rows=1 loops=100) Index Cond: (id = "ANY_subquery".id) Buffers: shared hit=300 Planning: Buffers: shared hit=7 Planning Time: 0.092 ms Trigger for constraint dag_run_note_dr_fkey: time=0.277 calls=100 Trigger for constraint task_instance_dag_run_fkey: time=0.328 calls=100 Trigger for constraint bdr_dag_run_fkey: time=0.273 calls=100 Trigger for constraint dagrun_asset_event_dag_run_id_fkey: time=0.359 calls=100 Trigger for constraint deadline_dagrun_id_fkey: time=0.232 calls=100 Trigger for constraint apdr_created_dag_run_id_fkey: time=0.420 calls=100 Trigger for constraint task_state_store_dag_run_fkey: time=0.305 calls=100 Execution Time: 2.472 ms ``` </details> <details><summary>Seed SQL</summary> ```sql insert into dag_bundle(name) values ('bench') on conflict do nothing; insert into dag (dag_id, max_active_tasks, has_task_concurrency_limits, is_paused, is_stale, max_consecutive_failed_dag_runs, bundle_name, timetable_type, partition_mapper_info, timetable_partitioned) values ('consumer', 16, false, false, false, 0, 'bench', 'PartitionedAssetTimetable', '{}', true); -- 99,500 fired partition runs, one dag_run each insert into dag_run (dag_id, run_id, run_type, run_after, state, partition_key) select 'consumer', 'p_' || i, 'asset_triggered', now() - (100000 - i) * interval '1 second', 'success', 'acct_' || i from generate_series(1, 99500) i; -- 100,000 APDR rows: 99,500 fired (created_dag_run_id set), the newest 500 pending insert into asset_partition_dag_run (target_dag_id, partition_key, created_at, updated_at, created_dag_run_id) select 'consumer', 'acct_' || i, now() - (100000 - i) * interval '1 second', now() - (100000 - i) * interval '1 second', dr.id from generate_series(1, 100000) i left join dag_run dr on dr.dag_id = 'consumer' and dr.run_id = 'p_' || i; analyze dag; analyze dag_run; analyze asset_partition_dag_run; ``` </details> -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
