SakshamKapoor2911 commented on PR #70545:
URL: https://github.com/apache/airflow/pull/70545#issuecomment-6031415797

   Rebased onto current `main` (`62917eab`) and re-verified.
   
   **Correction:** an earlier comment described this as a *shutdown* deadlock 
fix. That was inaccurate — this PR addresses the **runtime / dispatch path** 
(`_process_workloads()`), which is the variant reported in #70526 and which 
#67881 (the `end()` drain) does not cover. #67881 fixed the shutdown/join path; 
this adds a `result_queue` drain **before each `activity_queue.put()`** during 
normal dispatch.
   
   **Why:** workers can block writing results into a full `result_queue` pipe → 
they stop consuming `activity_queue` → the scheduler blocks in 
`activity_queue.put()` and, being single-threaded at that point, never reaches 
`_read_results()` → circular wait.
   
   **Change:** one line in `local_executor.py` — `self._read_results()` before 
`self.activity_queue.put(workload)` in `_process_workloads`. The branch had 
gone stale against `main`'s workload-API refactor, so I rebased and adapted the 
tests to the new API (`get_workload_key` / `executor_queues` / 
`_dispatch_counts`).
   
   **Tests (run locally on this branch, `fork` start method, Python 3.12):**
   - `pytest airflow-core/tests/unit/executors/test_local_executor.py` → **59 
passed, 1 skipped**.
   - Added `test_process_workloads_drains_results_before_put` (asserts 
drain-before-dispatch ordering).
   - Added a bounded real-queue regression test 
`test_process_workloads_progresses_under_concurrent_result_refill` that runs 
the actual `_process_workloads` with real `multiprocessing.SimpleQueue`s and a 
worker concurrently refilling results. Under identical conditions it **fails 
(times out) without the patch and completes with it** (measured pipe capacity 
`F_GETPIPE_SZ = 65536`, kernel 6.8.0).
   
   I have **not** verified a production deployment — this is a synthetic 
reproduction that passes with the patch. Feedback on the approach is welcome.
   


-- 
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