GCP Workflows捕获PostgreSQL导入HTTP响应错误问题求助
GCP Workflows 导入BQ数据到PostgreSQL的错误处理问题
问题背景
通过GCP Workflows调用PostgreSQL的importContext接口将BigQuery数据导入PostgreSQL,流程中通过HTTP Get轮询导入状态。但导入出错时,接口返回状态码200且status字段仍为DONE,仅响应体包含error节点,导致现有错误判断逻辑完全失效。需要调整流程逻辑,让导入出错时流程主动失败,同时解决尝试将响应存入map未成功的问题。
现有代码与响应示例
导入操作代码
steps: - callImport: call: http.post args: url: ${"https://sqladmin.googleapis.com/v1/projects/" + projectid + "/instances/" + instance + "/import"} auth: type: OAuth2 body: importContext: uri: ${file} database: ${databaseschema} fileType: CSV csvImportOptions: table: ${importtable} columns : [a,b,c,b] result: operation
成功响应体
{"body":{"endTime":"2023-07-26T12:15:55.629Z","importContext":{"csvImportOptions":{"columns":[strings],"table":"table_name"},"database":"postgres","fileType":"CSV","kind":"sql#importContext","uri":"gs://workflow/2023-07-26000000000000.csv"},"insertTime":"2023-07-26T12:15:44.791Z","kind":"sql#operation","name":"af439df3-21ea-45b2-a92c-a4de00000024","operationType":"IMPORT","selfLink":"https://sqladmin.googleapis.com/v1/projects/","startTime":"2023-07-26T12:15:45.027Z","status":"DONE","targetId":"pricing-dev-master","targetLink":"https://sqladmin.googleapis.com/v1/projects/","targetProject":"project-dev","user":"sa@.iam.gserviceaccount.com"},"code":200,"headers":{"Alt-Svc":"h3=\":443\"; ma=2592000,h3-29=\":443\"; ma=2592000","Cache-Control":"private","Content-Length":"1179","Content-Type":"application/json; charset=UTF-8","Date":"Wed, 26 Jul 2023 12:16:00 GMT","Server":"ESF","Vary":"Origin, X-Origin, Referer","X-Content-Type-Options":"nosniff","X-Frame-Options":"SAMEORIGIN","X-Xss-Protection":"0"}}
错误响应体
{"body":{"endTime":"2023-07-26T10:19:49.527Z","error":{"errors":[{"code":"ERROR_RDBMS","kind":"sql#operationError","message":"generic::failed_precondition: ERROR: invalid input syntax for type integer: \"2023-07-11\"\nCONTEXT: COPY table_name, line 1, column id: \"2023-07-11\"\n"}],"kind":"sql#operationErrors"},"importContext":{"csvImportOptions":{"table":"table_name"},"database":"postgres","fileType":"CSV","kind":"sql#importContext","uri":"gs://workflow/2023-07-26000000000000.csv"},"insertTime":"2023-07-26T10:19:38.945Z","kind":"sql#operation","name":"375f1c18-c2ac-4b1e-800c-650f00000024","operationType":"IMPORT","selfLink":"https://sqladmin.googleapis.com/v1/projects/","startTime":"2023-07-26T10:19:39.118Z","status":"DONE","targetId":"pricing-dev-master","targetLink":"https://sqladmin.googleapis.com/v1/projects/","targetProject":"project-dev","user":"sa@.iam.gserviceaccount.com"},"code":200,"headers":{"Alt-Svc":"h3=\":443\"; ma=2592000,h3-29=\":443\"; ma=2592000","Cache-Control":"private","Content-Length":"1290","Content-Type":"application/json; charset=UTF-8","Date":"Wed, 26 Jul 2023 10:19:49 GMT","Server":"ESF","Vary":"Origin, X-Origin, Referer","X-Content-Type-Options":"nosniff","X-Frame-Options":"SAMEORIGIN","X-Xss-Protection":"0"}}
现有无效判断逻辑
- chekoperation: switch: - condition: ${operation.body.status != "DONE"} next: wait - condition: ${"sql#operationError" in progress_Map} raise: "Error in the load" next : completed - completed: return: "done" - wait: call: sys.sleep args: seconds: 5 next: getoperation - getoperation: call: http.get args: url: ${operation.body.selfLink} auth: type: OAuth2 result: operation next: trackprogress # Creating Maps to add the response body into progress_Map - trackprogress: assign: - progress_Map[file]: ${operation.body} raise : $(progress_Map) - returnoutput: return: ${progress_Map} next: chekoperation
- chekoperation: switch: - condition: ${operation.body.status != "DONE"} # Tried this as well #- condition: ${not("DONE" in progress_Map)} next: wait - condition: ${"sql#operationError" in progress_Map} raise: "Error in the load" next : completed
解决方案
核心问题分析
- 无需额外存储map:轮询时直接检查当前HTTP Get返回的
operation.body即可,额外存储map完全多余,还会增加逻辑复杂度。 - 错误判断顺序错误:应该先检查响应是否包含
error字段,再判断是否完成;现有逻辑先判断状态是否为DONE,导致出错时直接进入completed分支,跳过错误检查。 - 条件判断写法错误:
"sql#operationError" in progress_Map的写法完全错误,应该直接检查operation.body.error是否存在,或者检查错误列表中的kind字段。
修正后完整流程代码
steps: - callImport: call: http.post args: url: ${"https://sqladmin.googleapis.com/v1/projects/" + projectid + "/instances/" + instance + "/import"} auth: type: OAuth2 body: importContext: uri: ${file} database: ${databaseschema} fileType: CSV csvImportOptions: table: ${importtable} columns : [a,b,c,b] result: operation next: checkOperation - checkOperation: switch: # 优先检查是否存在错误:如果响应体有error节点,直接抛出错误 - condition: ${operation.body.error != null} raise: ${"导入失败:" + operation.body.error.errors[0].message} # 状态未完成则等待后轮询 - condition: ${operation.body.status != "DONE"} next: waitAndPoll # 无错误且状态完成,返回成功 next: completed - waitAndPoll: call: sys.sleep args: seconds: 5 next: getOperationStatus - getOperationStatus: call: http.get args: url: ${operation.body.selfLink} auth: type: OAuth2 result: operation next: checkOperation - completed: return: "导入完成"
关键修正点说明
- 错误优先检查:在
checkOperation步骤中,第一个条件就判断operation.body.error != null,只要存在错误节点就直接抛出包含具体错误信息的异常,终止流程。 - 简化轮询逻辑:移除多余的
progress_Map存储,直接复用operation变量存储每次轮询的响应,减少不必要的状态维护。 - 正确的条件判断:直接检查响应体中的
error字段是否存在,而不是错误地去map中查找字符串。 - 完整的轮询闭环:从
checkOperation到waitAndPoll再到getOperationStatus,形成完整的轮询循环,直到状态完成或发现错误。
内容的提问来源于stack exchange,提问作者Amdone
相关产品推荐
相关产品推荐

