Fix triggerer deadlocks - #51279

Closed
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock
Closed

Fix triggerer deadlocks#51279
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock

Conversation

@gopidesupavan

@gopidesupavangopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
Member

The triggers getting deadlock when using sync functions with sync_to_async. To avoid that we have couple of solutions discussed in here #50185.

Use the ThreadPoolExecutor to read trigger workloads and the future object will be used to wait in get_message, this will we can avoid collisions as described here #50185 (comment)


^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in airflow-core/newsfragments.

Comment threadtask-sdk/src/airflow/sdk/execution_time/task_runner.py Outdated
@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

Need to add some tests. added

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor
image

@x42005e1f

Copy link
Copy Markdown

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

 async with SUPERVISOR_COMMS.lock:
self.requests_sock.write(msg.model_dump_json(exclude_none=True).encode() + b"\n")
TRIGGERER_SUPERVISOR_COMMS_FUTURE = self._stdin_threadpool_executor.submit(
SUPERVISOR_COMMS._read_stdin_line
)
line = await asyncio.wrap_future(TRIGGERER_SUPERVISOR_COMMS_FUTURE)
TRIGGERER_SUPERVISOR_COMMS_FUTURE = None # type: ignore[assignment]

?

@x42005e1f

Copy link
Copy Markdown

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

@x42005e1f

Copy link
Copy Markdown

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

Yeah agree :)

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch 2 times, most recently from 7e108a1 to 8a9433cCompareJune 2, 2025 12:30
yield TriggerEvent({"count": dag_run_states_count, "dag_run_state": dag_run_state})


@pytest.mark.xfail(

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Tests are passing now, xfail not required .

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch from 34d1ba3 to 3701d04CompareJune 3, 2025 09:16
Comment on lines -808 to +812
async def connect_stdin() -> asyncio.StreamReader:
reader = asyncio.StreamReader()
protocol = asyncio.StreamReaderProtocol(reader)
await loop.connect_read_pipe(lambda: protocol, sys.stdin)
return reader

self.response_sock = await connect_stdin()

line = await self.response_sock.readline()
msg = comms_decoder.get_message()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why was this changed?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sys.stdin is already configured to comms here https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR805.

I think its fine to read from get_message?

global TRIGGERER_SUPERVISOR_COMMS_FUTURE
line = None

if TRIGGERER_SUPERVISOR_COMMS_FUTURE is not None:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I really don't love this being in here. It feels like a massive abstraction leak. I think we should instead subclass CommsDecoder into a new class defined/living somewhere with the triggerer code.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agree, happy to do with subclass..

@ashb

ashb commented Jun 3, 2025

Copy link
Copy Markdown
Member

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

@x42005e1f

Copy link
Copy Markdown

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

Let me try to explain, since I was the initiator of this change.

The problem is that synchronous and asynchronous lock calls can coexist in an asynchronous context. When an asynchronous task, holding the lock asynchronously, switches contexts, another task may try to acquire the lock synchronously (for some other request). The result is a deadlock - the attempt to acquire the lock synchronously cannot complete until the asynchronous task completes, and the asynchronous task cannot complete because the event loop is blocked by the synchronous call. ThreadPoolExecutor allows to delegate the first (asynchronous) call to a worker thread, and as a result it will be able to complete without switching to the asynchronous task, which will allow to bypass deadlock. Calling future.result() for a future object created by an asynchronous task in the same thread is necessary to ensure no collisions.

There are two cleaner solutions. The first one is to use ThreadPoolExecutor for each send_request() + get_message() - in this case we can get rid of the lock altogether. The second one is to make an asynchronous version of each method that calls send_request() + get_message(), but this is harder to implement and may not always be possible.

I will also clarify that this PR is incomplete without using a more specific type of lock that allows the async -> sync case (synchronous acquiring after asynchronous one in the same thread).

@x42005e1f

Copy link
Copy Markdown

Also note that it is possible to use sync_to_async() (where the lock will be acquired in the worker thread) instead of an explicit ThreadPoolExecutor. This method has less flexibility, because for communication it will be necessary to use only synchronous send_request() + get_message(), but it does not require storing a future object and using a special type of lock (moreover, it can be downgraded to threading.Lock). The method used in this PR can be used with asyncio tools, but to do so you need to access what's under their hood.

@x42005e1f

Copy link
Copy Markdown

In general, it is impossible to solve the synchronization problem between synchronous and asynchronous code in the same thread when synchronous code refers to threading, due to the specifics of cooperative multitasking implementation. Synchronous calls will always block the event loop, and the blocked event loop will prevent asynchronous tasks from executing that could have completed these synchronous calls. Turning synchronous calls into implicitly asynchronous ones (eventlet and gevent approach) leads to coroutine-safety violation. So the solutions are either to reduce this type of synchronization or to delegate execution to a worker thread.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

@x42005e1f Thanks for the response :).

@ashb is that looks fine ? you have any suggestions are alternatives for this please?

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Trigger runner process locked with multiple Workflow triggers

3 participants

@gopidesupavan@x42005e1f@ashb
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Fix triggerer deadlocks - #51279

Closed
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock
Closed

Fix triggerer deadlocks#51279
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock

Conversation

@gopidesupavan

@gopidesupavangopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
Member

