使用DataprocInstantiateInlineWorkflowTemplateOperator运行PySpark作业报错
问题:DataprocInstantiateInlineWorkflowTemplateOperator 字段报错
使用DataprocInstantiateInlineWorkflowTemplateOperator运行PySpark作业时,先后遇到两个报错:
- 包含
stepID字段时:ValueError: Protocol message OrderedJob has no "stepID" field - 移除
stepID后:ValueError: Protocol message OrderedJob has no "pysparkJob" field
解决方案
问题根源是字段命名格式不匹配:Airflow的Dataproc操作符要求使用蛇形命名法(snake_case),而你使用了GCP API文档中的驼峰命名法(camelCase)。需要统一修改所有字段为蛇形命名:
修正后的工作流模板JSON
{ "id": "my-workflow-template", "jobs": [ { "step_id": "123456dfgy", "pyspark_job": { "main_python_file_uri": "gs://my-bucket/app.py" } } ], "name": "My Workflow Template", "placement": { "managed_cluster": { "cluster_name": "my-managed-cluster", "config": { "master_config": { "disk_config": { "boot_disk_size_gb": 1024, "boot_disk_type": "pd-standard" }, "machine_type_uri": "n1-standard-4", "num_instances": 1 }, "worker_config": { "disk_config": { "boot_disk_size_gb": 1024, "boot_disk_type": "pd-standard" }, "machine_type_uri": "n1-standard-4", "num_instances": 2 } } } } }
关键修改点
stepID→step_idpysparkJob→pyspark_jobmainPythonFileUri→main_python_file_urimanagedCluster→managed_clusterclusterName→cluster_name
原因说明
Airflow的GCP操作符内部会将传入的字典转换为Protocol Buffers对象,这些对象的字段遵循蛇形命名规则,和GCP REST API的驼峰命名不一致。直接照搬API文档的字段名会导致字段无法被识别,从而触发报错。
内容的提问来源于stack exchange,提问作者Aman Saurav
相关产品推荐
相关产品推荐

