Repository navigation
Avoid extra_dejson in ADF and Synapse async hooks - #72130
Conversation
extra_dejson can mask secrets via a sync send on the triggerer event loop, which raises AsyncToSync. Parse extras with json.loads instead, matching the MSGraph workaround. closes: apache#55728
This comment was marked as spam.
This comment was marked as spam.
There is no async The is also a PR open which addresses another generic issue regarding the async path with connections as well. |
potiuk
left a comment
There was a problem hiding this comment.
Nice. Simple to accept without proof of running it on real Synapse
|
But @Vamsi-klu - please test it when we release the provider - you will be `@-mentioned' |
extra_dejson can mask secrets via a sync send on the triggerer event loop, which raises AsyncToSync. Parse extras with json.loads instead, matching the MSGraph workaround. closes: apache#55728
get_conn() resolves the connection with get_connection() and reads extra_dejson, whose secret masking sends to the supervisor synchronously; inside an async task with another async SDK call in flight both raise DeadlockImminentError. aget_conn() returns an AsyncOpenAI client configured exactly like get_conn()'s, from get_async_connection() with the extra from Connection.aextra_dejson() (apache#71890, Airflow 3.3.2+), or the raw extra via json.loads on older versions (as in apache#72130). The client settings move into _client_kwargs(), shared by both paths. acreate_embeddings() is the async counterpart of create_embeddings(). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
get_conn() resolves the hook's connection and every fallback connection with get_connection() and reads extra_dejson, whose secret masking sends to the supervisor synchronously; inside an async task with another async SDK call in flight both raise DeadlockImminentError, although the agent the hook builds is async-native. aget_conn() fetches the connection and its fallbacks with get_async_connection() (extra from Connection.aextra_dejson(), apache#71890, Airflow 3.3.2+, or the raw extra via json.loads on older versions as in apache#72130), seeds the existing connection cache and builds the model with the unchanged get_conn(); acreate_agent() does the same before create_agent(), spec_file included. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
get_conn() resolves the connection with get_connection() and reads extra_dejson, whose secret masking sends to the supervisor synchronously; inside an async task with another async SDK call in flight both raise DeadlockImminentError. aget_conn() returns an AsyncOpenAI client configured exactly like get_conn()'s, from get_async_connection() with the extra from Connection.aextra_dejson() (apache#71890, Airflow 3.3.2+), or the raw extra via json.loads on older versions (as in apache#72130). The client settings move into _client_kwargs(), shared by both paths. acreate_embeddings() is the async counterpart of create_embeddings(). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
get_conn() resolves the hook's connection and every fallback connection with get_connection() and reads extra_dejson, whose secret masking sends to the supervisor synchronously; inside an async task with another async SDK call in flight both raise DeadlockImminentError, although the agent the hook builds is async-native. aget_conn() fetches the connection and its fallbacks with get_async_connection() (extra from Connection.aextra_dejson(), apache#71890, Airflow 3.3.2+, or the raw extra via json.loads on older versions as in apache#72130), seeds the existing connection cache and builds the model with the unchanged get_conn(); acreate_agent() does the same before create_agent(), spec_file included. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
get_conn() resolves the hook's connection and every fallback connection with get_connection() and reads extra_dejson, whose secret masking sends to the supervisor synchronously; inside an async task with another async SDK call in flight both raise DeadlockImminentError, although the agent the hook builds is async-native. aget_conn() fetches the connection and its fallbacks with get_async_connection() (extra from Connection.aextra_dejson(), apache#71890, Airflow 3.3.2+, or the raw extra via json.loads on older versions as in apache#72130), seeds the existing connection cache and builds the model with the unchanged get_conn(); acreate_agent() does the same before create_agent(), spec_file included. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…74147) Connection.extra_dejson masks the extra's secrets with a synchronous call to the supervisor, which raises DeadlockImminentError when a hook reads it on an event loop with another async call in flight. Airflow 3.3.2+ has Connection.aextra_dejson() (#71890), but providers that still support older versions cannot call it directly, so async hooks fall back to json.loads(conn.extra) and skip the masking (#72130). get_async_extra_dejson(conn) awaits Connection.aextra_dejson() when it exists, and otherwise runs extra_dejson in a worker thread, where blocking on the supervisor is safe, as get_async_connection() does for get_connection(). Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
* Avoid the worker-thread fallback of get_async_extra_dejson on Airflow 3.0 and 3.1 get_async_extra_dejson() falls back to running extra_dejson in a worker thread when Connection.aextra_dejson() is missing. Its secret masking sends to the supervisor, and on Airflow 3.0 and 3.1 the supervisor channel has no thread lock: a send from a worker thread could interleave with an asend() in flight on the event loop, e.g. in the triggerer. On those versions the helper now deserializes Connection.extra without the supervisor, as async hooks did before (#55179, #72130). The worker thread stays for Airflow 3.2 to 3.3.1, whose channel is thread-safe, and Airflow 2, which has no supervisor. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Test get_async_extra_dejson against the real supervisor comms instead Revert the Airflow 3.0/3.1 json.loads branch: on those versions async code only runs in the triggerer, whose channel serializes every request, a worker thread's synchronous send included, through asend() and an asyncio.Lock, so the worker-thread fallback was safe and kept the masking. Instead, add a test that forks a real task process talking to the real supervisor of the installed Airflow and reads extras concurrently from an event loop, so each compat job checks the helper against that version's comms rather than mocks. It needs Airflow 3.2+, the first version where a task process can call the supervisor from async code. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> * Import AIRFLOW_V_3_2_PLUS from tests_common in the supervisor comms test Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
closes: #55728
What is the change?
The async paths in the Azure Data Factory and Synapse pipeline hooks stop reading
conn.extra_dejsonand parse extras withjson.loads(conn.extra)instead.AzureSynapsePipelineAsyncHook.get_async_connnow fetches the connection throughget_async_connectionrather than the syncget_connection, and ADFrefresh_conncloses the cached_async_connbefore rebuilding it.Why did I do it?
extra_dejsoncan callmask_secret, which does a sync send on the triggerer event loop and raisesRuntimeError: You cannot use AsyncToSync in the same thread as an async event loop. The stack users see atexecute_completeis that trigger error re-raised. This is the same leftover that #55179 fixed for MSGraph.How did I do it?
In
hooks/data_factory.pyandhooks/synapse.pyI swappedconn.extra_dejsonforjson.loads(conn.extra) if conn.extra else {}on the async paths only, and maderefresh_connawaitself.close(). The sync hooks keepextra_dejsonso worker-side secret masking is unchanged.What's the impact?
Deferrable ADF and Synapse pipeline runs stop crashing on the triggerer with the AsyncToSync error. Sync hooks and
execute_completebehave exactly as before. Azure provider only.What's the test plan?
New unit tests make
extra_dejsonraise the real AsyncToSync error and assert the async hooks still build their clients; they fail ifextra_dejsoncomes back.test_refresh_connnow asserts the aio client gets closed. Reviewers can re-run the two touched test files orbreeze testing providers-tests --test-type "Providers[microsoft.azure]"; provider CI covers the same.Was generative AI tooling used to co-author this PR?
Generated-by: Cursor Grok 4.6 following the guidelines