The triggers getting deadlock when using sync functions with sync_to_async. To avoid that we have couple of solutions discussed in here #50185.

Use the ThreadPoolExecutor to read trigger workloads and the future object will be used to wait in get_message, this will we can avoid collisions as described here #50185 (comment)


^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in airflow-core/newsfragments.

Comment threadtask-sdk/src/airflow/sdk/execution_time/task_runner.py Outdated
@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

Need to add some tests. added

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor
image

@x42005e1f

Copy link
Copy Markdown

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

 async with SUPERVISOR_COMMS.lock:
self.requests_sock.write(msg.model_dump_json(exclude_none=True).encode() + b"\n")
TRIGGERER_SUPERVISOR_COMMS_FUTURE = self._stdin_threadpool_executor.submit(
SUPERVISOR_COMMS._read_stdin_line
)
line = await asyncio.wrap_future(TRIGGERER_SUPERVISOR_COMMS_FUTURE)
TRIGGERER_SUPERVISOR_COMMS_FUTURE = None # type: ignore[assignment]

?

@x42005e1f

Copy link
Copy Markdown

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

@x42005e1f

Copy link
Copy Markdown

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

Yeah agree :)

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch 2 times, most recently from 7e108a1 to 8a9433cCompareJune 2, 2025 12:30
yield TriggerEvent({"count": dag_run_states_count, "dag_run_state": dag_run_state})


@pytest.mark.xfail(

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Tests are passing now, xfail not required .

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch from 34d1ba3 to 3701d04CompareJune 3, 2025 09:16
Comment on lines -808 to +812
async def connect_stdin() -> asyncio.StreamReader:
reader = asyncio.StreamReader()
protocol = asyncio.StreamReaderProtocol(reader)
await loop.connect_read_pipe(lambda: protocol, sys.stdin)
return reader

self.response_sock = await connect_stdin()

line = await self.response_sock.readline()
msg = comms_decoder.get_message()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why was this changed?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sys.stdin is already configured to comms here https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR805.

I think its fine to read from get_message?

global TRIGGERER_SUPERVISOR_COMMS_FUTURE
line = None

if TRIGGERER_SUPERVISOR_COMMS_FUTURE is not None:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I really don't love this being in here. It feels like a massive abstraction leak. I think we should instead subclass CommsDecoder into a new class defined/living somewhere with the triggerer code.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agree, happy to do with subclass..

@ashb

ashb commented Jun 3, 2025

Copy link
Copy Markdown
Member

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

@x42005e1f

Copy link
Copy Markdown

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

Let me try to explain, since I was the initiator of this change.

The problem is that synchronous and asynchronous lock calls can coexist in an asynchronous context. When an asynchronous task, holding the lock asynchronously, switches contexts, another task may try to acquire the lock synchronously (for some other request). The result is a deadlock - the attempt to acquire the lock synchronously cannot complete until the asynchronous task completes, and the asynchronous task cannot complete because the event loop is blocked by the synchronous call. ThreadPoolExecutor allows to delegate the first (asynchronous) call to a worker thread, and as a result it will be able to complete without switching to the asynchronous task, which will allow to bypass deadlock. Calling future.result() for a future object created by an asynchronous task in the same thread is necessary to ensure no collisions.

There are two cleaner solutions. The first one is to use ThreadPoolExecutor for each send_request() + get_message() - in this case we can get rid of the lock altogether. The second one is to make an asynchronous version of each method that calls send_request() + get_message(), but this is harder to implement and may not always be possible.

I will also clarify that this PR is incomplete without using a more specific type of lock that allows the async -> sync case (synchronous acquiring after asynchronous one in the same thread).

@x42005e1f

Copy link
Copy Markdown

Also note that it is possible to use sync_to_async() (where the lock will be acquired in the worker thread) instead of an explicit ThreadPoolExecutor. This method has less flexibility, because for communication it will be necessary to use only synchronous send_request() + get_message(), but it does not require storing a future object and using a special type of lock (moreover, it can be downgraded to threading.Lock). The method used in this PR can be used with asyncio tools, but to do so you need to access what's under their hood.

@x42005e1f

Copy link
Copy Markdown

In general, it is impossible to solve the synchronization problem between synchronous and asynchronous code in the same thread when synchronous code refers to threading, due to the specifics of cooperative multitasking implementation. Synchronous calls will always block the event loop, and the blocked event loop will prevent asynchronous tasks from executing that could have completed these synchronous calls. Turning synchronous calls into implicitly asynchronous ones (eventlet and gevent approach) leads to coroutine-safety violation. So the solutions are either to reduce this type of synchronization or to delegate execution to a worker thread.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

@x42005e1f Thanks for the response :).

@ashb is that looks fine ? you have any suggestions are alternatives for this please?

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Trigger runner process locked with multiple Workflow triggers

3 participants

@gopidesupavan@x42005e1f@ashb
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Fix triggerer deadlocks - #51279

Closed
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock
Closed

Fix triggerer deadlocks#51279
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock

Conversation

@gopidesupavan

@gopidesupavangopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
Member

The triggers getting deadlock when using sync functions with sync_to_async. To avoid that we have couple of solutions discussed in here #50185.

Use the ThreadPoolExecutor to read trigger workloads and the future object will be used to wait in get_message, this will we can avoid collisions as described here #50185 (comment)


^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in airflow-core/newsfragments.

