Portable Flink runner does not support flatten with different input/output coders in different languages (FlattenTest.testFlattenWithDifferentInputAndOutputCoders2)
- Dominant language
- Java
- Stars
- 8.7k
- Forks
- 4.7k
- Avg merge
- 1d 20h
- Merged PRs (30d)
- 196
Description
Both beam_PostCommit_Java_PVR_Flink_Batch and beam_PostCommit_Java_PVR_Flink_Streaming are failing newly added test org.apache.beam.sdk.transforms.FlattenTest.testFlattenWithDifferentInputAndOutputCoders2.
SEVERE: Error in task code: CHAIN MapPartition (MapPartition at [6]{Values, FlatMapElements, PAssert$0}) -\> FlatMap (FlatMap at ExtractOutput[0]) -\> Map (Key Extractor) -\> GroupCombine (GroupCombine at GroupCombine: PAssert$0/GroupGlobally/GatherAllOutputs/GroupByKey) -\> Map (Key Extractor) (2/2) java.lang.ClassCastException: org.apache.beam.sdk.values.KV cannot be cast to [B
at org.apache.beam.sdk.coders.ByteArrayCoder.encode(ByteArrayCoder.java:41)
at org.apache.beam.sdk.coders.LengthPrefixCoder.encode(LengthPrefixCoder.java:56)
at org.apache.beam.sdk.coders.Coder.encode(Coder.java:136)
at org.apache.beam.sdk.util.WindowedValue$FullWindowedValueCoder.encode(WindowedValue.java:590)
at org.apache.beam.sdk.util.WindowedValue$FullWindowedValueCoder.encode(WindowedValue.java:581)
at org.apache.beam.sdk.util.WindowedValue$FullWindowedValueCoder.encode(WindowedValue.java:541)
at org.apache.beam.sdk.fn.data.BeamFnDataSizeBasedBufferingOutboundObserver.accept(BeamFnDataSizeBasedBufferingOutboundObserver.java:109)
at org.apache.beam.runners.fnexecution.control.SdkHarnessClient$CountingFnDataReceiver.accept(SdkHarnessClient.java:667)
at org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunction.processElements(FlinkExecutableStageFunction.java:271)
at org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunction.mapPartition(FlinkExecutableStageFunction.java:203)
at org.apache.flink.runtime.operators.MapPartitionDriver.run(MapPartitionDriver.java:103)
at org.apache.flink.runtime.operators.BatchTask.run(BatchTask.java:504)
at org.apache.flink.runtime.operators.BatchTask.invoke(BatchTask.java:369)
at org.apache.flink.runtime.taskmanager.Task.doRun(Task.java:705)
at org.apache.flink.runtime.taskmanager.Task.run(Task.java:530)
at java.lang.Thread.run(Thread.java:748)
Imported from Jira [BEAM-10016](https://issues.apache.org/jira/browse/BEAM-10016). Original Jira may contain additional context.
Reported by: ibzib.
This issue has child subcomponents which were not migrated over. See the original Jira for more information.
Contributor guide
Assessment
This issue has not been assessed yet.