viniciusdsmello opened a new issue, #74373:
URL: https://github.com/apache/airflow/issues/74373
### Under which category would you file this issue?
Providers
### Apache Airflow version
3.3.1
### What happened and how to reproduce it?
`KubernetesStartKueueJobOperator` creates the Job suspended with the Kueue
queue label, then inherits the waiting of `KubernetesJobOperator`. Pod
discovery is already covered by #70585 / #70636. The completion wait has two
further problems under Kueue. Line numbers refer to cncf.kubernetes 10.23.0
(the latest release); the code is the same in 10.21.0, which we run.
**1. Completion fails with 403 when the Role grants `get` on `jobs` but not
on `jobs/status`.**
`KubernetesHook.get_job_status` (`hooks/kubernetes.py` 623-631) and its
async version (1315-1324, used by the deferrable trigger) call
`read_namespaced_job_status`, which is the `jobs/status` subresource.
Kubernetes RBAC checks subresources separately, so a Role that grants `get` on
`batch/jobs` still gets `Forbidden` here, although `read_namespaced_job`
returns the same object including `status`. Least-privilege Roles for Airflow
on shared Kueue clusters commonly grant only `jobs`.
Reproduction (namespace `team-a`, Airflow connection `team_a` with a
kubeconfig for ServiceAccount `airflow` in that namespace):
```yaml
apiVersion: rbac.authorization.k8s.io/v1
kind: Role
metadata: {name: airflow-runner, namespace: team-a}
rules:
- {apiGroups: [""], resources: [pods], verbs: [create, get, list, watch,
delete, patch]}
- {apiGroups: [""], resources: [pods/log], verbs: [get]}
- {apiGroups: [batch], resources: [jobs], verbs: [create, get, list,
watch, delete]}
```
```bash
# Same identity, same Job: the object read works, the status subresource
does not.
kubectl -n team-a get --raw /apis/batch/v1/namespaces/team-a/jobs/<job>
# 200
kubectl -n team-a get --raw
/apis/batch/v1/namespaces/team-a/jobs/<job>/status # Forbidden
```
Any `KubernetesJobOperator` / `KubernetesStartKueueJobOperator` task with
`wait_until_job_complete=True` on that connection then fails in
`wait_until_job_complete` with `ApiException(403)` once the pods are found. We
observed the 403 on a live cluster (Kubernetes 1.35) with exactly this Role.
**2. A Job that Kueue preempts is polled until the task times out, and a
readmitted Job fails on log fetching.**
When Kueue preempts an admitted Job (for example a higher-priority workload
in a ClusterQueue with `preemption.withinClusterQueue: LowerPriority`), it sets
`spec.suspend: true` again and the Job controller deletes the Job's pods. Kueue
later readmits the Job and new pods start.
- `KubernetesHook.wait_until_job_complete` (`hooks/kubernetes.py` 635-649)
loops until `is_job_complete` sees a Complete or Failed condition. A suspended
Job has neither, so the task polls until `execution_timeout`, even if the Job
is never readmitted. The deferrable trigger (`triggers/job.py` 141, async hook
1330-1345) has the same loop.
- If the Job is readmitted and completes, `execute()` (`operators/job.py`
224-253) fetches logs from `self.pods`, the pods found before the preemption.
They no longer exist, so the log read returns 404 and the task fails although
the Job succeeded.
This part is from reading the code paths above and from unit tests against a
scripted cluster. We have not yet triggered a live preemption against the stock
operator. Expected reproduction:
```yaml
apiVersion: kueue.x-k8s.io/v1beta1
kind: WorkloadPriorityClass
metadata: {name: low}
value: 1000
---
apiVersion: kueue.x-k8s.io/v1beta1
kind: WorkloadPriorityClass
metadata: {name: high}
value: 10000
---
apiVersion: kueue.x-k8s.io/v1beta1
kind: ClusterQueue
metadata: {name: cq}
spec:
namespaceSelector: {}
preemption: {withinClusterQueue: LowerPriority}
resourceGroups:
- coveredResources: [cpu, memory]
flavors:
- name: default-flavor
resources:
- {name: cpu, nominalQuota: 1}
- {name: memory, nominalQuota: 1Gi}
---
apiVersion: kueue.x-k8s.io/v1beta1
kind: LocalQueue
metadata: {name: lq, namespace: team-a}
spec: {clusterQueue: cq}
```
```python
from datetime import datetime
from airflow.providers.cncf.kubernetes.operators.kueue import
KubernetesStartKueueJobOperator
from airflow.sdk import DAG
from kubernetes.client import models as k8s
with DAG("kueue_preemption_repro", schedule=None, start_date=datetime(2026,
1, 1)):
KubernetesStartKueueJobOperator(
task_id="low_priority_job",
queue_name="lq",
namespace="team-a",
image="busybox",
cmds=["sh", "-c", "sleep 300; echo done"],
labels={"kueue.x-k8s.io/priority-class": "low"},
container_resources=k8s.V1ResourceRequirements(
requests={"cpu": "1", "memory": "512Mi"}, limits={"cpu": "1",
"memory": "512Mi"}
),
wait_until_job_complete=True,
get_logs=True,
)
```
1. Trigger the DAG and wait until the Job's pod is running.
2. Create a Job in `team-a` with labels `kueue.x-k8s.io/queue-name: lq` and
`kueue.x-k8s.io/priority-class: high`, requesting 1 CPU, running `sleep 60`.
Kueue preempts the first Job: it is suspended and its pod deleted.
3. The Airflow task keeps logging "The job ... is incomplete. Sleeping for
10 sec." while the first Job is suspended.
4. After the high-priority Job finishes, Kueue readmits the first Job, which
completes. The task then fails fetching logs from the deleted pod.
### What you think should happen instead?
1. Read the Job with `read_namespaced_job` in `get_job_status` (sync and
async). It returns the same `status` and needs only `get` on `jobs`.
Alternatively, fall back to it on a 403.
2. In the completion wait, treat "Job suspended with no active pods" after
pod discovery as a preemption: log it, then go back to pod discovery (bounded
by the same timeout #70636 introduces) instead of polling for a final
condition. Fetch logs from the Job's pods as listed at the end, for example by
the `batch.kubernetes.io/job-name` label, not from the pods discovered before
the preemption.
### Operating System
Not Applicable (the operator runs in the official Airflow image)
### Deployment
Official Apache Airflow Helm Chart
### Apache Airflow Provider(s)
cncf-kubernetes
### Versions of Apache Airflow Providers
apache-airflow-providers-cncf-kubernetes==10.21.0 (also checked against
10.23.0, unchanged code)
### Official Helm Chart version
Not Applicable
### Kubernetes Version
1.35
### Helm Chart configuration
Not Applicable
### Docker Image customizations
Not Applicable (no change to the provider or the Kubernetes client)
### Anything else?
- Happens every time a Kueue Job is preempted (problem 2), and on every
completion wait with a Role that lacks `jobs/status` (problem 1).
- Related: #70585 (issue) and #70636 (PR) for the pod discovery wait, and
#44568, which added the Kueue operators.
- We run a subclass of `KubernetesStartKueueJobOperator` that handles both
points plus the discovery wait, with unit tests that drive `execute()` through
queued, running, preempted, readmitted and failed states. We can adapt it into
a PR if the approach is acceptable.
### Are you willing to submit PR?
- [x] Yes I am willing to submit a PR!
### Code of Conduct
- [x] I agree to follow this project's [Code of
Conduct](https://github.com/apache/airflow/blob/main/CODE_OF_CONDUCT.md)
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]