Comment threadtask-sdk/src/airflow/sdk/execution_time/task_runner.py Outdated
@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

Need to add some tests. added

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor
image

@x42005e1f

Copy link
Copy Markdown

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

 async with SUPERVISOR_COMMS.lock:
self.requests_sock.write(msg.model_dump_json(exclude_none=True).encode() + b"\n")
TRIGGERER_SUPERVISOR_COMMS_FUTURE = self._stdin_threadpool_executor.submit(
SUPERVISOR_COMMS._read_stdin_line
)
line = await asyncio.wrap_future(TRIGGERER_SUPERVISOR_COMMS_FUTURE)
TRIGGERER_SUPERVISOR_COMMS_FUTURE = None # type: ignore[assignment]

?

@x42005e1f

Copy link
Copy Markdown

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

@x42005e1f

Copy link
Copy Markdown

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

Yeah agree :)

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch 2 times, most recently from 7e108a1 to 8a9433cCompareJune 2, 2025 12:30
yield TriggerEvent({"count": dag_run_states_count, "dag_run_state": dag_run_state})


@pytest.mark.xfail(

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Tests are passing now, xfail not required .

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch from 34d1ba3 to 3701d04CompareJune 3, 2025 09:16
Comment on lines -808 to +812
async def connect_stdin() -> asyncio.StreamReader:
reader = asyncio.StreamReader()
protocol = asyncio.StreamReaderProtocol(reader)
await loop.connect_read_pipe(lambda: protocol, sys.stdin)
return reader

self.response_sock = await connect_stdin()

line = await self.response_sock.readline()
msg = comms_decoder.get_message()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why was this changed?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sys.stdin is already configured to comms here https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR805.

I think its fine to read from get_message?

global TRIGGERER_SUPERVISOR_COMMS_FUTURE
line = None

if TRIGGERER_SUPERVISOR_COMMS_FUTURE is not None:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I really don't love this being in here. It feels like a massive abstraction leak. I think we should instead subclass CommsDecoder into a new class defined/living somewhere with the triggerer code.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agree, happy to do with subclass..

@ashb

ashb commented Jun 3, 2025

Copy link
Copy Markdown
Member

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

@x42005e1f

Copy link
Copy Markdown

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

Let me try to explain, since I was the initiator of this change.

The problem is that synchronous and asynchronous lock calls can coexist in an asynchronous context. When an asynchronous task, holding the lock asynchronously, switches contexts, another task may try to acquire the lock synchronously (for some other request). The result is a deadlock - the attempt to acquire the lock synchronously cannot complete until the asynchronous task completes, and the asynchronous task cannot complete because the event loop is blocked by the synchronous call. ThreadPoolExecutor allows to delegate the first (asynchronous) call to a worker thread, and as a result it will be able to complete without switching to the asynchronous task, which will allow to bypass deadlock. Calling future.result() for a future object created by an asynchronous task in the same thread is necessary to ensure no collisions.

There are two cleaner solutions. The first one is to use ThreadPoolExecutor for each send_request() + get_message() - in this case we can get rid of the lock altogether. The second one is to make an asynchronous version of each method that calls send_request() + get_message(), but this is harder to implement and may not always be possible.

I will also clarify that this PR is incomplete without using a more specific type of lock that allows the async -> sync case (synchronous acquiring after asynchronous one in the same thread).

@x42005e1f

Copy link
Copy Markdown

Also note that it is possible to use sync_to_async() (where the lock will be acquired in the worker thread) instead of an explicit ThreadPoolExecutor. This method has less flexibility, because for communication it will be necessary to use only synchronous send_request() + get_message(), but it does not require storing a future object and using a special type of lock (moreover, it can be downgraded to threading.Lock). The method used in this PR can be used with asyncio tools, but to do so you need to access what's under their hood.

@x42005e1f

Copy link
Copy Markdown

In general, it is impossible to solve the synchronization problem between synchronous and asynchronous code in the same thread when synchronous code refers to threading, due to the specifics of cooperative multitasking implementation. Synchronous calls will always block the event loop, and the blocked event loop will prevent asynchronous tasks from executing that could have completed these synchronous calls. Turning synchronous calls into implicitly asynchronous ones (eventlet and gevent approach) leads to coroutine-safety violation. So the solutions are either to reduce this type of synchronization or to delegate execution to a worker thread.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

@x42005e1f Thanks for the response :).

@ashb is that looks fine ? you have any suggestions are alternatives for this please?

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Trigger runner process locked with multiple Workflow triggers

3 participants

@gopidesupavan@x42005e1f@ashb
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Fix triggerer deadlocks - #51279

Closed
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock
Closed

Fix triggerer deadlocks#51279
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock

Conversation

@gopidesupavan

@gopidesupavangopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
Member

The triggers getting deadlock when using sync functions with sync_to_async. To avoid that we have couple of solutions discussed in here #50185.

Use the ThreadPoolExecutor to read trigger workloads and the future object will be used to wait in get_message, this will we can avoid collisions as described here #50185 (comment)


^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in airflow-core/newsfragments.

Comment threadtask-sdk/src/airflow/sdk/execution_time/task_runner.py Outdated
@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

Need to add some tests. added

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor
image

@x42005e1f

Copy link
Copy Markdown

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

 async with SUPERVISOR_COMMS.lock:
