Airflow 3.1.5中SparkSubmitOperator pod_overwrite JSON验证失败求助
问题场景
我在Airflow 3.1.5环境中使用SparkSubmitOperator运行DAG,自定义了Pod配置的包装器代码如下:
self.executor_config = { "pod_override": k8s.V1Pod( metadata=k8s.V1ObjectMeta(labels={"spark-role": "driver"}), spec=k8s.V1PodSpec( init_containers=[self._build_init_container_spec(pod_env)], containers=[ k8s.V1Container( name="base", env=pod_env, image=self.image, image_pull_policy="IfNotPresent", resources=k8s.V1ResourceRequirements( limits={ "cpu": str(driver_cores), "memory": "32Gi", }, requests={ "cpu": str(driver_cores), "memory": driver_memory.to_k8s_spec(), }, ), security_context=k8s.V1SecurityContext( run_as_user=1001, run_as_group=1000, ), ), ], ), ), }
Pod启动后验证工作负载失败,实际执行的命令中--json-string参数的JSON丢失了双引号,格式无效:
exec /usr/bin/tini -s -- python -m airflow.sdk.execution_time.execute_workload \ --json-string '{token:some_value,somekyes:somevalues}'
但Pod的YAML文件中args字段的JSON是正确带双引号的:
'{"token":"some_value","somekyes":"somevalues"}'
问题出在YAML到容器CLI的转换过程中双引号被剥离。我考虑改用--json-path参数,但不清楚如何创建并挂载包含工作负载的文件,同时希望execute_workload能兼容两种JSON传入方式。
解决方案
方案一:修复双引号被剥离问题
在构建Pod配置时,显式指定容器的command和args,通过单引号包裹JSON字符串,避免YAML解析时剥离双引号:
# 修改V1Container配置,添加显式的command和args containers=[ k8s.V1Container( name="base", env=pod_env, image=self.image, image_pull_policy="IfNotPresent", # 手动构造执行命令,用单引号包裹JSON字符串 command=["/usr/bin/tini", "-s", "--", "python", "-m", "airflow.sdk.execution_time.execute_workload"], args=["--json-string", '{"token":"some_value","somekyes":"somevalues"}'], resources=k8s.V1ResourceRequirements( limits={ "cpu": str(driver_cores), "memory": "32Gi", }, requests={ "cpu": str(driver_cores), "memory": driver_memory.to_k8s_spec(), }, ), security_context=k8s.V1SecurityContext( run_as_user=1001, run_as_group=1000, ), ), ],
这种方式跳过Airflow自动生成命令的逻辑,直接传递带单引号包裹的有效JSON,确保CLI执行时双引号不会丢失。
方案二:使用--json-path挂载文件传递JSON
如果要改用文件方式传递,需要完成以下三步:
通过InitContainer生成JSON文件
修改初始化容器的逻辑,生成包含工作负载JSON的文件:def _build_init_container_spec(self, pod_env): return k8s.V1Container( name="init-workload", image="busybox:latest", command=["sh", "-c"], args=[ 'echo \'{"token":"some_value","somekyes":"somevalues"}\' > /workload/workload.json' ], volume_mounts=[ k8s.V1VolumeMount( name="workload-volume", mount_path="/workload" ) ] )添加共享Volume
在PodSpec中定义emptyDir类型的Volume,用于在InitContainer和主容器之间共享文件:spec=k8s.V1PodSpec( volumes=[ k8s.V1Volume( name="workload-volume", empty_dir=k8s.V1EmptyDirVolumeSource() ) ], init_containers=[self._build_init_container_spec(pod_env)], containers=[ k8s.V1Container( name="base", env=pod_env, image=self.image, image_pull_policy="IfNotPresent", # 挂载Volume到主容器 volume_mounts=[ k8s.V1VolumeMount( name="workload-volume", mount_path="/workload" ) ], # 修改命令为使用--json-path参数 command=["/usr/bin/tini", "-s", "--", "python", "-m", "airflow.sdk.execution_time.execute_workload"], args=["--json-path", "/workload/workload.json"], resources=k8s.V1ResourceRequirements(...), security_context=k8s.V1SecurityContext(...), ), ], ),修改
execute_workload支持--json-path
如果当前脚本不支持该参数,需要添加文件读取逻辑:import argparse import json parser = argparse.ArgumentParser() parser.add_argument("--json-string", type=str) parser.add_argument("--json-path", type=str) args = parser.parse_args() if args.json_string: workload = json.loads(args.json_string) elif args.json_path: with open(args.json_path, 'r') as f: workload = json.load(f) else: raise ValueError("必须提供--json-string或--json-path参数")
内容的提问来源于stack exchange,提问作者twierdzenie twierdzenie
相关产品推荐
相关产品推荐

