Managed Airflow (3e génération) | Managed Airflow (2e génération) | Managed Airflow (1re génération héritée)
Cette page explique comment activer CeleryKubernetesExecutor dans Managed Airflow et comment utiliser KubernetesExecutor dans vos DAG.
À propos de CeleryKubernetesExecutor
CeleryKubernetesExecutor est un type d'exécuteur qui peut utiliser CeleryExecutor et KubernetesExecutor en même temps. Airflow sélectionne l'exécuteur en fonction de la file d'attente que vous définissez pour la tâche. Dans un DAG, vous pouvez exécuter certaines tâches avec CeleryExecutor et d'autres avec KubernetesExecutor :
- CeleryExecutor est optimisé pour une exécution rapide et évolutive des tâches.
- KubernetesExecutor est conçu pour l'exécution de tâches gourmandes en ressources et l'exécution de tâches de manière isolée.
CeleryKubernetesExecutor dans Managed Airflow
CeleryKubernetesExecutor dans Managed Airflow vous permet d'utiliser KubernetesExecutor pour vos tâches. Il n'est pas possible d'utiliser KubernetesExecutor dans Managed Airflow séparément de CeleryKubernetesExecutor.
Managed Airflow exécute les tâches que vous exécutez avec KubernetesExecutor dans le cluster de votre environnement, dans le même espace de noms que les nœuds de calcul Airflow. Ces tâches ont les mêmes liaisons que les nœuds de calcul Airflow et peuvent accéder aux ressources de votre projet.
Les tâches que vous exécutez avec KubernetesExecutor utilisent le modèle de tarification Managed Airflow, car les pods contenant ces tâches s'exécutent dans le cluster de votre environnement. Les références SKU de calcul Managed Airflow (pour le processeur, la mémoire et le stockage) s'appliquent à ces pods.
Nous vous recommandons d'exécuter les tâches avec CeleryExecutor dans les cas suivants :
- Le temps de démarrage des tâches est important.
- Les tâches ne nécessitent pas d'isolation au moment de l'exécution et ne sont pas gourmandes en ressources.
Nous vous recommandons d'exécuter les tâches avec KubernetesExecutor dans les cas suivants :
- Les tâches nécessitent une isolation au moment de l'exécution. Par exemple, pour que les tâches ne soient pas en concurrence pour la mémoire et le processeur, car elles s'exécutent dans leurs propres pods.
- Les tâches sont gourmandes en ressources et vous souhaitez contrôler les ressources de processeur et de mémoire disponibles.
Comparaison entre KubernetesExecutor et KubernetesPodOperator
L'exécution de tâches avec KubernetesExecutor est semblable à l'exécution de tâches à l'aide de KubernetesPodOperator. Les tâches sont exécutées dans des pods, ce qui permet d'isoler les tâches au niveau des pods et d'améliorer la gestion des ressources.
Toutefois, il existe quelques différences clés :
- KubernetesExecutor n'exécute les tâches que dans l'espace de noms Managed Airflow versionné de votre environnement. Il n'est pas possible de modifier cet espace de noms dans Managed Airflow. Vous pouvez spécifier un espace de noms dans lequel KubernetesPodOperator exécute les tâches de pod.
- KubernetesExecutor peut utiliser n'importe quel opérateur Airflow intégré. KubernetesPodOperator n'exécute qu'un script fourni défini par le point d'entrée du conteneur.
- KubernetesExecutor utilise l'image Docker Managed Airflow par défaut avec les mêmes remplacements d'options de configuration Python et Airflow, variables d'environnement et packages PyPI que ceux définis dans votre environnement Managed Airflow.
À propos des images Docker
Par défaut, KubernetesExecutor lance les tâches à l'aide de la même image Docker que celle utilisée par Managed Airflow pour les nœuds de calcul Celery. Il s'agit de l' image Managed Airflow pour votre environnement, avec toutes les modifications que vous avez spécifiées pour votre environnement, telles que les packages PyPI personnalisés ou les variables d'environnement.
Avant de commencer
Vous pouvez utiliser CeleryKubernetesExecutor dans Managed Airflow (3e génération).
Il n'est pas possible d'utiliser un exécuteur autre que CeleryKubernetesExecutor dans Managed Airflow (3e génération). Cela signifie que vous pouvez exécuter des tâches à l'aide de CeleryExecutor, de KubernetesExecutor ou des deux dans un DAG, mais il n'est pas possible de configurer votre environnement pour qu'il n'utilise que KubernetesExecutor ou CeleryExecutor.
Configurer CeleryKubernetesExecutor
Vous pouvez remplacer les options de configuration Airflow existantes liées à KubernetesExecutor :
[kubernetes]worker_pods_creation_batch_sizeCette option définit le nombre d'appels de création de pods de nœuds de calcul Kubernetes par boucle de programmeur. La valeur par défaut est
1. Par conséquent, un seul pod est lancé par pulsation de programmeur. Si vous utilisez KubernetesExecutor de manière intensive, nous vous recommandons d'augmenter cette valeur.[kubernetes]worker_pods_pending_timeoutCette option définit, en secondes, la durée pendant laquelle un nœud de calcul peut rester à l'état
Pending(Pod en cours de création) avant d'être considéré comme ayant échoué. La valeur par défaut est de 5 minutes.
Exécuter des tâches avec KubernetesExecutor ou CeleryExecutor
Vous pouvez exécuter des tâches à l'aide de CeleryExecutor, de KubernetesExecutor ou des deux dans un DAG :
Airflow 3
- Pour exécuter une tâche avec KubernetesExecutor, spécifiez la valeur
KubernetesExecutordans le paramètreexecutord'une tâche. - Pour exécuter une tâche avec CeleryExecutor, omettez le paramètre
executor.
Airflow 2
- Pour exécuter une tâche avec KubernetesExecutor, spécifiez la valeur
kubernetesdans le paramètrequeued'une tâche. - Pour exécuter une tâche avec CeleryExecutor, omettez le paramètre
queue.
L'exemple suivant exécute la tâche task-kubernetes à l'aide de KubernetesExecutor et la tâche task-celery à l'aide de CeleryExecutor :
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
Exécuter des commandes de CLI Airflow liées à KubernetesExecutor
Vous pouvez exécuter plusieurs
commandes de CLI Airflow liées à KubernetesExecutor
à l'aide de gcloud.
Personnaliser les spécifications des pods de nœuds de calcul
Vous pouvez personnaliser les spécifications des pods de nœuds de calcul en les transmettant dans le paramètre executor_config d'une tâche. Vous pouvez l'utiliser pour définir des exigences personnalisées en matière de processeur et de mémoire.
Vous pouvez remplacer l'intégralité des spécifications des pods de nœuds de calcul utilisées pour exécuter une tâche. Pour
récupérer les spécifications des pods d'une tâche utilisée par KubernetesExecutor, vous pouvez
exécuter la kubernetes generate-dag-yaml commande de CLI Airflow.
Pour en savoir plus sur la personnalisation des spécifications des pods de nœuds de calcul, consultez la documentation Airflow.
Managed Airflow (3e génération) accepte les valeurs suivantes pour les exigences en matière de ressources :
| Ressource | Minimum | Maximum | Étape |
|---|---|---|---|
| Processeur | 0,25 | 32 | Valeurs d'étape : 0,25, 0,5, 1, 2, 4, 6, 8, 10, ..., 32. Les valeurs demandées sont arrondies à la valeur d'étape compatible la plus proche (par exemple, 5 à 6). |
| Mémoire | 2 G (Go) | 128 G (Go) | Valeurs d'étape : 2, 3, 4, 5, ..., 128. Les valeurs demandées sont arrondies à la valeur d'étape compatible la plus proche (par exemple, 3,5 G à 4 G). |
| Stockage | - | 100 G (Go) | N'importe quelle valeur. Si plus de 100 Go sont demandés, seuls 100 Go sont fournis. |
Pour en savoir plus sur les unités de ressources dans Kubernetes, consultez Unités de ressources dans Kubernetes.
L'exemple suivant illustre une tâche qui utilise des spécifications de pods de nœuds de calcul personnalisées :
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',
})
),
],
),
)
},
)
Afficher les journaux des tâches
Les journaux des tâches exécutées par KubernetesExecutor sont disponibles dans l'onglet Journaux , ainsi que les journaux des tâches exécutées par CeleryExecutor :
Dans la Google Cloud console, accédez à la page Environnements.
Dans la liste des environnements, cliquez sur le nom de votre environnement. La page Détails de l'environnement s'ouvre.
Accédez à l'onglet Journaux.
Accédez à Tous les journaux > Journaux Airflow > Nœuds de calcul.
Les nœuds de calcul nommés
airflow-k8s-workerexécutent les tâches KubernetesExecutor. Pour rechercher les journaux d'une tâche spécifique, vous pouvez utiliser un ID de DAG ou un ID de tâche comme mot clé dans la recherche.
Étape suivante
- Résoudre les problèmes liés à KubernetesExecutor
- Utiliser KubernetesPodOperator
- Utiliser des opérateurs GKE
- Remplacer les options de configuration Airflow