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]

Reply via email to