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

Apache Beam PCollection转Row报IllegalStateException问题求助

问题:将String类型PCollection转换为Row类型PCollection失败

问题背景

  • 需求:需要将String/String数组类型的PCollection转换为Row类型的PCollection
  • 尝试:即使将Beam Schema所有字段设为String类型,仍出现相同异常
  • 环境:Java 11 + Maven 3.8.5 + Apache Beam Java SDK 2.41.0;Java 1.8 + Beam 2.40.0环境下问题复现

代码示例

public class beamRowPractise {

    public static void main(String[] args){

        PipelineOptions opts = PipelineOptionsFactory.create();
        opts.setRunner(DirectRunner.class);
        Pipeline p = Pipeline.create(opts);
        PCollection<String> pc1 = p.apply(TextIO.read().from("data/indata.csv"));
        PCollection<Row> pc2 = pc1.apply(MapElements.via(new mapString())).setRowSchema(getSchema()) ;
        System.out.println(pc2.getSchema().toString());
        p.run();
        }
    public static class mapString extends SimpleFunction<String, Row> {
        @Override

        public  Row apply(String record){
            String arr[] = record.split(",");

            Row.Builder row = Row.withSchema(getSchema()) ;

            row.withFieldValue("name",arr[0]);
            row.withFieldValue("id1",arr[1]);
            row.withFieldValue("id2",arr[2]);
            row.withFieldValue("id3",arr[3]);
            row.withFieldValue("id4",arr[4]);

            return  row.build();

        }
    }

    public  static  Schema getSchema() {
        org.apache.beam.sdk.schemas.Schema.Builder typed_schema_builder = org.apache.beam.sdk.schemas.Schema.builder();
        typed_schema_builder.addField("name", org.apache.beam.sdk.schemas.Schema.FieldType.STRING);
        typed_schema_builder.addField("id1", Schema.FieldType.INT64);
        typed_schema_builder.addField("id2", org.apache.beam.sdk.schemas.Schema.FieldType.INT64);
        typed_schema_builder.addField("id3", org.apache.beam.sdk.schemas.Schema.FieldType.INT64);
        typed_schema_builder.addField("id4", org.apache.beam.sdk.schemas.Schema.FieldType.INT64);
        org.apache.beam.sdk.schemas.Schema typed_beam_schema = typed_schema_builder.build();
        org.apache.beam.sdk.schemas.Schema schema = typed_beam_schema;
        return  schema;
    }
}

报错信息

Exception in thread "main" org.apache.beam.sdk.Pipeline$PipelineExecutionException: java.lang.IllegalStateException
    at org.apache.beam.runners.direct.DirectRunner$DirectPipelineResult.waitUntilFinish(DirectRunner.java:374)
    at org.apache.beam.runners.direct.DirectRunner$DirectPipelineResult.waitUntilFinish(DirectRunner.java:342)
    at org.apache.beam.runners.direct.DirectRunner.run(DirectRunner.java:218)
    at org.apache.beam.runners.direct.DirectRunner.run(DirectRunner.java:67)
    at org.apache.beam.sdk.Pipeline.run(Pipeline.java:323)
    at org.apache.beam.sdk.Pipeline.run(Pipeline.java:309)
    at com.bhargav.beamFirst.beamRowPractise.main(beamRowPractise.java:25)
