无法将GCP Composer的Airflow任务提交至Dataproc的问题求助
解决Airflow提交Dataproc任务报错:Protocol message JobReference has no "projectId" field
核心问题
Google Cloud Dataproc的Python客户端库采用**蛇形命名(snake_case)**定义proto消息字段,而非JSON常用的驼峰命名(camelCase)。你代码中job['reference']里的projectId是错误字段名,正确应为project_id。
修复方案
1. 统一修正字段命名
将所有驼峰命名的字段替换为蛇形命名:
projectId→project_idclusterName→cluster_namesparkJob→spark_jobmainClass→main_classjarFileUris→jar_file_uris
2. 补全字符串常量的引号
代码中my-project-id、my-cluster如果是字符串常量,必须用引号包裹,否则会被Python识别为未定义变量引发额外错误。
3. 可选简化:复用Operator的project_id配置
你已经在DataprocSubmitJobOperator层面指定了project_id='my-project',无需在job['reference']中重复配置,客户端会自动继承该值,进一步简化代码。
修正后的完整代码
submit_dataproc_task = DataprocSubmitJobOperator( task_id='submit_spark_job_to_dataproc', project_id='my-project', region='my-region', # 替换为实际Dataproc区域 job={ 'placement': {'cluster_name': 'my-cluster'}, 'spark_job': { 'main_class': '', 'jar_file_uris': ['{{ task_instance.xcom_pull("download_jar_task") }}'], 'args': [], 'properties': { 'spark.corebanking.input.dir': 'file:///home/sample.avro' } } }, dag=dag, )
内容的提问来源于stack exchange,提问作者luckyman80
相关产品推荐
相关产品推荐

