You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.30 23:42:38