Caused by: java.lang.IllegalStateException
    at org.apache.beam.vendor.guava.v26_0_jre.com.google.common.base.Preconditions.checkState(Preconditions.java:491)
    at org.apache.beam.sdk.coders.RowCoderGenerator$EncodeInstruction.encodeDelegate(RowCoderGenerator.java:313)
    at org.apache.beam.sdk.coders.Coder$ByteBuddy$hZNCN9ub.encode(Unknown Source)
    at org.apache.beam.sdk.coders.Coder$ByteBuddy$hZNCN9ub.encode(Unknown Source)
    at org.apache.beam.sdk.schemas.SchemaCoder.encode(SchemaCoder.java:124)
    at org.apache.beam.sdk.coders.Coder.encode(Coder.java:136)
    at org.apache.beam.sdk.util.CoderUtils.encodeToSafeStream(CoderUtils.java:86)
    at org.apache.beam.sdk.util.CoderUtils.encodeToByteArray(CoderUtils.java:70)
    at org.apache.beam.sdk.util.CoderUtils.encodeToByteArray(CoderUtils.java:55)
    at org.apache.beam.sdk.util.CoderUtils.clone(CoderUtils.java:168)
    at org.apache.beam.sdk.util.MutationDetectors$CodedValueMutationDetector.<init>(MutationDetectors.java:118)
    at org.apache.beam.sdk.util.MutationDetectors.forValueWithCoder(MutationDetectors.java:49)
    at org.apache.beam.runners.direct.ImmutabilityCheckingBundleFactory$ImmutabilityEnforcingBundle.add(ImmutabilityCheckingBundleFactory.java:115)
    at org.apache.beam.runners.direct.ParDoEvaluator$BundleOutputManager.output(ParDoEvaluator.java:305)
    at org.apache.beam.repackaged.direct_java.runners.core.SimpleDoFnRunner.outputWindowedValue(SimpleDoFnRunner.java:275)
    at org.apache.beam.repackaged.direct_java.runners.core.SimpleDoFnRunner.access$900(SimpleDoFnRunner.java:85)
    at org.apache.beam.repackaged.direct_java.runners.core.SimpleDoFnRunner$DoFnProcessContext.output(SimpleDoFnRunner.java:423)
    at org.apache.beam.sdk.transforms.DoFnOutputReceivers$WindowedContextOutputReceiver.output(DoFnOutputReceivers.java:76)
    at org.apache.beam.sdk.transforms.MapElements$1.processElement(MapElements.java:142)

Process finished with exit code 1

问题原因与解决方案

核心问题

代码中Schema定义id1-id4为INT64类型,但在mapString的apply方法中,直接将CSV分割后的String类型值(arr[1]-arr[4])赋值给这些字段,类型不匹配导致Row编码时触发异常。即使后续将Schema全部改为String类型,若代码未同步调整或存在其他类型不一致问题,仍会报错。

解决方案

有两种可行的修正方向:

方案1:将String类型转换为Long类型(匹配INT64字段)

修改mapString类的apply方法,把分割后的String转为Long:

public static class mapString extends SimpleFunction<String, Row> {
    @Override
    public Row apply(String record){
        String arr[] = record.split(",");
        Row.Builder row = Row.withSchema(getSchema()) ;

        row.withFieldValue("name", arr[0]);
        // 将String转为Long,适配INT64类型字段
        row.withFieldValue("id1", Long.parseLong(arr[1]));
        row.withFieldValue("id2", Long.parseLong(arr[2]));
        row.withFieldValue("id3", Long.parseLong(arr[3]));
        row.withFieldValue("id4", Long.parseLong(arr[4]));

        return row.build();
    }
}

方案2:将Schema字段全部改为String类型

调整getSchema方法,把所有字段设为STRING类型,同时确保赋值的类型一致:

public static Schema getSchema() {
    org.apache.beam.sdk.schemas.Schema.Builder typed_schema_builder = org.apache.beam.sdk.schemas.Schema.builder();
    typed_schema_builder.addField("name", org.apache.beam.sdk.schemas.Schema.FieldType.STRING);
    typed_schema_builder.addField("id1", org.apache.beam.sdk.schemas.Schema.FieldType.STRING);
    typed_schema_builder.addField("id2", org.apache.beam.sdk.schemas.Schema.FieldType.STRING);
    typed_schema_builder.addField("id3", org.apache.beam.sdk.schemas.Schema.FieldType.STRING);
    typed_schema_builder.addField("id4", org.apache.beam.sdk.schemas.Schema.FieldType.STRING);
    return typed_schema_builder.build();
}

验证

修正后运行代码,即可成功将CSV的String行转换为Row类型的PCollection,不再抛出异常。


内容的提问来源于stack exchange,提问作者Bhargav Velisetti

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 18:15:49