Conversation
a3d8b31 to
e6838f1
Compare
Miretpl
left a comment
There was a problem hiding this comment.
Could you adjust your PR description to match our project guidelines - https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst?
|
I would appreciate a review from a Kubernetes Google maintainer on that part of the code. |
e6838f1 to
0bf66bf
Compare
Besides the AI disclosure block what's missing or not part of the guideline? |
0bf66bf to
4ec850e
Compare
4ec850e to
a708b4b
Compare
Only this was missing here. Thanks for adding it. I've rerun the whole CI as there were a couple of unrelated failures. Let's see how it will be now. |
a708b4b to
2f15574
Compare
2f15574 to
73adafa
Compare
|
Since the ApiClient is now reused for the lifetime of the trigger, can we explicitly document how credential refresh is handled independently of the client lifecycle? In particular, I'm wondering about long-running triggers where credentials may expire while the client is still alive. |
|
What happens if client initialization succeeds but a later poll fails due to an authentication/configuration error? Should that error cause the cached client to be invalidated so that the next poll can rebuild it, or is the expectation that the trigger will terminate/retry? It would be good to define this now that the client is no longer recreated on every poll. |
Credential refresh is handled separately from the cached client:
|
The cached client isn’t invalidated when a poll fails. Existing retries still handle transient errors, but authentication errors such as 401/403 cause the pod trigger to report an error, and cleanup closes the client. Whether the task retries afterward depends on its retry settings. |
73adafa to
f64718a
Compare
We couldn’t guarantee that before. The regression test I added showed that retries could reuse the rejected token, so I also added a fix. The GKE client now refreshes credentials and retries once after a 401, using the same client. If the refresh or retry fails, the error still propagates. |
d304e6f to
13c3bd3
Compare
aaron-y-chen
left a comment
There was a problem hiding this comment.
Overall, the Google provider parts LGTM. However, it now relies on _cached_kube_client and close(), which are added to AsyncKubernetesHook in this PR. Without marking this cross-provider dependency, the Google provider could still be installed with an older Kubernetes provider and fail with an AttributeError on the first GKE trigger poll.
Should we add # use next version at here?
airflow/providers/google/pyproject.toml
Line 166 in b4d0129
I added # use next version 👍 |
aaron-y-chen
left a comment
There was a problem hiding this comment.
Kubernetes Google part LGTM.
bc15bfd to
ca71fc5
Compare
Miretpl
left a comment
There was a problem hiding this comment.
Looking at the changes in the code, I'm reverting my approve to prevent accidental merge of it. I would recommend separate this PR to cncf and google provider - it should be easier to merge at least one of them.
| try: | ||
| return await super().call_api(*args, **kwargs) | ||
| except async_client.ApiException as error: | ||
| if error.status != 401: |
There was a problem hiding this comment.
We can receive a 401 code for more than just an expired token, so I don't think that this handling is correct.
There was a problem hiding this comment.
How about following google-auth’s refresh-on-401 approach, capped at one attempt? I’d retry only if the token changes and propagate any refresh or retry failure. This preserves the refresh attempt requested by @soupam05
ca71fc5 to
cc2b9b8
Compare
@Miretpl I opened #73399 for the cncf changes. The Google changes depend on them, so I’ll keep them here until that PR merges, then rebase this PR to leave only the Google changes. |
cc2b9b8 to
d8ac3f1
Compare
d8ac3f1 to
c53b2a1
Compare
…ests for token expiry scenarios
VladaZakharova
left a comment
There was a problem hiding this comment.
get_conn() previously closed its ApiClient when the async context exited. For cached configurations it now requires a separate await hook.close(), so existing direct hook users and trigger subclasses can silently leak aiohttp sessions. Can client reuse be opt-in for the built-in triggers, preserving the existing context-manager lifecycle by default? If the lifecycle change is intentional, it should at least be documented as a provider behavioral change.
@VladaZakharova I'll address this as soon as #73399 is merged |
I've been chasing triggerer issues on our production deployment: sustained
triggers.blocked_main_threadbursts whenever a wave of deferrableKubernetesPodOperator/GKEStartPodOperatortasks started or finished, with the async thread blocked 2–10s at a time. #69661 helped the log-parsing path, but the metric barely moved. After reading the source and profiling the TriggerRunner with py-spy during blocking windows, the loop thread showed up insidessl.create_default_context()reached fromAsyncKubernetesHook.get_conn(): every API method builds a freshApiClientper call, which parses the CA bundle synchronously on the event loop and opens a new aiohttp session, a full TCP+TLS handshake per status poll (default 2s per trigger) with zero reuse. The GKE hook additionally fetched a new OAuth token and wrote the cluster CA to a temp file each call.The change itself is simple: instead of building a new API client for every call, the hook now builds one and keeps it for its whole lifetime. For in-cluster and static-auth kubeconfigs Airflow loads the configuration once per hook, and anything that does rotate (the in-cluster service-account token) is read from the configuration per request rather than baked into the client, so one client can serve the hook's lifetime. Exec-based auth (EKS/GKE kubeconfigs) is the exception and keeps the old per-call behavior on purpose, reloading the config on every call is exactly what refreshes its short-lived credentials, so those clients can't be reused. To pair every cached client with an owner, hooks get an idempotent
close()and the triggers call it fromcleanup(), so the client lives exactly as long as its trigger. The GKE hook gets the same treatment plus token caching: it re-sets the bearer header every time the connection is used, so tokens refresh on schedule instead of being fetched from Google on every call (and since all that was left of its_load_config()was building the client, it's now called_build_client()).A few things follow from the caching that are worth knowing. If you instantiate these hooks yourself or subclass the triggers and override
cleanup(), callhook.close()when you're done, otherwise the client is only cleaned up at garbage collection and aiohttp logs an "unclosed session" warning; calling the hook again afterclose()just builds a fresh client. mTLS client certs from static kubeconfigs are now read once per hook rather than on every call, which I think is fine for hooks that live only as long as a trigger (if cert rotation mid-trigger ever becomes a real problem, rebuilding the client on auth errors would be the fix). The pre-existing gap left by #65212 the default-kubeconfig path sets_config_loadedwithout exec-auth detection, is untouched. And one expectation to set: each new trigger still builds its one client on the event loop, soblocked_main_threadshould drop dramatically but won't be exactly zero when a batch of triggers starts.To make sure this was really the problem, I reproduced it locally on a minikube cluster: real pod triggers polling real pods, with Airflow's own
block_watchdogrunning, counting how many SSL contexts and API clients got created. I ran the exact same load twice, once with the old per-callget_connpatched back in as a baseline, once with this change. The baseline created a new client (and SSL context) for every single API call; with this change each trigger creates exactly one and releases it on cleanup. Unit tests in both providers assert the same behavior.related: #69661
Was generative AI tooling used to co-author this PR?
No