使用Apache Beam的DoFn解析JSON时遇到异常行为求助
我之前也踩过这个坑!你遇到的问题本质是Apache Beam读取JSON文件的方式和JSON格式不匹配导致的——当你用格式化后的多行JSON时,Beam默认按行读取,每一行都不是完整的JSON对象,所以解析JsonHolder这个POJO时会报错;而单行JSON是完整的对象,自然能正常解析。
下面给你几个针对性的解决方案:
1. 读取整个文件作为单个元素解析
如果你的每个GCS文件都是一个独立的完整JSON对象,推荐用FileIO读取整个文件内容,而不是按行拆分:
Pipeline pipeline = Pipeline.create(options); pipeline.apply(FileIO.match().filepattern("gs://your-bucket/path/*.json")) .apply(FileIO.readMatches()) .apply(MapElements.via(new SimpleFunction<ReadableFile, String>() { @Override public String apply(ReadableFile file) { try { // 读取整个文件的完整JSON内容 return file.readFullyAsUTF8String(); } catch (IOException e) { throw new RuntimeException("读取文件失败: " + file.getMetadata().resourceId(), e); } } })) .apply(ParseJsons.of(JsonHolder.class)) // 后续写入Cloud SQL的逻辑 .apply(JdbcIO.<JsonHolder>write() .withDataSourceConfiguration(JdbcIO.DataSourceConfiguration.create( "com.mysql.cj.jdbc.Driver", "jdbc:mysql://your-cloud-sql-ip:3306/db-name") .withUsername("your-username") .withPassword("your-password")) .withStatement("INSERT INTO target_table (col1, col2) VALUES (?, ?)") .withPreparedStatementSetter((element, stmt) -> { stmt.setString(1, element.getCol1()); stmt.setString(2, element.getCol2()); })); pipeline.run().waitUntilFinish();
这种方式会把整个文件的内容作为一个完整的JSON字符串传递给解析器,完美适配格式化后的多行JSON。
2. 转用JSON Lines格式存储数据
如果你的数据是多个JSON对象集合,推荐把文件改成JSON Lines格式(每行一个独立的JSON对象):
{"col1": "value1", "col2": "value2"} {"col1": "value3", "col2": "value4"}
这种情况下,你原来用TextIO.read()的代码就能正常工作,无需大改:
pipeline.apply(TextIO.read().from("gs://your-bucket/path/*.json")) .apply(ParseJsons.of(JsonHolder.class)) // 后续写入Cloud SQL的逻辑
3. 自定义逻辑拼接多行内容(不推荐大文件)
如果无法修改文件格式,也可以自定义DoFn把同一文件的所有行拼接成完整JSON后再解析:
pipeline.apply(TextIO.read().from("gs://your-bucket/path/*.json") .withFilenamePolicy(DefaultFilenamePolicy.fromStandardParameters()) ) .apply(GroupByKey.create()) // 按文件名分组,把同一文件的所有行聚合起来 .apply(MapElements.via(new SimpleFunction<KV<String, Iterable<String>>, JsonHolder>() { @Override public JsonHolder apply(KV<String, Iterable<String>> input) { // 拼接所有行成完整JSON String fullJson = String.join("\n", input.getValue()); ObjectMapper mapper = new ObjectMapper(); try { return mapper.readValue(fullJson, JsonHolder.class); } catch (JsonProcessingException e) { throw new RuntimeException("解析文件JSON失败: " + input.getKey(), e); } } })) // 后续写入Cloud SQL的逻辑
注意:这种方式对大文件不友好,会占用较多内存,优先推荐前两种方案。
额外注意事项
- 确保你的
JsonHolder类字段名和JSON键完全一致,或者添加了正确的Jackson注解(比如@JsonProperty),否则也会出现解析异常。 - 验证Cloud SQL的连接配置(驱动、地址、账号密码)是否正确,避免解析成功后写入环节出问题。
内容的提问来源于stack exchange,提问作者R Wri
相关产品推荐
相关产品推荐

