From bcf4fbb6b94c184a96eebef837356902b762d87e Mon Sep 17 00:00:00 2001 From: Andrew Kenworthy Date: Wed, 30 Sep 2026 11:09:07 +0200 Subject: [PATCH 1/2] fix: Airflow demo improvements --- .../01-airflow-demo-clusterrole.yaml | 42 ----- .../01-airflow-demo-role.yaml | 45 ++++++ .../02-airflow-demo-clusterrolebinding.yaml | 13 -- .../02-airflow-demo-rolebinding.yaml | 26 ++++ .../03-enable-and-run-spark-dag.yaml | 1 + .../04-enable-and-run-date-dag.yaml | 1 + .../05-enable-and-run-kafka-dag.yaml | 1 + .../06-create-opa-users.yaml | 1 + .../airflow-scheduled-job/serviceaccount.yaml | 45 +----- stacks/airflow/airflow.yaml | 143 ++++++++++++++---- stacks/airflow/minio.yaml | 6 +- 11 files changed, 194 insertions(+), 130 deletions(-) delete mode 100644 demos/airflow-scheduled-job/01-airflow-demo-clusterrole.yaml create mode 100644 demos/airflow-scheduled-job/01-airflow-demo-role.yaml delete mode 100644 demos/airflow-scheduled-job/02-airflow-demo-clusterrolebinding.yaml create mode 100644 demos/airflow-scheduled-job/02-airflow-demo-rolebinding.yaml diff --git a/demos/airflow-scheduled-job/01-airflow-demo-clusterrole.yaml b/demos/airflow-scheduled-job/01-airflow-demo-clusterrole.yaml deleted file mode 100644 index bf4c978e..00000000 --- a/demos/airflow-scheduled-job/01-airflow-demo-clusterrole.yaml +++ /dev/null @@ -1,42 +0,0 @@ ---- -apiVersion: rbac.authorization.k8s.io/v1 -kind: ClusterRole -metadata: - name: airflow-demo-clusterrole -rules: -- apiGroups: - - spark.stackable.tech - resources: - - sparkapplications - verbs: - - create - - get - - list -- apiGroups: - - apps - resources: - - statefulsets - verbs: - - get - - watch - - list -- apiGroups: - - "" - resources: - - persistentvolumeclaims - verbs: - - list -- apiGroups: - - "" - resources: - - pods - verbs: - - get - - watch - - list -- apiGroups: - - "" - resources: - - pods/exec - verbs: - - create diff --git a/demos/airflow-scheduled-job/01-airflow-demo-role.yaml b/demos/airflow-scheduled-job/01-airflow-demo-role.yaml new file mode 100644 index 00000000..63f44bea --- /dev/null +++ b/demos/airflow-scheduled-job/01-airflow-demo-role.yaml @@ -0,0 +1,45 @@ +--- +# Used by the Airflow cluster itself: the sparkapp_dag DAG creates a SparkApplication +# (SparkKubernetesOperator) and polls its status (SparkKubernetesSensor). Reading the Spark +# driver/executor pods and their logs is already covered by the RoleBinding the Airflow operator +# creates for airflow-serviceaccount. +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: airflow-demo-spark +rules: +- apiGroups: + - spark.stackable.tech + resources: + - sparkapplications + verbs: + - create + - get +--- +# Used by the setup Jobs (start-*-job): "kubectl rollout status" on the Airflow and Kafka +# StatefulSets and "kubectl exec" into kafka-broker-default-0 and airflow-webserver-default-0. +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: airflow-demo-setup +rules: +- apiGroups: + - apps + resources: + - statefulsets + verbs: + - get + - list + - watch +- apiGroups: + - "" + resources: + - pods + verbs: + - get +- apiGroups: + - "" + resources: + - pods/exec + verbs: + - create diff --git a/demos/airflow-scheduled-job/02-airflow-demo-clusterrolebinding.yaml b/demos/airflow-scheduled-job/02-airflow-demo-clusterrolebinding.yaml deleted file mode 100644 index 1eccf63e..00000000 --- a/demos/airflow-scheduled-job/02-airflow-demo-clusterrolebinding.yaml +++ /dev/null @@ -1,13 +0,0 @@ ---- -apiVersion: rbac.authorization.k8s.io/v1 -kind: ClusterRoleBinding -metadata: - name: airflow-demo-clusterrole-binding -roleRef: - apiGroup: rbac.authorization.k8s.io - kind: ClusterRole - name: airflow-demo-clusterrole -subjects: -- apiGroup: rbac.authorization.k8s.io - kind: Group - name: system:serviceaccounts diff --git a/demos/airflow-scheduled-job/02-airflow-demo-rolebinding.yaml b/demos/airflow-scheduled-job/02-airflow-demo-rolebinding.yaml new file mode 100644 index 00000000..edd01682 --- /dev/null +++ b/demos/airflow-scheduled-job/02-airflow-demo-rolebinding.yaml @@ -0,0 +1,26 @@ +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: airflow-demo-spark +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: airflow-demo-spark +subjects: +- kind: ServiceAccount + name: airflow-serviceaccount + namespace: {{ NAMESPACE }} +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: airflow-demo-setup +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: airflow-demo-setup +subjects: +- kind: ServiceAccount + name: demo-serviceaccount + namespace: {{ NAMESPACE }} diff --git a/demos/airflow-scheduled-job/03-enable-and-run-spark-dag.yaml b/demos/airflow-scheduled-job/03-enable-and-run-spark-dag.yaml index f7b52582..1afb2357 100644 --- a/demos/airflow-scheduled-job/03-enable-and-run-spark-dag.yaml +++ b/demos/airflow-scheduled-job/03-enable-and-run-spark-dag.yaml @@ -6,6 +6,7 @@ metadata: spec: template: spec: + serviceAccountName: demo-serviceaccount containers: - name: start-pyspark-job image: oci.stackable.tech/sdp/tools:1.0.0-stackable0.0.0-dev diff --git a/demos/airflow-scheduled-job/04-enable-and-run-date-dag.yaml b/demos/airflow-scheduled-job/04-enable-and-run-date-dag.yaml index 94165cb6..c1faaab9 100644 --- a/demos/airflow-scheduled-job/04-enable-and-run-date-dag.yaml +++ b/demos/airflow-scheduled-job/04-enable-and-run-date-dag.yaml @@ -6,6 +6,7 @@ metadata: spec: template: spec: + serviceAccountName: demo-serviceaccount containers: - name: start-date-job image: oci.stackable.tech/sdp/tools:1.0.0-stackable0.0.0-dev diff --git a/demos/airflow-scheduled-job/05-enable-and-run-kafka-dag.yaml b/demos/airflow-scheduled-job/05-enable-and-run-kafka-dag.yaml index 84a26c6d..c2efcd97 100644 --- a/demos/airflow-scheduled-job/05-enable-and-run-kafka-dag.yaml +++ b/demos/airflow-scheduled-job/05-enable-and-run-kafka-dag.yaml @@ -6,6 +6,7 @@ metadata: spec: template: spec: + serviceAccountName: demo-serviceaccount containers: - name: start-kafka-job image: oci.stackable.tech/sdp/tools:1.0.0-stackable0.0.0-dev diff --git a/demos/airflow-scheduled-job/06-create-opa-users.yaml b/demos/airflow-scheduled-job/06-create-opa-users.yaml index 3a77e437..85b78c29 100644 --- a/demos/airflow-scheduled-job/06-create-opa-users.yaml +++ b/demos/airflow-scheduled-job/06-create-opa-users.yaml @@ -6,6 +6,7 @@ metadata: spec: template: spec: + serviceAccountName: demo-serviceaccount containers: - name: start-users-job image: oci.stackable.tech/sdp/tools:1.0.0-stackable0.0.0-dev diff --git a/demos/airflow-scheduled-job/serviceaccount.yaml b/demos/airflow-scheduled-job/serviceaccount.yaml index e42b2e26..f4ed842c 100644 --- a/demos/airflow-scheduled-job/serviceaccount.yaml +++ b/demos/airflow-scheduled-job/serviceaccount.yaml @@ -1,48 +1,7 @@ --- +# Runs the setup Jobs. Its permissions are granted by the airflow-demo-setup Role +# (01-airflow-demo-role.yaml). apiVersion: v1 kind: ServiceAccount metadata: name: demo-serviceaccount ---- -apiVersion: rbac.authorization.k8s.io/v1 -kind: ClusterRoleBinding -metadata: - name: demo-clusterrolebinding -subjects: - - kind: ServiceAccount - name: demo-serviceaccount - namespace: {{ NAMESPACE }} -roleRef: - kind: ClusterRole - name: demo-clusterrole - apiGroup: rbac.authorization.k8s.io ---- -apiVersion: rbac.authorization.k8s.io/v1 -kind: ClusterRole -metadata: - name: demo-clusterrole -rules: - - apiGroups: - - "" - resources: - - pods - verbs: - - get - - list - - watch - - apiGroups: - - apps - resources: - - statefulsets - verbs: - - get - - list - - watch - - apiGroups: - - batch - resources: - - jobs - verbs: - - get - - list - - watch diff --git a/stacks/airflow/airflow.yaml b/stacks/airflow/airflow.yaml index 71ee799f..0482577a 100644 --- a/stacks/airflow/airflow.yaml +++ b/stacks/airflow/airflow.yaml @@ -79,7 +79,7 @@ spec: spec: containers: - name: airflow - image: oci.stackable.tech/sdp/airflow:3.1.6-stackable0.0.0-dev + image: oci.stackable.tech/sdp/airflow:3.3.1-stackable0.0.0-dev imagePullPolicy: IfNotPresent env: - name: NAMESPACE @@ -107,7 +107,7 @@ spec: spec: containers: - name: base - image: oci.stackable.tech/sdp/airflow:3.1.6-stackable0.0.0-dev + image: oci.stackable.tech/sdp/airflow:3.3.1-stackable0.0.0-dev imagePullPolicy: IfNotPresent env: # Needed by the DAGs run by KubernetesPodOperator @@ -386,6 +386,8 @@ data: import yaml from airflow.utils import yaml import os + import threading + import time if TYPE_CHECKING: from airflow.utils.context import Context @@ -430,8 +432,14 @@ data: template_fields = ("application_name", "namespace") # See https://github.com/stackabletech/spark-k8s-operator/pull/460/files#diff-d737837121132af6b60f50279a78464b05dcfd06c05d1d090f4198a5e962b5f6R371 # Unknown is set immediately so it must be excluded from the failed states. - FAILURE_STATES = ("Failed") - SUCCESS_STATES = ("Succeeded") + FAILURE_STATES = ("Failed",) + SUCCESS_STATES = ("Succeeded",) + + # The Stackable Spark operator sets spark.kubernetes.{driver,executor}.podTemplateContainerName, + # so the main container is called "spark" in both driver and executor pods. + SPARK_CONTAINER = "spark" + # How long the terminal poke waits for the log followers to drain. + LOG_DRAIN_TIMEOUT_SECONDS = 30 def __init__( self, @@ -454,32 +462,103 @@ data: self.api_group = api_group self.api_version = api_version self.poke_interval = poke_interval - - def _log_driver(self, application_state: str, response: dict) -> None: - if not self.attach_log: - return - status_info = response["status"] - if "driverInfo" not in status_info: - return - driver_info = status_info["driverInfo"] - if "podName" not in driver_info: - return - driver_pod_name = driver_info["podName"] - namespace = response["metadata"]["namespace"] - log_method = self.log.error if application_state in self.FAILURE_STATES else self.log.info + # One thread per Spark pod, keyed by pod name. This relies on mode="poke" + # (the default): with mode="reschedule" every poke gets a new instance. + self._log_followers: Dict[str, threading.Thread] = {} + + def _start_log_followers(self, namespace: str) -> None: + # The Stackable Spark operator's SparkApplicationStatus only exposes + # `phase` and `resolvedTemplateRef` - driverInfo is not populated - so + # find the driver and executor pods by the labels Spark sets on them. + label_selector = ( + f"app.kubernetes.io/instance={self.application_name}," + "spark-role in (driver,executor)" + ) try: - log = "" - for line in self.hook.get_pod_logs(driver_pod_name, namespace=namespace): - log += line.decode() - log_method(log) + pods = self.hook.core_v1_client.list_namespaced_pod( + namespace=namespace, label_selector=label_selector + ) except client.rest.ApiException as e: - self.log.warning( - "Could not read logs for pod %s. It may have been disposed.\n" - "Make sure timeToLiveSeconds is set on your SparkApplication spec.\n" - "underlying exception: %s", - driver_pod_name, - e, + self.log.warning("Could not list Spark pods for %s: %s", self.application_name, e) + return + for pod in pods.items: + pod_name = pod.metadata.name + if pod_name in self._log_followers: + continue + role = (pod.metadata.labels or {}).get("spark-role") + prefix = "[spark driver]" if role == "driver" else f"[spark executor {pod_name}]" + self.log.info("Following logs of %s pod %s", role, pod_name) + follower = threading.Thread( + target=self._follow_pod_log, + args=(namespace, pod_name, prefix), + name=f"log-{pod_name}", + daemon=True, ) + self._log_followers[pod_name] = follower + follower.start() + + def _container_finished(self, namespace: str, pod_name: str) -> bool: + try: + pod = self.hook.core_v1_client.read_namespaced_pod(name=pod_name, namespace=namespace) + except client.rest.ApiException as e: + return e.status == 404 + for status in pod.status.container_statuses or []: + if status.name == self.SPARK_CONTAINER: + return status.state.terminated is not None + return False + + def _follow_pod_log(self, namespace: str, pod_name: str, prefix: str) -> None: + # Streams the log with follow=True. The stream only ends when the + # container exits, so the last lines are read before the operator + # deletes the pod. If the stream breaks early it is reopened; the log + # then starts from the beginning again, so already emitted lines are skipped. + emitted = 0 + while True: + try: + # _preload_content=False returns the raw HTTP response. Without it, + # newer kubernetes clients (seen with 36.0.3) return str(bytes), i.e. + # the whole log as a single "b'...'" string. + resp = self.hook.core_v1_client.read_namespaced_pod_log( + name=pod_name, + namespace=namespace, + container=self.SPARK_CONTAINER, + follow=True, + _preload_content=False, + ) + except client.rest.ApiException as e: + if e.status == 404: + return + # 400 while the container is still being created. + time.sleep(1) + continue + seen = 0 + buffer = b"" + try: + for chunk in resp.stream(decode_content=True): + buffer += chunk + *lines, buffer = buffer.split(b"\n") + for line in lines: + seen += 1 + if seen > emitted: + self.log.info("%s %s", prefix, line.decode(errors="replace")) + emitted = seen + if buffer and seen + 1 > emitted: + self.log.info("%s %s", prefix, buffer.decode(errors="replace")) + emitted = seen + 1 + except Exception as e: + self.log.debug("Log stream for pod %s interrupted: %s", pod_name, e) + finally: + resp.release_conn() + if self._container_finished(namespace, pod_name): + return + time.sleep(1) + + def _wait_for_log_followers(self) -> None: + deadline = time.monotonic() + self.LOG_DRAIN_TIMEOUT_SECONDS + for pod_name, follower in self._log_followers.items(): + follower.join(timeout=max(0.0, deadline - time.monotonic())) + if follower.is_alive(): + self.log.warning("Log follower for pod %s did not finish in time", pod_name) def poke(self, context: Dict) -> bool: self.log.info("Poking: %s", self.application_name) @@ -495,8 +574,13 @@ data: except KeyError: self.log.debug(f"SparkApplication status could not be established: {response}") return False - if self.attach_log and application_state in self.FAILURE_STATES + self.SUCCESS_STATES: - self._log_driver(application_state, response) + + if self.attach_log: + if application_state in self.FAILURE_STATES + self.SUCCESS_STATES: + self._wait_for_log_followers() + else: + self._start_log_followers(response["metadata"]["namespace"]) + if application_state in self.FAILURE_STATES: raise AirflowException(f"SparkApplication failed with state: {application_state}") elif application_state in self.SUCCESS_STATES: @@ -547,6 +631,7 @@ data: namespace=ns, application_name="{{ task_instance.xcom_pull(task_ids='spark_pi_submit')['metadata']['name'] }}", poke_interval=5, + attach_log=True, dag=dag, ) diff --git a/stacks/airflow/minio.yaml b/stacks/airflow/minio.yaml index 34089aac..b34cfe3a 100644 --- a/stacks/airflow/minio.yaml +++ b/stacks/airflow/minio.yaml @@ -523,7 +523,7 @@ spec: serviceAccountName: minio-sa containers: - name: minio - image: "quay.io/minio/minio:RELEASE.2024-12-18T13-15-44Z" + image: "docker.io/pgsty/minio:RELEASE.2026-08-04T00-00-00Z" imagePullPolicy: IfNotPresent command: - "/bin/sh" @@ -652,7 +652,7 @@ spec: serviceAccountName: minio-sa containers: - name: minio-make-bucket - image: "quay.io/minio/mc:RELEASE.2024-11-21T17-21-54Z" + image: "docker.io/pgsty/mc:RELEASE.2026-09-16T00-00-00Z" imagePullPolicy: IfNotPresent command: - "/bin/sh" @@ -683,7 +683,7 @@ spec: requests: memory: 128Mi - name: minio-make-user - image: "quay.io/minio/mc:RELEASE.2024-11-21T17-21-54Z" + image: "docker.io/pgsty/mc:RELEASE.2026-09-16T00-00-00Z" imagePullPolicy: IfNotPresent command: - "/bin/sh" From 13a145c391b0515a251a2e298b7c6932cd0ef963 Mon Sep 17 00:00:00 2001 From: Andrew Kenworthy Date: Wed, 30 Sep 2026 16:12:17 +0200 Subject: [PATCH 2/2] place service account first in the order of manifests --- demos/demos-v2.yaml | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/demos/demos-v2.yaml b/demos/demos-v2.yaml index 5360da21..8bc03804 100644 --- a/demos/demos-v2.yaml +++ b/demos/demos-v2.yaml @@ -46,13 +46,13 @@ demos: - airflow - job-scheduling manifests: - - plainYaml: https://raw.githubusercontent.com/stackabletech/demos/main/demos/airflow-scheduled-job/01-airflow-demo-clusterrole.yaml - - plainYaml: https://raw.githubusercontent.com/stackabletech/demos/main/demos/airflow-scheduled-job/02-airflow-demo-clusterrolebinding.yaml + - plainYaml: https://raw.githubusercontent.com/stackabletech/demos/main/demos/airflow-scheduled-job/serviceaccount.yaml + - plainYaml: https://raw.githubusercontent.com/stackabletech/demos/main/demos/airflow-scheduled-job/01-airflow-demo-role.yaml + - plainYaml: https://raw.githubusercontent.com/stackabletech/demos/main/demos/airflow-scheduled-job/02-airflow-demo-rolebinding.yaml - plainYaml: https://raw.githubusercontent.com/stackabletech/demos/main/demos/airflow-scheduled-job/03-enable-and-run-spark-dag.yaml - plainYaml: https://raw.githubusercontent.com/stackabletech/demos/main/demos/airflow-scheduled-job/04-enable-and-run-date-dag.yaml - plainYaml: https://raw.githubusercontent.com/stackabletech/demos/main/demos/airflow-scheduled-job/05-enable-and-run-kafka-dag.yaml - plainYaml: https://raw.githubusercontent.com/stackabletech/demos/main/demos/airflow-scheduled-job/06-create-opa-users.yaml - - plainYaml: https://raw.githubusercontent.com/stackabletech/demos/main/demos/airflow-scheduled-job/serviceaccount.yaml - plainYaml: https://raw.githubusercontent.com/stackabletech/demos/main/demos/airflow-scheduled-job/create-trino-tables.yaml supportedNamespaces: [] resourceRequests: