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
相关产品推荐
相关产品推荐

