如何在同一个K8s Pod中运行多个Airflow任务?
如何在同一个K8s Pod中运行多个Airflow任务?
嗨,我来帮你解决这个问题!你现在用KubernetesPodOperator创建了两个任务,每个任务都会启动独立的Pod,想让它们共享同一个Pod运行对吧?这里有几种不同的方案,你可以根据自己的需求选择:
方法一:合并脚本到单个任务(最简单直接)
如果不需要在Airflow UI里把两个脚本拆成独立的任务节点,直接把两个脚本的执行命令串联起来,放到一个KubernetesPodOperator里就行。这样只会创建一个Pod,先后执行两个脚本:
from airflow.providers.cncf.kubernetes.operators.kubernetes_pod import KubernetesPodOperator from airflow.utils.dates import days_ago from datetime import timedelta default_args = { 'owner': 'airflow', 'start_date': days_ago(1), } with DAG( 'test-1', default_args=default_args, schedule_interval=timedelta(minutes=100) ) as dag: combined_test_task = KubernetesPodOperator( task_id="combined_test", image="demo111:1.91", cmds=["/bin/bash", "-c"], arguments=["./test-1.sh && ./test-2.sh"], # 记得补充其他必要参数,比如pod名称、命名空间等 )
方法二:复用Pod运行独立任务节点(适合需要UI区分任务的场景)
如果你希望在Airflow UI里看到两个独立的任务节点(test1和test2),但它们共享同一个Pod执行,可以利用Airflow 2.x+的Pod复用特性。
首先需要在Airflow配置中开启Pod复用(或者在DAG/任务级别单独配置):
# airflow.cfg 中的配置 [kubernetes] pod_reuse_mode = ReusePodMode.REUSE_POD
然后你的DAG可以保持原来的任务结构,任务会自动复用同一个Pod:
from airflow.providers.cncf.kubernetes.operators.kubernetes_pod import KubernetesPodOperator from airflow.utils.dates import days_ago from datetime import timedelta default_args = { 'owner': 'airflow', 'start_date': days_ago(1), } with DAG( 'test-1', default_args=default_args, schedule_interval=timedelta(minutes=100), # 在DAG级别配置Pod复用 executor_config={"pod_reuse_mode": "ReusePodMode.REUSE_POD"} ) as dag: test1 = KubernetesPodOperator( task_id="test1", image="demo111:1.91", cmds=["/bin/bash", "-c"], arguments=["./test-1.sh"], # 其他必要参数 ) test2 = KubernetesPodOperator( task_id="test2", image="demo111:1.91", cmds=["/bin/bash", "-c"], arguments=["./test-2.sh"], # 其他必要参数 ) test1 >> test2
注意:这种方式要求两个任务之间没有资源冲突(比如不会互相占用文件或端口),Pod会在所有任务执行完成后才销毁。
方法三:使用多容器Pod(并行执行脚本)
如果想让两个脚本在同一个Pod里并行执行,可以定义一个包含两个容器的Pod,每个容器对应一个脚本:
from airflow.providers.cncf.kubernetes.operators.kubernetes_pod import KubernetesPodOperator from airflow.providers.cncf.kubernetes.pod_generator import PodGenerator from kubernetes.client.models import V1Pod, V1Container, V1PodSpec from airflow.utils.dates import days_ago from datetime import timedelta default_args = { 'owner': 'airflow', 'start_date': days_ago(1), } with DAG( 'test-1', default_args=default_args, schedule_interval=timedelta(minutes=100) ) as dag: # 定义第一个容器,执行test-1.sh test1_container = V1Container( name="test1-container", image="demo111:1.91", command=["/bin/bash", "-c"], args=["./test-1.sh"] ) # 定义第二个容器,执行test-2.sh test2_container = V1Container( name="test2-container", image="demo111:1.91", command=["/bin/bash", "-c"], args=["./test-2.sh"] ) # 组装Pod Spec pod_spec = V1PodSpec(containers=[test1_container, test2_container]) pod = V1Pod(spec=pod_spec) multi_container_task = KubernetesPodOperator( task_id="multi_container_test", pod_template=PodGenerator.generate_pod(pod), # 其他必要参数 )
这种方式下两个容器会在同一个Pod里同时运行脚本,但Airflow UI只会显示一个任务节点。如果需要串行执行,你得在容器命令里添加依赖逻辑(比如让第二个容器等待第一个容器执行完成)。
备注:内容来源于stack exchange,提问作者oTolev
相关产品推荐
相关产品推荐

