apache / apache/beam

PubsubLiteIO.read() fails in DataflowRunner v1 due to lack of BundleFinalizer support

Open
#20,905 0 comments 0 reactions 0 assignees View on GitHub
dataflow extensions gcp java new feature P3 pubsub
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

Reading from a Pub/Sub Lite subscription using PubsubLiteIO fails with DataflowRunner. It works in DirectRunner. 

```

import org.apache.beam.sdk.io.gcp.pubsublite.PubsubLiteIO;
//..
pipeline
.apply("Read
From Lite", PubsubLiteIO.read(subscriberOpitons))
.apply("Convert and print", MapElements.into(TypeDescriptors.strings()).via(

(SequencedMessage sequencedMessage) -> {
String data = sequencedMessage.getMessage().getData().toStringUtf8();

LOG.info("Received: " + data);
return data;
}
));

```

```

java.lang.UnsupportedOperationException: BundleFinalizer unsupported by non-portable Dataflow.
at
org.apache.beam.runners.dataflow.worker.SplittableProcessFnFactory$SplittableDoFnRunnerFactory.lambda$createRunner$2(SplittableProcessFnFactory.java:170)
at
org.apache.beam.runners.dataflow.worker.repackaged.org.apache.beam.runners.core.OutputAndTimeBoundedSplittableProcessElementInvoker$1.bundleFinalizer(OutputAndTimeBoundedSplittableProcessElementInvoker.java:195)
at
org.apache.beam.sdk.io.gcp.pubsublite.PerSubscriptionPartitionSdf$DoFnInvoker.invokeProcessElement(Unknown
Source)
at org.apache.beam.runners.dataflow.worker.repackaged.org.apache.beam.runners.core.OutputAndTimeBoundedSplittableProcessElementInvoker.invokeProcessElement(OutputAndTimeBoundedSplittableProcessElementInvoker.java:123)
at
org.apache.beam.runners.dataflow.worker.repackaged.org.apache.beam.runners.core.SplittableParDoViaKeyedWorkItems$ProcessFn.processElement(SplittableParDoViaKeyedWorkItems.java:523)

```

 

Imported from Jira [BEAM-12085](https://issues.apache.org/jira/browse/BEAM-12085). Original Jira may contain additional context.
Reported by: tianzi.

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.