apache / apache/beam

[BUG]: Java `Wait.OnSignal` does not set output coder.

Open
#23,084 1 comment 2 reactions 0 assignees View on GitHub
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

### What happened?

I spent a significant amount of time where the coder is missing or type is being erased, and located a single point.

The issue is a combination of `ProtoCoder` and `MapElements`, as
- `ProtoCoder` does not have `getEncodedTypeDescriptor()`, so it returns just `TypeDescriptor` instead of actual type.
- `MapElements` only takes `TypeDescriptor` as a parameter, and try to re-infer the coder.

Proposed Fix:
- set input coder to output coder in `Wait.OnSignal` after `MapElements`

How to reproduce.

```
import org.apache.beam.sdk.options.PipelineOptions.CheckEnabled;
import org.apache.beam.sdk.testing.TestPipeline;
import org.apache.beam.sdk.transforms.Create;
import org.apache.beam.sdk.transforms.Wait;
import org.junit.Rule;
import org.junit.Test;

import com.google.protobuf.Empty;

public class EmptyMessageTest {
@Rule
public final TestPipeline pipeline = TestPipeline.create();

@Test
public void test() {
pipeline.getOptions().setStableUniqueNames(CheckEnabled.OFF);

var emptyMessage = pipeline.apply(Create.of(Empty.getDefaultInstance()));
var singleNull = pipeline.apply(Create.of((Void) null));

emptyMessage.apply(Wait.on(singleNull));
pipeline.run();
}
}
```

```
Caused by: java.lang.NoSuchMethodException: com.google.protobuf.Message.getDefaultInstance()
at java.base/java.lang.Class.getMethod(Class.java:2108)
at org.apache.beam.sdk.extensions.protobuf.ProtoCoder.getParser(ProtoCoder.java:297)
at org.apache.beam.sdk.extensions.protobuf.ProtoCoder.decode(ProtoCoder.java:205)
at org.apache.beam.sdk.extensions.protobuf.ProtoCoder.decode(ProtoCoder.java:108)
at org.apache.beam.sdk.util.CoderUtils.decodeFromSafeStream(CoderUtils.java:118)
at org.apache.beam.sdk.util.CoderUtils.decodeFromByteArray(CoderUtils.java:101)
at org.apache.beam.sdk.util.CoderUtils.decodeFromByteArray(CoderUtils.java:95)
at org.apache.beam.sdk.util.CoderUtils.clone(CoderUtils.java:144)
at org.apache.beam.sdk.util.MutationDetectors$CodedValueMutationDetector.(MutationDetectors.java:118)
at org.apache.beam.sdk.util.MutationDetectors.forValueWithCoder(MutationDetectors.java:49)
```

### Issue Priority
Priority: 3

### Issue Component
Component: sdk-java-core

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.