apache / apache/beam

[Streaming][PubSub Lite][DataflowRunner] PubSub Lite IO doesn't sink message to PubSub Lite topic

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

Description

We are currently using PubSub Lite IO with Dataflow Runner.

Our Beam job is on streaming mode.

The read from a PubSub Lite subscription works correctly.

The sink to a PubSub topic doesn't work with the runner. 

When we take a look  on the job graph for the PubSubLite Write step (transform : org.apache.beam.sdk.io.gcp.pubsublite.PubsubLiteSink ) in the Google Cloud Console we don't see any writes.

When we check the topic we don't see any outputs.

Code works well on 2.27, 2.28, 2.29 Beam version.

Here is the code we used to do the check on version:
* 2.27 
* 2.28
* 2.29
* 2.30
* 2.31
* 2.33
* 2.33

 
```

// PubSubStreamingWriteJobOptions options =
PipelineOptionsFactory.fromArgs(args).withValidation().as(PubSubStreamingWriteJobOptions.class);

options.setStreaming(true);

//set
up file system
FileSystems.setDefaultPipelineOptions(options);

TopicPath topicPath = TopicPath.newBuilder()

.setProject(ProjectId.of("[PROJECT ID]"))
.setLocation(CloudZone.of(CloudRegion.of("[REGION]"),
"[ZONE CHAR]"))
.setName(TopicName.of("[TOPIC ID]"))
.build();

PublisherOptions
publisherOptions =
PublisherOptions.newBuilder()
.setTopicPath(topicPath)

.build();

Pipeline pipeline = Pipeline.create(options);

pipeline.apply(TextIO.read()

.from("gs://[BUCKET]/[OBJECT_PREFIX]*")
.watchForNewFiles(

Duration.standardMinutes(1),
Watch.Growth.afterTimeSinceNewOutput(Duration.standardHours(1))))

.apply(CREATE_PUB_SUB_LITE_MESSAGE_STEP, MapElements.into(TypeDescriptor.of(PubSubMessage.class)).via(file
-> {
Instant instant = Instant.now();
Message message =

Message.builder()
.setData(ByteString.copyFromUtf8("message " + file))

.setEventTime(Timestamp.newBuilder()

.setNanos(instant.getNano())
.setSeconds(instant.getEpochSecond())

.build())
.build();

return
message.toProto();
}))
.apply(SINK_PUB_SUB_LITE_MESSAGES_STEP, PubsubLiteIO.write(publisherOptions));
pipeline.run();

```

Can you help us to found the issue and fix the Beam version please?

 

Best regards,

David Duarte

 

 

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

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.