self.requests_sock.write(msg.model_dump_json(exclude_none=True).encode() + b"\n")
TRIGGERER_SUPERVISOR_COMMS_FUTURE = self._stdin_threadpool_executor.submit(
SUPERVISOR_COMMS._read_stdin_line
)
line = await asyncio.wrap_future(TRIGGERER_SUPERVISOR_COMMS_FUTURE)
TRIGGERER_SUPERVISOR_COMMS_FUTURE = None # type: ignore[assignment]

?

@x42005e1f

Copy link
Copy Markdown

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

@x42005e1f

Copy link
Copy Markdown

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

Yeah agree :)

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch 2 times, most recently from 7e108a1 to 8a9433cCompareJune 2, 2025 12:30
yield TriggerEvent({"count": dag_run_states_count, "dag_run_state": dag_run_state})


@pytest.mark.xfail(

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Tests are passing now, xfail not required .

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch from 34d1ba3 to 3701d04CompareJune 3, 2025 09:16
Comment on lines -808 to +812
async def connect_stdin() -> asyncio.StreamReader:
reader = asyncio.StreamReader()
protocol = asyncio.StreamReaderProtocol(reader)
await loop.connect_read_pipe(lambda: protocol, sys.stdin)
return reader

self.response_sock = await connect_stdin()

line = await self.response_sock.readline()
msg = comms_decoder.get_message()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why was this changed?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sys.stdin is already configured to comms here https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR805.

I think its fine to read from get_message?

global TRIGGERER_SUPERVISOR_COMMS_FUTURE
line = None

if TRIGGERER_SUPERVISOR_COMMS_FUTURE is not None:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I really don't love this being in here. It feels like a massive abstraction leak. I think we should instead subclass CommsDecoder into a new class defined/living somewhere with the triggerer code.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agree, happy to do with subclass..

@ashb

ashb commented Jun 3, 2025

Copy link
Copy Markdown
Member

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

@x42005e1f

Copy link
Copy Markdown

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

Let me try to explain, since I was the initiator of this change.

The problem is that synchronous and asynchronous lock calls can coexist in an asynchronous context. When an asynchronous task, holding the lock asynchronously, switches contexts, another task may try to acquire the lock synchronously (for some other request). The result is a deadlock - the attempt to acquire the lock synchronously cannot complete until the asynchronous task completes, and the asynchronous task cannot complete because the event loop is blocked by the synchronous call. ThreadPoolExecutor allows to delegate the first (asynchronous) call to a worker thread, and as a result it will be able to complete without switching to the asynchronous task, which will allow to bypass deadlock. Calling future.result() for a future object created by an asynchronous task in the same thread is necessary to ensure no collisions.

There are two cleaner solutions. The first one is to use ThreadPoolExecutor for each send_request() + get_message() - in this case we can get rid of the lock altogether. The second one is to make an asynchronous version of each method that calls send_request() + get_message(), but this is harder to implement and may not always be possible.

I will also clarify that this PR is incomplete without using a more specific type of lock that allows the async -> sync case (synchronous acquiring after asynchronous one in the same thread).

@x42005e1f

Copy link
Copy Markdown

Also note that it is possible to use sync_to_async() (where the lock will be acquired in the worker thread) instead of an explicit ThreadPoolExecutor. This method has less flexibility, because for communication it will be necessary to use only synchronous send_request() + get_message(), but it does not require storing a future object and using a special type of lock (moreover, it can be downgraded to threading.Lock). The method used in this PR can be used with asyncio tools, but to do so you need to access what's under their hood.

@x42005e1f

Copy link
Copy Markdown

In general, it is impossible to solve the synchronization problem between synchronous and asynchronous code in the same thread when synchronous code refers to threading, due to the specifics of cooperative multitasking implementation. Synchronous calls will always block the event loop, and the blocked event loop will prevent asynchronous tasks from executing that could have completed these synchronous calls. Turning synchronous calls into implicitly asynchronous ones (eventlet and gevent approach) leads to coroutine-safety violation. So the solutions are either to reduce this type of synchronization or to delegate execution to a worker thread.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

@x42005e1f Thanks for the response :).

@ashb is that looks fine ? you have any suggestions are alternatives for this please?

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Trigger runner process locked with multiple Workflow triggers

3 participants

@gopidesupavan@x42005e1f@ashb
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Fix triggerer deadlocks - #51279

Closed
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock
Closed

Fix triggerer deadlocks#51279
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock

Conversation

@gopidesupavan

@gopidesupavangopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
Member

The triggers getting deadlock when using sync functions with sync_to_async. To avoid that we have couple of solutions discussed in here #50185.

Use the ThreadPoolExecutor to read trigger workloads and the future object will be used to wait in get_message, this will we can avoid collisions as described here #50185 (comment)


^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in airflow-core/newsfragments.

Comment threadtask-sdk/src/airflow/sdk/execution_time/task_runner.py Outdated
@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

Need to add some tests. added

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor
image

@x42005e1f

Copy link
Copy Markdown

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

 async with SUPERVISOR_COMMS.lock:
