如何使用Apache Beam与Jackson解析非换行分隔的普通JSON并转换为CSV格式?
问题分析
你碰到的这个错误根源在于ParseJsons.of(Person.class)的默认行为:它期望输入的每个字符串元素都是完整的、能独立解析为Person的JSON对象。你提到去掉空格后能正常解析,其实是因为去掉空格后的JSON变成了单行,TextIO.read()刚好读取到完整的一行(也就是完整的JSON对象),但这只是巧合,不是通用解决方案。
错误信息里提示“expected ']'”,大概率是因为当TextIO.read()把多行的JSON对象拆成了多个不完整的行片段时,Jackson尝试解析这些片段,误以为应该是数组结构,从而抛出了格式错误。
解决方案
下面给你两种实用的解决方法,覆盖不同的场景:
方案1:读取整个文件作为单个字符串(适合单个JSON对象的文件)
如果你的输入文件就是单个独立的JSON对象(不管是单行还是多行),可以用WholeFileIO替代TextIO来读取整个文件内容,确保传入解析器的是完整的JSON:
import org.apache.beam.sdk.io.WholeFileIO; import org.apache.beam.sdk.transforms.MapElements; import org.apache.beam.sdk.values.TypeDescriptors; public class DataToModel { public static void main(String[] args) { PipelineOptions options = PipelineOptionsFactory.create(); options.setRunner(DirectRunner.class); Pipeline p = Pipeline.create(options); // 读取整个文件为一个字符串 PCollection<String> json = p.apply(WholeFileIO.read().from("src/main/resources/test.json")) .apply(MapElements.into(TypeDescriptors.strings()) .via(fileResult -> new String(fileResult.readFullyAsBytes()))); PCollection<Person> person = json .apply(ParseJsons.of(Person.class)) .setCoder(SerializableCoder.of(Person.class)); // 后续提取和写入逻辑不变 PCollection<String> names = person.apply(MapElements .into(TypeDescriptors.strings()) .via(Person::getFirstName) ); names.apply(TextIO.write().to("src/main/resources/test_out")); p.run().waitUntilFinish(); } }
方案2:自定义解析逻辑(支持多格式输入)
如果你的场景需要兼容单个JSON对象、JSON数组、每行一个JSON对象(JSON Lines)等多种格式,可以自定义一个ParDo来灵活处理:
import com.fasterxml.jackson.databind.ObjectMapper; import com.fasterxml.jackson.databind.type.CollectionType; import org.apache.beam.sdk.transforms.ParDo; import org.apache.beam.sdk.transforms.DoFn; public class DataToModel { public static void main(String[] args) { PipelineOptions options = PipelineOptionsFactory.create(); options.setRunner(DirectRunner.class); Pipeline p = Pipeline.create(options); PCollection<String> json = p.apply(TextIO.read().from("src/main/resources/test.json")); // 自定义Jackson解析逻辑,兼容单个对象和数组 ObjectMapper mapper = new ObjectMapper(); CollectionType personListType = mapper.getTypeFactory() .constructCollectionType(java.util.List.class, Person.class); PCollection<Person> person = json .apply(ParDo.of(new DoFn<String, Person>() { @ProcessElement public void processElement(@Element String jsonStr, OutputReceiver<Person> out) throws Exception { try { // 先尝试解析为单个Person对象 out.output(mapper.readValue(jsonStr, Person.class)); } catch (Exception e) { // 解析单个对象失败时,尝试解析为数组并批量输出 java.util.List<Person> persons = mapper.readValue(jsonStr, personListType); out.outputAll(persons); } } })) .setCoder(SerializableCoder.of(Person.class)); // 后续处理逻辑不变 PCollection<String> names = person.apply(MapElements .into(TypeDescriptors.strings()) .via(Person::getFirstName) ); names.apply(TextIO.write().to("src/main/resources/test_out")); p.run().waitUntilFinish(); } }
这种方法更通用,不管你的输入是哪种JSON格式,都能正确解析出Person对象。
另外补充一句:如果你的实际业务是处理批量的JSON对象(比如每行一个),只需要确保每个JSON对象单独占一行,TextIO.read()会按行读取,ParseJsons.of(Person.class)就能直接正常工作,不需要额外修改。
内容的提问来源于stack exchange,提问作者rocksNwaves
相关产品推荐
相关产品推荐

