如何将GCP Workflows中的参数传递至Dataproc?
GCP Workflows传递参数到Dataproc Batch的正确方法
在GCP Workflows中定义数据工程流程时,若要将Workflow变量作为参数传递给Dataproc Batch,需要确保表达式被正确解析,避免直接传递纯文本。
问题场景
已通过Workflow API传入参数,初始化了变量:
main: params: [input] steps: - init: assign: - customer_id: ${default(map.get(input, "customer_id"), "default_customer")} - client_id: ${default(map.get(input, "client_id"), "default_client")}
日志打印变量正常,但传递到Dataproc Batch的args时,变量未被解析,直接以原字符串形式发送(如text.concat("--customer_id=", customer_id))。
错误写法示例
- submit_dataproc_batch: call: http.post args: url: "https://dataproc.googleapis.com/v1/projects/MY_PROJECT/locations/LOCATION/batches" auth: type: OAuth2 headers: Content-Type: "application/json" body: pysparkBatch: mainPythonFileUri: "gs://MY_BUCKET/MY_FOLDER/123123123-CA42-4773-AE9E-84BB0CED4057/main.py" jarFileUris: - "gs://MY_BUCKET/jars/delta-spark_2.13-3.1.0.jar" pythonFileUris: - "gs://MY_BUCKET/MY_FOLDER/123123123-CA42-4773-AE9E-84BB0CED4057/code.zip" args: - text.concat("--customer_id=", customer_id) - text.concat("--customer_id2=", ${customer_id}) - "--customer_id=${customer_id}"
正确解决方法
GCP Workflows仅会解析${}包裹的内容作为表达式执行,因此需要将参数拼接逻辑放在${}内:
方法1:直接在args中拼接
- submit_dataproc_batch: call: http.post args: url: "https://dataproc.googleapis.com/v1/projects/MY_PROJECT/locations/LOCATION/batches" auth: type: OAuth2 headers: Content-Type: "application/json" body: pysparkBatch: mainPythonFileUri: "gs://MY_BUCKET/MY_FOLDER/123123123-CA42-4773-AE9E-84BB0CED4057/main.py" jarFileUris: - "gs://MY_BUCKET/jars/delta-spark_2.13-3.1.0.jar" pythonFileUris: - "gs://MY_BUCKET/MY_FOLDER/123123123-CA42-4773-AE9E-84BB0CED4057/code.zip" args: - ${text.concat("--customer_id=", customer_id)} - ${text.concat("--client_id=", client_id)}
方法2:提前拼接参数字符串
若参数较多,可在初始化步骤提前拼接好参数,再传入args,提升可读性:
main: params: [input] steps: - init: assign: - customer_id: ${default(map.get(input, "customer_id"), "default_customer")} - client_id: ${default(map.get(input, "client_id"), "default_client")} - customer_arg: ${text.concat("--customer_id=", customer_id)} - client_arg: ${text.concat("--client_id=", client_id)} - submit_dataproc_batch: call: http.post args: url: "https://dataproc.googleapis.com/v1/projects/MY_PROJECT/locations/LOCATION/batches" auth: type: OAuth2 headers: Content-Type: "application/json" body: pysparkBatch: mainPythonFileUri: "gs://MY_BUCKET/MY_FOLDER/123123123-CA42-4773-AE9E-84BB0CED4057/main.py" jarFileUris: - "gs://MY_BUCKET/jars/delta-spark_2.13-3.1.0.jar" pythonFileUris: - "gs://MY_BUCKET/MY_FOLDER/123123123-CA42-4773-AE9E-84BB0CED4057/code.zip" args: - ${customer_arg} - ${client_arg}
内容的提问来源于stack exchange,提问作者54m
相关产品推荐
相关产品推荐

