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]
