Cloud Composer中DataprocCreateBatchOperator不识别PySpark作业参数
问题根因
该报错由pyspark_batch.args字段的传参格式不符合Dataproc Serverless API要求导致:
- gcloud命令提交作业时,shell会自动将
--后空格分隔的参数名、参数值拆分为独立项,由gcloud封装为符合API规范的格式发送,因此可以正常运行。 - 直接通过算子调用API时,
args数组要求每个参数名、每个参数值都必须作为单独的数组元素存在,不能将参数名和对应值拼接在同一个字符串中。 - 当前配置将
--run_timestamp "{{ ts }}"这类“参数名+空格+参数值”的整串作为单个数组元素传入,Dataproc会将整串识别为一个完整参数名传给PySpark作业,底层argparse解析器无法识别这类异常参数,就会抛出未识别参数的错误。 - 此前尝试给参数值加/移除双引号无法解决问题,是因为核心矛盾是参数项未正确拆分,和值外层的引号无关:API传参直接传递原生字符串,不需要像shell那样用引号处理带空格的值,额外添加的双引号反而会被当作值的一部分传入作业。
修复方案
调整pyspark_batch下的args数组配置,将每个参数名、参数值拆分为独立的数组元素,移除参数值外层多余的双引号即可,修正后的完整配置如下(同时补全了原配置中subnetwork_uri字段缺失的闭合引号):
create_batch = DataprocCreateBatchOperator( task_id="hourly_pipeline", project_id="dev", region="us-west1", batch_id="".join(random.choice(string.ascii_lowercase + string.digits + "-") for i in range(40)), batch={ "environment_config": { "execution_config": { "service_account": "<service_account>", "subnetwork_uri": "<uri>" } }, "pyspark_batch": { "main_python_file_uri": "gs://pipeline/code/pipeline_feat_creation.py", "args": [ '--run_timestamp', '{{ ts }}', '--temp_bucket', 'gs://pipeline/spark_temp_bucket/hourly/', '--bucket', 'pipeline', '--pipeline', 'hourly' ], "jar_file_uris": [ "gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.25.0.jar" ], } } )
注意:如果参数值本身包含空格,只需要把带空格的完整值作为单个数组元素传入即可,不需要额外加双引号包裹。
内容的提问来源于stack exchange,提问作者Frank Pinto
相关产品推荐
相关产品推荐