self.requests_sock.write(msg.model_dump_json(exclude_none=True).encode() + b"\n")
TRIGGERER_SUPERVISOR_COMMS_FUTURE = self._stdin_threadpool_executor.submit(
SUPERVISOR_COMMS._read_stdin_line
)
line = await asyncio.wrap_future(TRIGGERER_SUPERVISOR_COMMS_FUTURE)
TRIGGERER_SUPERVISOR_COMMS_FUTURE = None # type: ignore[assignment]

?

@x42005e1f

Copy link
Copy Markdown

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

@x42005e1f

Copy link
Copy Markdown

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

Yeah agree :)

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch 2 times, most recently from 7e108a1 to 8a9433cCompareJune 2, 2025 12:30
yield TriggerEvent({"count": dag_run_states_count, "dag_run_state": dag_run_state})


@pytest.mark.xfail(

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Tests are passing now, xfail not required .

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch from 34d1ba3 to 3701d04CompareJune 3, 2025 09:16
Comment on lines -808 to +812
async def connect_stdin() -> asyncio.StreamReader:
reader = asyncio.StreamReader()
protocol = asyncio.StreamReaderProtocol(reader)
await loop.connect_read_pipe(lambda: protocol, sys.stdin)
return reader

self.response_sock = await connect_stdin()

line = await self.response_sock.readline()
msg = comms_decoder.get_message()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why was this changed?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sys.stdin is already configured to comms here https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR805.

I think its fine to read from get_message?

global TRIGGERER_SUPERVISOR_COMMS_FUTURE
line = None

if TRIGGERER_SUPERVISOR_COMMS_FUTURE is not None:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I really don't love this being in here. It feels like a massive abstraction leak. I think we should instead subclass CommsDecoder into a new class defined/living somewhere with the triggerer code.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agree, happy to do with subclass..

@ashb

ashb commented Jun 3, 2025

Copy link
Copy Markdown
Member

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

@x42005e1f

Copy link
Copy Markdown

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

Let me try to explain, since I was the initiator of this change.

The problem is that synchronous and asynchronous lock calls can coexist in an asynchronous context. When an asynchronous task, holding the lock asynchronously, switches contexts, another task may try to acquire the lock synchronously (for some other request). The result is a deadlock - the attempt to acquire the lock synchronously cannot complete until the asynchronous task completes, and the asynchronous task cannot complete because the event loop is blocked by the synchronous call. ThreadPoolExecutor allows to delegate the first (asynchronous) call to a worker thread, and as a result it will be able to complete without switching to the asynchronous task, which will allow to bypass deadlock. Calling future.result() for a future object created by an asynchronous task in the same thread is necessary to ensure no collisions.

There are two cleaner solutions. The first one is to use ThreadPoolExecutor for each send_request() + get_message() - in this case we can get rid of the lock altogether. The second one is to make an asynchronous version of each method that calls send_request() + get_message(), but this is harder to implement and may not always be possible.

I will also clarify that this PR is incomplete without using a more specific type of lock that allows the async -> sync case (synchronous acquiring after asynchronous one in the same thread).

@x42005e1f

Copy link
Copy Markdown

Also note that it is possible to use sync_to_async() (where the lock will be acquired in the worker thread) instead of an explicit ThreadPoolExecutor. This method has less flexibility, because for communication it will be necessary to use only synchronous send_request() + get_message(), but it does not require storing a future object and using a special type of lock (moreover, it can be downgraded to threading.Lock). The method used in this PR can be used with asyncio tools, but to do so you need to access what's under their hood.

@x42005e1f

Copy link
Copy Markdown

In general, it is impossible to solve the synchronization problem between synchronous and asynchronous code in the same thread when synchronous code refers to threading, due to the specifics of cooperative multitasking implementation. Synchronous calls will always block the event loop, and the blocked event loop will prevent asynchronous tasks from executing that could have completed these synchronous calls. Turning synchronous calls into implicitly asynchronous ones (eventlet and gevent approach) leads to coroutine-safety violation. So the solutions are either to reduce this type of synchronization or to delegate execution to a worker thread.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

@x42005e1f Thanks for the response :).

@ashb is that looks fine ? you have any suggestions are alternatives for this please?

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Trigger runner process locked with multiple Workflow triggers

3 participants

@gopidesupavan@x42005e1f@ashb
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Fix triggerer deadlocks - #51279

Closed
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock
Closed

Fix triggerer deadlocks#51279
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock

Conversation

@gopidesupavan

@gopidesupavangopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
Member

The triggers getting deadlock when using sync functions with sync_to_async. To avoid that we have couple of solutions discussed in here #50185.

Use the ThreadPoolExecutor to read trigger workloads and the future object will be used to wait in get_message, this will we can avoid collisions as described here #50185 (comment)


^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in airflow-core/newsfragments.

Comment threadtask-sdk/src/airflow/sdk/execution_time/task_runner.py Outdated
@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

Need to add some tests. added

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor
image

@x42005e1f

Copy link
Copy Markdown

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

 async with SUPERVISOR_COMMS.lock:
self.requests_sock.write(msg.model_dump_json(exclude_none=True).encode() + b"\n")
TRIGGERER_SUPERVISOR_COMMS_FUTURE = self._stdin_threadpool_executor.submit(
SUPERVISOR_COMMS._read_stdin_line
)
line = await asyncio.wrap_future(TRIGGERER_SUPERVISOR_COMMS_FUTURE)
TRIGGERER_SUPERVISOR_COMMS_FUTURE = None # type: ignore[assignment]

