Apache Beam BigQuery IO是否支持流式插入JSON类型字段?
处理BigQuery JSON类型字段的Apache Beam流式插入问题
问题描述
我使用Apache Beam(Dataflow)处理数据并插入Google BigQuery表,该表包含一个JSON类型字段test_jsontype_field。流水线从GCP PubSub读取JSON数据,通过TableRowJsonCoder将JSON字符串解码为TableRow对象,再用BigQueryIO.Write.Method.STREAMING_INSERTS流式插入BigQuery时失败,报错:
The field <FIELD> is not a record.
其中<FIELD>就是表中的JSON类型字段test_jsontype_field。
核心问题
Apache Beam向BigQuery流式插入时是否支持JSON类型字段?如果支持,需要哪些特定配置或编码实践确保兼容性?
转换代码
static TableRow convertJsonToTableRow(String json) { TableRow row; try (InputStream inputStream = new ByteArrayInputStream(json.getBytes())) { row = TableRowJsonCoder.of().decode(inputStream, Coder.Context.OUTER); logger.debug("message {}", row.toPrettyString()); } catch (IOException e) { throw new RuntimeException("failed to serialize json to table row: " + json, e); } return row; }
错误日志
org.apache.beam.sdk.util.UserCodeException: java.lang.RuntimeException: java.io.IOException: Insert failed: [{"errors":[{"debugInfo":"","location":"test_jsontype_field","message":"This field: test_jsontype_field is not a record.","reason":"invalid"}],"index":0}] at org.apache.beam.sdk.util.UserCodeException.wrap(UserCodeException.java:39) at org.apache.beam.sdk.io.gcp.bigquery.BatchedStreamingWrite$BatchAndInsertElements$DoFnInvoker.invokeFinishBundle(Unknown Source) at org.apache.beam.fn.harness.FnApiDoFnRunner.finishBundle(FnApiDoFnRunner.java:1776) at org.apache.beam.fn.harness.data.PTransformFunctionRegistry.lambda$register$0(PTransformFunctionRegistry.java:116) at org.apache.beam.fn.harness.control.ProcessBundleHandler.processBundle(ProcessBundleHandler.java:560) at org.apache.beam.fn.harness.control.BeamFnControlClient.delegateOnInstructionRequestType(BeamFnControlClient.java:150) at org.apache.beam.fn.harness.control.BeamFnControlClient$InboundObserver.lambda$onNext$0(BeamFnControlClient.java:115) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at org.apache.beam.sdk.util.UnboundedScheduledExecutorService$ScheduledFutureTask.run(UnboundedScheduledExecutorService.java:163) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:840) Caused by: java.lang.RuntimeException: java.io.IOException: Insert failed: [{"errors":[{"debugInfo":"","location":"test_jsontype_field","message":"This field: test_jsontype_field is not a record.","reason":"invalid"}],"index":0}] at org.apache.beam.sdk.io.gcp.bigquery.BatchedStreamingWrite.flushRows(BatchedStreamingWrite.java:416) at org.apache.beam.sdk.io.gcp.bigquery.BatchedStreamingWrite.access$900(BatchedStreamingWrite.java:67) at org.apache.beam.sdk.io.gcp.bigquery.BatchedStreamingWrite$BatchAndInsertElements.finishBundle(BatchedStreamingWrite.java:286) Caused by: java.io.IOException: Insert failed: [{"errors":[{"debugInfo":"","location":"test_jsontype_field","message":"This field: test_jsontype_field is not a record.","reason":"invalid"}],"index":0}] at org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:1262) at org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.insertAll(BigQueryServicesImpl.java:1281) at org.apache.beam.sdk.io.gcp.bigquery.BatchedStreamingWrite.flushRows(BatchedStreamingWrite.java:403) at org.apache.beam.sdk.io.gcp.bigquery.BatchedStreamingWrite.access$900(BatchedStreamingWrite.java:67) at org.apache.beam.sdk.io.gcp.bigquery.BatchedStreamingWrite$BatchAndInsertElements.finishBundle(BatchedStreamingWrite.java:286) at org.apache.beam.sdk.io.gcp.bigquery.BatchedStreamingWrite$BatchAndInsertElements$DoFnInvoker.invokeFinishBundle(Unknown Source) at org.apache.beam.fn.harness.FnApiDoFnRunner.finishBundle(FnApiDoFnRunner.java:1776) at org.apache.beam.fn.harness.data.PTransformFunctionRegistry.lambda$register$0(PTransformFunctionRegistry.java:116) at org.apache.beam.fn.harness.control.ProcessBundleHandler.processBundle(ProcessBundleHandler.java:560) at org.apache.beam.fn.harness.control.BeamFnControlClient.delegateOnInstructionRequestType(BeamFnControlClient.java:150) at org.apache.beam.fn.harness.control.BeamFnControlClient$InboundObserver.lambda$onNext$0(BeamFnControlClient.java:115) at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539) at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264) at org.apache.beam.sdk.util.UnboundedScheduledExecutorService$ScheduledFutureTask.run(UnboundedScheduledExecutorService.java:163) at java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) at java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) at java.base/java.lang.Thread.run(Thread.java:840)
解决方案
Apache Beam流式插入BigQuery完全支持JSON类型字段,问题出在TableRowJsonCoder的解码逻辑:它会把JSON字符串解析成嵌套的TableRow对象,但BigQuery的JSON类型字段需要接收原始JSON字符串,而非结构化的TableRow记录。
你需要修改转换逻辑,确保test_jsontype_field的值是原始JSON字符串:
方案1:手动解析JSON,保留原始字段值
static TableRow convertJsonToTableRow(String json) { try { JSONObject jsonObj = new JSONObject(json); TableRow row = new TableRow(); // 处理其他非JSON类型字段,示例为id字段 if (jsonObj.has("id")) { row.set("id", jsonObj.getLong("id")); } // 直接将test_jsontype_field的原始JSON字符串存入TableRow if (jsonObj.has("test_jsontype_field")) { row.set("test_jsontype_field", jsonObj.get("test_jsontype_field").toString()); } logger.debug("message {}", row.toPrettyString()); return row; } catch (JSONException e) { throw new RuntimeException("failed to process json: " + json, e); } }
方案2:解码后将嵌套TableRow转回JSON字符串
如果已经用TableRowJsonCoder完成解码,可以把嵌套的TableRow重新序列化为字符串:
static TableRow convertJsonToTableRow(String json) { TableRow row; try (InputStream inputStream = new ByteArrayInputStream(json.getBytes())) { row = TableRowJsonCoder.of().decode(inputStream, Coder.Context.OUTER); // 检查并转换JSON类型字段 Object jsonFieldValue = row.get("test_jsontype_field"); if (jsonFieldValue instanceof TableRow) { row.set("test_jsontype_field", ((TableRow) jsonFieldValue).toString()); } logger.debug("message {}", row.toPrettyString()); } catch (IOException e) { throw new RuntimeException("failed to serialize json to table row: " + json, e); } return row; }
额外验证
务必确认BigQuery表的schema中,test_jsontype_field的类型确实设置为JSON,而非RECORD类型。
内容的提问来源于stack exchange,提问作者Akhilesh
相关产品推荐
相关产品推荐

