Managed Airflow (Gen 3) | Managed Airflow (Gen 2) | Managed Airflow (Legacy Gen 1)
Auf dieser Seite wird erläutert, wie Sie CeleryKubernetesExecutor in Managed Airflow aktivieren und KubernetesExecutor in Ihren DAGs verwenden.
Informationen zu CeleryKubernetesExecutor
CeleryKubernetesExecutor ist ein Executor-Typ, der CeleryExecutor und KubernetesExecutor gleichzeitig verwenden kann. Airflow wählt den Executor basierend auf der Warteschlange aus, die Sie für die Aufgabe definieren. In einem DAG können Sie einige Aufgaben mit CeleryExecutor und andere Aufgaben mit KubernetesExecutor ausführen:
- CeleryExecutor ist für die schnelle und skalierbare Ausführung von Aufgaben optimiert.
- KubernetesExecutor wurde für die Ausführung ressourcenintensiver Aufgaben und die Ausführung von Aufgaben in Isolation entwickelt.
CeleryKubernetesExecutor in Managed Airflow
Mit CeleryKubernetesExecutor in Managed Airflow können Sie KubernetesExecutor für Ihre Aufgaben verwenden. Es ist nicht möglich, KubernetesExecutor in Managed Airflow getrennt von CeleryKubernetesExecutor zu verwenden.
Managed Airflow führt Aufgaben, die Sie mit KubernetesExecutor ausführen, im Cluster Ihrer Umgebung im selben Namespace wie Airflow-Worker aus. Solche Aufgaben haben dieselben Bindungen wie Airflow Worker und können auf Ressourcen in Ihrem Projekt zugreifen.
Für Aufgaben, die Sie mit KubernetesExecutor ausführen, gilt das Managed Airflow-Preismodell, da Pods mit diesen Aufgaben im Cluster Ihrer Umgebung ausgeführt werden. Für diese Pods gelten die Managed Airflow-Compute-SKUs (für CPU, Arbeitsspeicher und Speicher).
Wir empfehlen, Aufgaben mit dem CeleryExecutor auszuführen, wenn:
- die Startzeit der Aufgabe wichtig ist.
- Aufgaben keine Laufzeitisolation erfordern und nicht ressourcenintensiv sind.
Wir empfehlen, Aufgaben mit dem KubernetesExecutor auszuführen, wenn:
- Aufgaben eine Laufzeitisolation erfordern. Beispielsweise damit Aufgaben nicht um Arbeitsspeicher und CPU konkurrieren, da sie in eigenen Pods ausgeführt werden.
- Aufgaben ressourcenintensiv sind und Sie die verfügbaren CPU- und Arbeitsspeicherressourcen steuern möchten.
KubernetesExecutor im Vergleich zu KubernetesPodOperator
Das Ausführen von Aufgaben mit KubernetesExecutor ähnelt dem Ausführen von Aufgaben mit KubernetesPodOperator. Aufgaben werden in Pods ausgeführt, wodurch eine Aufgabenisolation auf Pod-Ebene und eine bessere Ressourcenverwaltung ermöglicht werden.
Es gibt jedoch einige wichtige Unterschiede:
- KubernetesExecutor führt Aufgaben nur im versionierten Managed Airflow-Namespace Ihrer Umgebung aus. Dieser Namespace kann in Managed Airflow nicht geändert werden. Sie können einen Namespace angeben, in dem KubernetesPodOperator Pod-Aufgaben ausführt.
- KubernetesExecutor kann jeden integrierten Airflow-Operator verwenden. KubernetesPodOperator führt nur ein bereitgestelltes Skript aus, das durch den Einstiegspunkt des Containers definiert wird.
- KubernetesExecutor verwendet das standardmäßige Managed Airflow-Docker-Image mit denselben Python-Versionen, Überschreibungen von Airflow-Konfigurationsoptionen, Umgebungsvariablen und PyPI-Paketen, die in Ihrer Managed Airflow-Umgebung definiert sind.
Informationen zu Docker-Images
Standardmäßig startet KubernetesExecutor Aufgaben mit demselben Docker-Image, das Managed Airflow für Celery-Worker verwendet. Dies ist das Managed Airflow-Image für Ihre Umgebung mit allen Änderungen, die Sie für Ihre Umgebung angegeben haben, z. B. benutzerdefinierte PyPI Pakete oder Umgebungsvariablen.
.Hinweis
Sie können CeleryKubernetesExecutor in Managed Airflow (Gen 3) verwenden.
In Managed Airflow (Gen 3) kann kein anderer Executor als CeleryKubernetesExecutor verwendet werden. Das bedeutet, dass Sie Aufgaben mit CeleryExecutor, KubernetesExecutor oder beiden in einem DAG ausführen können. Es ist jedoch nicht möglich, Ihre Umgebung so zu konfigurieren, dass nur KubernetesExecutor oder CeleryExecutor verwendet wird.
CeleryKubernetesExecutor konfigurieren
Möglicherweise möchten Sie vorhandene Airflow-Konfigurationsoptionen überschreiben, die sich auf KubernetesExecutor beziehen:
[kubernetes]worker_pods_creation_batch_sizeMit dieser Option wird die Anzahl der Erstellungsaufrufe für Kubernetes-Worker-Pods pro Planer-Loop definiert. Der Standardwert ist
1, sodass nur ein einzelner Pod pro Planer-Heartbeat gestartet wird. Wenn Sie KubernetesExecutor häufig verwenden, empfehlen wir, diesen Wert zu erhöhen.[kubernetes]worker_pods_pending_timeoutMit dieser Option wird in Sekunden definiert, wie lange ein Worker im Status
Pending(Pod wird erstellt) bleiben kann, bevor er als fehlgeschlagen gilt. Der Standardwert ist 5 Minuten.
Aufgaben mit KubernetesExecutor oder CeleryExecutor ausführen
Sie können Aufgaben mit CeleryExecutor, KubernetesExecutor oder beiden in einem DAG ausführen:
Airflow 3
- Wenn Sie eine Aufgabe mit KubernetesExecutor ausführen möchten, geben Sie den Wert
KubernetesExecutorim Parameterexecutoreiner Aufgabe an. - Wenn Sie eine Aufgabe mit CeleryExecutor ausführen möchten, lassen Sie den Parameter
executorweg.
Airflow 2
- Wenn Sie eine Aufgabe mit KubernetesExecutor ausführen möchten, geben Sie den Wert
kubernetesim Parameterqueueeiner Aufgabe an. - Wenn Sie eine Aufgabe mit CeleryExecutor ausführen möchten, lassen Sie den Parameter
queueweg.
Im folgenden Beispiel wird die Aufgabe task-kubernetes mit KubernetesExecutor und die Aufgabe task-celery mit CeleryExecutor ausgeführt:
Airflow 3
import datetime
import airflow
from airflow.providers.standard.operators.python import PythonOperator
with airflow.DAG(
"composer_sample_celery_kubernetes",
start_date=datetime.datetime(2026, 1, 1),
schedule=None) as dag:
def kubernetes_example():
print("This task runs using KubernetesExecutor")
def celery_example():
print("This task runs using CeleryExecutor")
# To run with KubernetesExecutor, set queue to kubernetes
task_kubernetes = PythonOperator(
task_id='task-kubernetes',
python_callable=kubernetes_example,
dag=dag,
executor='KubernetesExecutor')
# To run with CeleryExecutor, omit the queue argument
task_celery = PythonOperator(
task_id='task-celery',
python_callable=celery_example,
dag=dag)
task_kubernetes >> task_celery
Airflow 2
import datetime
import airflow
from airflow.operators.python_operator import PythonOperator
with airflow.DAG(
"composer_sample_celery_kubernetes",
start_date=datetime.datetime(2026, 1, 1),
schedule=None) as dag:
def kubernetes_example():
print("This task runs using KubernetesExecutor")
def celery_example():
print("This task runs using CeleryExecutor")
# To run with KubernetesExecutor, set queue to kubernetes
task_kubernetes = PythonOperator(
task_id='task-kubernetes',
python_callable=kubernetes_example,
dag=dag,
queue='kubernetes')
# To run with CeleryExecutor, omit the queue argument
task_celery = PythonOperator(
task_id='task-celery',
python_callable=celery_example,
dag=dag)
task_kubernetes >> task_celery
Befehle der Airflow-Befehlszeile im Zusammenhang mit KubernetesExecutor ausführen
Sie können mehrere
Befehle der Airflow-Befehlszeile im Zusammenhang mit KubernetesExecutor
mit gcloud ausführen.
Spezifikation für Worker-Pod anpassen
Sie können die Spezifikation für Worker-Pods anpassen, indem Sie sie im Parameter executor_config einer Aufgabe übergeben. So können Sie benutzerdefinierte CPU- und Arbeitsspeicheranforderungen definieren.
Sie können die gesamte Spezifikation für Worker-Pods überschreiben, die zum Ausführen einer Aufgabe verwendet wird. Wenn Sie die Pod-Spezifikation einer Aufgabe abrufen möchten, die von KubernetesExecutor verwendet wird, können Sie den
Befehl kubernetes generate-dag-yaml der Airflow-Befehlszeile
ausführen.
Weitere Informationen zum Anpassen der Spezifikation für Worker-Pods finden Sie in der Airflow-Dokumentation.
Managed Airflow (Gen 3) unterstützt die folgenden Werte für Ressourcenanforderungen:
| Ressource | Minimum | Maximum | Schritt |
|---|---|---|---|
| CPU | 0,25 | 32 | Schrittwerte: 0,25, 0,5, 1, 2, 4, 6, 8, 10, ..., 32. Angefragte Werte werden auf den nächstgelegenen unterstützten Schrittwert aufgerundet (z. B. 5 auf 6). |
| Arbeitsspeicher | 2G (GB) | 128G (GB) | Schrittwerte: 2, 3, 4, 5, ..., 128. Angefragte Werte werden auf den nächstgelegenen unterstützten Schrittwert aufgerundet (z. B. 3, 5 G auf 4 G). |
| Speicher | - | 100G (GB) | Beliebiger Wert. Wenn mehr als 100 GB angefordert werden, werden nur 100 GB bereitgestellt. |
Weitere Informationen zu Ressourceneinheiten in Kubernetes finden Sie unter Ressourceneinheiten in Kubernetes.
Das folgende Beispiel zeigt eine Aufgabe, die eine benutzerdefinierte Spezifikation für Worker-Pods verwendet:
PythonOperator(
task_id='custom-spec-example',
python_callable=f,
dag=dag,
queue='kubernetes',
executor_config={
'pod_override': k8s.V1Pod(
spec=k8s.V1PodSpec(
containers=[
k8s.V1Container(
name='base',
resources=k8s.V1ResourceRequirements(requests={
'cpu': '0.5',
'memory': '2G',
})
),
],
),
)
},
)
Aufgabenlogs ansehen
Logs von Aufgaben, die von KubernetesExecutor ausgeführt werden, sind auf dem Tab Logs zusammen mit Logs von Aufgaben verfügbar, die von CeleryExecutor ausgeführt werden:
Rufen Sie in der Google Cloud Console die Seite Umgebungen auf.
Klicken Sie in der Liste der Umgebungen auf den Namen Ihrer Umgebung. Die Seite Umgebungsdetails wird geöffnet.
Wechseln Sie zum Tab Logs.
Gehen Sie zu Alle Logs > Airflow-Logs > Worker.
Worker mit dem Namen
airflow-k8s-workerführen KubernetesExecutor-Aufgaben aus. Wenn Sie nach Logs einer bestimmten Aufgabe suchen möchten, können Sie eine DAG-ID oder eine Aufgaben-ID als Keyword in der Suche verwenden.
Nächste Schritte
- Fehlerbehebung für KubernetesExecutor
- KubernetesPodOperator verwenden
- GKE-Operatoren verwenden
- Airflow-Konfigurationsoptionen überschreiben