?

@x42005e1f

Copy link
Copy Markdown

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

@x42005e1f

Copy link
Copy Markdown

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

Yeah agree :)

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch 2 times, most recently from 7e108a1 to 8a9433cCompareJune 2, 2025 12:30
yield TriggerEvent({"count": dag_run_states_count, "dag_run_state": dag_run_state})


@pytest.mark.xfail(

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Tests are passing now, xfail not required .

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch from 34d1ba3 to 3701d04CompareJune 3, 2025 09:16
Comment on lines -808 to +812
async def connect_stdin() -> asyncio.StreamReader:
reader = asyncio.StreamReader()
protocol = asyncio.StreamReaderProtocol(reader)
await loop.connect_read_pipe(lambda: protocol, sys.stdin)
return reader

self.response_sock = await connect_stdin()

line = await self.response_sock.readline()
msg = comms_decoder.get_message()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why was this changed?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sys.stdin is already configured to comms here https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR805.

I think its fine to read from get_message?

global TRIGGERER_SUPERVISOR_COMMS_FUTURE
line = None

if TRIGGERER_SUPERVISOR_COMMS_FUTURE is not None:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I really don't love this being in here. It feels like a massive abstraction leak. I think we should instead subclass CommsDecoder into a new class defined/living somewhere with the triggerer code.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agree, happy to do with subclass..

@ashb

ashb commented Jun 3, 2025

Copy link
Copy Markdown
Member

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

@x42005e1f

Copy link
Copy Markdown

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

Let me try to explain, since I was the initiator of this change.

The problem is that synchronous and asynchronous lock calls can coexist in an asynchronous context. When an asynchronous task, holding the lock asynchronously, switches contexts, another task may try to acquire the lock synchronously (for some other request). The result is a deadlock - the attempt to acquire the lock synchronously cannot complete until the asynchronous task completes, and the asynchronous task cannot complete because the event loop is blocked by the synchronous call. ThreadPoolExecutor allows to delegate the first (asynchronous) call to a worker thread, and as a result it will be able to complete without switching to the asynchronous task, which will allow to bypass deadlock. Calling future.result() for a future object created by an asynchronous task in the same thread is necessary to ensure no collisions.

There are two cleaner solutions. The first one is to use ThreadPoolExecutor for each send_request() + get_message() - in this case we can get rid of the lock altogether. The second one is to make an asynchronous version of each method that calls send_request() + get_message(), but this is harder to implement and may not always be possible.

I will also clarify that this PR is incomplete without using a more specific type of lock that allows the async -> sync case (synchronous acquiring after asynchronous one in the same thread).

@x42005e1f

Copy link
Copy Markdown

Also note that it is possible to use sync_to_async() (where the lock will be acquired in the worker thread) instead of an explicit ThreadPoolExecutor. This method has less flexibility, because for communication it will be necessary to use only synchronous send_request() + get_message(), but it does not require storing a future object and using a special type of lock (moreover, it can be downgraded to threading.Lock). The method used in this PR can be used with asyncio tools, but to do so you need to access what's under their hood.

@x42005e1f

Copy link
Copy Markdown

In general, it is impossible to solve the synchronization problem between synchronous and asynchronous code in the same thread when synchronous code refers to threading, due to the specifics of cooperative multitasking implementation. Synchronous calls will always block the event loop, and the blocked event loop will prevent asynchronous tasks from executing that could have completed these synchronous calls. Turning synchronous calls into implicitly asynchronous ones (eventlet and gevent approach) leads to coroutine-safety violation. So the solutions are either to reduce this type of synchronization or to delegate execution to a worker thread.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

@x42005e1f Thanks for the response :).

@ashb is that looks fine ? you have any suggestions are alternatives for this please?

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Trigger runner process locked with multiple Workflow triggers

3 participants

@gopidesupavan@x42005e1f@ashb
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Fix triggerer deadlocks - #51279

Closed
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock
Closed

Fix triggerer deadlocks#51279
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock

Conversation

@gopidesupavan

@gopidesupavangopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
Member

The triggers getting deadlock when using sync functions with sync_to_async. To avoid that we have couple of solutions discussed in here #50185.

Use the ThreadPoolExecutor to read trigger workloads and the future object will be used to wait in get_message, this will we can avoid collisions as described here #50185 (comment)


^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in airflow-core/newsfragments.

Comment threadtask-sdk/src/airflow/sdk/execution_time/task_runner.py Outdated
@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

Need to add some tests. added

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor
image

@x42005e1f

Copy link
Copy Markdown

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

 async with SUPERVISOR_COMMS.lock:
self.requests_sock.write(msg.model_dump_json(exclude_none=True).encode() + b"\n")
TRIGGERER_SUPERVISOR_COMMS_FUTURE = self._stdin_threadpool_executor.submit(
SUPERVISOR_COMMS._read_stdin_line
)
line = await asyncio.wrap_future(TRIGGERER_SUPERVISOR_COMMS_FUTURE)
TRIGGERER_SUPERVISOR_COMMS_FUTURE = None # type: ignore[assignment]

?

@x42005e1f

Copy link
Copy Markdown

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

@x42005e1f

Copy link
Copy Markdown

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

