Cloud Composer中DAG写入BigQuery失败报java.lang.ClassNotFoundException,如何配置DataprocSubmitJobOperator添加指定Jar包
解决Cloud Composer中Dataproc任务写入BigQuery的ClassNotFoundException问题
嗨,我懂你的烦恼!你在gcloud命令里靠--jars参数指定BigQuery依赖Jar包就能正常跑任务,但不知道怎么把这个配置搬到Airflow的DataprocSubmitJobOperator里对吧?其实特别简单,对应到Airflow的PySpark任务配置中,你只需要在PYSPARK_JOB里添加jar_file_uris参数就行。
修改后的Airflow代码示例
先更新你的PYSPARK_JOB定义,把BigQuery的Jar包路径加到jar_file_uris列表里:
PYSPARK_JOB = { "reference": {"project_id": PROJECT_ID}, "placement": {"cluster_name": "cluster-bc4b"}, # 替换成你的集群名称 "pyspark_job": { "main_python_file_uri": "gs://your-bucket-path/examen/ETL/loadBQ.py", # 替换成PySpark脚本在GCS的实际路径 "jar_file_uris": ["gs://spark-lib/bigquery/spark-bigquery-latest_2.12.jar"] # 对应gcloud命令中的--jars参数 } } submit_job = DataprocSubmitJobOperator( task_id="pyspark_task", job=PYSPARK_JOB, location=REGION, project_id=PROJECT_ID )
关键细节说明
jar_file_uris就是PySpark任务配置里用来指定依赖Jar包的参数,和你gcloud命令里的--jars功能完全一致,Dataproc集群运行任务时会自动加载这些Jar包。- 注意Airflow里不能用本地路径
~/...,得把你的PySpark脚本上传到GCS,然后把main_python_file_uri替换成脚本在GCS上的路径。
这样配置后,你的PySpark任务提交到Dataproc集群时就会带上BigQuery的依赖Jar包,java.lang.ClassNotFoundException: Failed to find data source: bigquery这个报错就会消失啦!
内容的提问来源于stack exchange,提问作者Nicolas Soria
相关产品推荐
相关产品推荐