Yeah agree :)

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch 2 times, most recently from 7e108a1 to 8a9433cCompareJune 2, 2025 12:30
yield TriggerEvent({"count": dag_run_states_count, "dag_run_state": dag_run_state})


@pytest.mark.xfail(

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Tests are passing now, xfail not required .

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch from 34d1ba3 to 3701d04CompareJune 3, 2025 09:16
Comment on lines -808 to +812
async def connect_stdin() -> asyncio.StreamReader:
reader = asyncio.StreamReader()
protocol = asyncio.StreamReaderProtocol(reader)
await loop.connect_read_pipe(lambda: protocol, sys.stdin)
return reader

self.response_sock = await connect_stdin()

line = await self.response_sock.readline()
msg = comms_decoder.get_message()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why was this changed?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sys.stdin is already configured to comms here https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR805.

I think its fine to read from get_message?

global TRIGGERER_SUPERVISOR_COMMS_FUTURE
line = None

if TRIGGERER_SUPERVISOR_COMMS_FUTURE is not None:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I really don't love this being in here. It feels like a massive abstraction leak. I think we should instead subclass CommsDecoder into a new class defined/living somewhere with the triggerer code.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agree, happy to do with subclass..

@ashb

ashb commented Jun 3, 2025

Copy link
Copy Markdown
Member

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

@x42005e1f

Copy link
Copy Markdown

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

Let me try to explain, since I was the initiator of this change.

The problem is that synchronous and asynchronous lock calls can coexist in an asynchronous context. When an asynchronous task, holding the lock asynchronously, switches contexts, another task may try to acquire the lock synchronously (for some other request). The result is a deadlock - the attempt to acquire the lock synchronously cannot complete until the asynchronous task completes, and the asynchronous task cannot complete because the event loop is blocked by the synchronous call. ThreadPoolExecutor allows to delegate the first (asynchronous) call to a worker thread, and as a result it will be able to complete without switching to the asynchronous task, which will allow to bypass deadlock. Calling future.result() for a future object created by an asynchronous task in the same thread is necessary to ensure no collisions.

There are two cleaner solutions. The first one is to use ThreadPoolExecutor for each send_request() + get_message() - in this case we can get rid of the lock altogether. The second one is to make an asynchronous version of each method that calls send_request() + get_message(), but this is harder to implement and may not always be possible.

I will also clarify that this PR is incomplete without using a more specific type of lock that allows the async -> sync case (synchronous acquiring after asynchronous one in the same thread).

@x42005e1f

Copy link
Copy Markdown

Also note that it is possible to use sync_to_async() (where the lock will be acquired in the worker thread) instead of an explicit ThreadPoolExecutor. This method has less flexibility, because for communication it will be necessary to use only synchronous send_request() + get_message(), but it does not require storing a future object and using a special type of lock (moreover, it can be downgraded to threading.Lock). The method used in this PR can be used with asyncio tools, but to do so you need to access what's under their hood.

@x42005e1f

Copy link
Copy Markdown

In general, it is impossible to solve the synchronization problem between synchronous and asynchronous code in the same thread when synchronous code refers to threading, due to the specifics of cooperative multitasking implementation. Synchronous calls will always block the event loop, and the blocked event loop will prevent asynchronous tasks from executing that could have completed these synchronous calls. Turning synchronous calls into implicitly asynchronous ones (eventlet and gevent approach) leads to coroutine-safety violation. So the solutions are either to reduce this type of synchronization or to delegate execution to a worker thread.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

@x42005e1f Thanks for the response :).

@ashb is that looks fine ? you have any suggestions are alternatives for this please?

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Trigger runner process locked with multiple Workflow triggers

3 participants

@gopidesupavan@x42005e1f@ashb
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Fix triggerer deadlocks - #51279

Closed
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock
Closed

Fix triggerer deadlocks#51279
gopidesupavan wants to merge 5 commits into
apache:mainfrom
gopidesupavan:fix-triggerer-comms-deadlock

Conversation

@gopidesupavan

@gopidesupavangopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
Member

The triggers getting deadlock when using sync functions with sync_to_async. To avoid that we have couple of solutions discussed in here #50185.

Use the ThreadPoolExecutor to read trigger workloads and the future object will be used to wait in get_message, this will we can avoid collisions as described here #50185 (comment)


^ Add meaningful description above
Read the Pull Request Guidelines for more information.
In case of fundamental code changes, an Airflow Improvement Proposal (AIP) is needed.
In case of a new dependency, check compliance with the ASF 3rd Party License Policy.
In case of backwards incompatible changes please leave a note in a newsfragment file, named {pr_number}.significant.rst or {issue_number}.significant.rst, in airflow-core/newsfragments.

Comment threadtask-sdk/src/airflow/sdk/execution_time/task_runner.py Outdated
@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

Need to add some tests. added

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor
image

@x42005e1f

Copy link
Copy Markdown

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

@gopidesupavan

gopidesupavan commented Jun 2, 2025

Copy link
Copy Markdown
MemberAuthor

The future pass itself looks right, however, here you still need to use the mentioned lock type, which I described in the linked comment. The approach without synchronization is appropriate only in case of full use of futures - when each send_request() + get_message() are executed in a worker thread.

I can write a separate lock implementation later if you would like.

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

 async with SUPERVISOR_COMMS.lock:
self.requests_sock.write(msg.model_dump_json(exclude_none=True).encode() + b"\n")
TRIGGERER_SUPERVISOR_COMMS_FUTURE = self._stdin_threadpool_executor.submit(
SUPERVISOR_COMMS._read_stdin_line
)
line = await asyncio.wrap_future(TRIGGERER_SUPERVISOR_COMMS_FUTURE)
TRIGGERER_SUPERVISOR_COMMS_FUTURE = None # type: ignore[assignment]

?

@x42005e1f

Copy link
Copy Markdown

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

ah you mean lock here before threadpool https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR963 ?

Yes, and in get_ti_count() too, since it can be used in separate threads.

The thread-level lock approach is special in that all uses of the lock remain, but the lock itself changes, special handling for async -> sync is added. Synchronization is still needed to eliminate collisions.

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

@x42005e1f

Copy link
Copy Markdown

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

Yeah you correct, i just ran with some multiple dags i could see difference, without lock :)

Multithreaded issues are usually hard to reproduce - it is often much easier to take a formal approach to them. This is why I would advise not to trust tests, at least not specialized ones - they can lie.

Yeah agree :)

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch 2 times, most recently from 7e108a1 to 8a9433cCompareJune 2, 2025 12:30
yield TriggerEvent({"count": dag_run_states_count, "dag_run_state": dag_run_state})


@pytest.mark.xfail(

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Tests are passing now, xfail not required .

@gopidesupavan
gopidesupavanforce-pushed the fix-triggerer-comms-deadlock branch from 34d1ba3 to 3701d04CompareJune 3, 2025 09:16
Comment on lines -808 to +812
async def connect_stdin() -> asyncio.StreamReader:
reader = asyncio.StreamReader()
protocol = asyncio.StreamReaderProtocol(reader)
await loop.connect_read_pipe(lambda: protocol, sys.stdin)
return reader

self.response_sock = await connect_stdin()

line = await self.response_sock.readline()
msg = comms_decoder.get_message()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why was this changed?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

sys.stdin is already configured to comms here https://github.com/apache/airflow/pull/51279/files#diff-e4cc497f1c786d142ce4c930f43e33b0bb4b53d375d274278fa82f4d5567608aR805.

I think its fine to read from get_message?

global TRIGGERER_SUPERVISOR_COMMS_FUTURE
line = None

if TRIGGERER_SUPERVISOR_COMMS_FUTURE is not None:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I really don't love this being in here. It feels like a massive abstraction leak. I think we should instead subclass CommsDecoder into a new class defined/living somewhere with the triggerer code.

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agree, happy to do with subclass..

@ashb

ashb commented Jun 3, 2025

Copy link
Copy Markdown
Member

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

@x42005e1f

Copy link
Copy Markdown

@gopidesupavan Can you explain your reason/thinking for switching to a thread pool? Generally I don't love the use of a threadpool in an async context, especially when we are just making requests (i.e. something asyncio should be really good at), so I'd really like to understand more about why this change was needed.

Let me try to explain, since I was the initiator of this change.

The problem is that synchronous and asynchronous lock calls can coexist in an asynchronous context. When an asynchronous task, holding the lock asynchronously, switches contexts, another task may try to acquire the lock synchronously (for some other request). The result is a deadlock - the attempt to acquire the lock synchronously cannot complete until the asynchronous task completes, and the asynchronous task cannot complete because the event loop is blocked by the synchronous call. ThreadPoolExecutor allows to delegate the first (asynchronous) call to a worker thread, and as a result it will be able to complete without switching to the asynchronous task, which will allow to bypass deadlock. Calling future.result() for a future object created by an asynchronous task in the same thread is necessary to ensure no collisions.

There are two cleaner solutions. The first one is to use ThreadPoolExecutor for each send_request() + get_message() - in this case we can get rid of the lock altogether. The second one is to make an asynchronous version of each method that calls send_request() + get_message(), but this is harder to implement and may not always be possible.

I will also clarify that this PR is incomplete without using a more specific type of lock that allows the async -> sync case (synchronous acquiring after asynchronous one in the same thread).

@x42005e1f

Copy link
Copy Markdown

Also note that it is possible to use sync_to_async() (where the lock will be acquired in the worker thread) instead of an explicit ThreadPoolExecutor. This method has less flexibility, because for communication it will be necessary to use only synchronous send_request() + get_message(), but it does not require storing a future object and using a special type of lock (moreover, it can be downgraded to threading.Lock). The method used in this PR can be used with asyncio tools, but to do so you need to access what's under their hood.

@x42005e1f

Copy link
Copy Markdown

In general, it is impossible to solve the synchronization problem between synchronous and asynchronous code in the same thread when synchronous code refers to threading, due to the specifics of cooperative multitasking implementation. Synchronous calls will always block the event loop, and the blocked event loop will prevent asynchronous tasks from executing that could have completed these synchronous calls. Turning synchronous calls into implicitly asynchronous ones (eventlet and gevent approach) leads to coroutine-safety violation. So the solutions are either to reduce this type of synchronization or to delegate execution to a worker thread.

@gopidesupavan

Copy link
Copy Markdown
MemberAuthor

@x42005e1f Thanks for the response :).

@ashb is that looks fine ? you have any suggestions are alternatives for this please?

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

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Trigger runner process locked with multiple Workflow triggers

3 participants

@gopidesupavan@x42005e1f@ashb