apache / apache/beam

Reduce Pub/Sub publishing latency

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

Description

The current implementation of the `PubsubUnboundedSink` uses a global window with a trigger on a fixed batch size of 1000 elements or a processing timespan of 2 seconds. After that, a random sharding of 100 is applied via a `GroupByKey` transform. The result is then pushed into a `DoFn` which performs the actual publishing step.

In case of low-latency (10s or 100s of milliseconds), this logic is quite bad, because it leads to a latency of  around 1.2 seconds, introduced by the transform steps described above.

There are several possibilities to improve the Pub/Sub sink, for example:

Let the upper parameters be configured via `PipelineOptions:`
* `pubsubBatchSize`: Approx. maximum number of elements in a Pub/Sub publishing batch
* `pubsubDelayThreshold`: Max. processing time duration before firing the sharding window
* `pubsubShardCount`: The number of shards to create before publishing

This would allow tweaking of the Pub/Sub sink for different scenarious of throughput and message size in the pipeline.

However, if the throughput is small (< 100 element/s), this mechanism is still quite slow. If we take a look at the Java client at `com.google.cloud:google-cloud-pubsub`, the `Publisher` class supports a wide range of options to optimize its batching behaviour. This would allow not to rely on a window with group by key functionality and let the publisher itself handle the batching.

Consider the following `DoFn` for publishing messages to Pub/Sub using that client:
```

class PublishFn extends DoFn {
private transient Publisher publisher;

private final ValueProvider topicPath;

public PublishFn(final ValueProvider
topicPath) {
this.topicPath = topicPath;
}

@Setup
public void setup() throws
IOException {
publisher = Publisher.defaultBuilder(TopicName.parse(topicPath.get()))

.setBatchingSettings(BatchingSettings.newBuilder()
.setRequestByteThreshold(40000L)

.setElementCountThreshold(1000L)
.setDelayThreshold(Duration.ofMillis(50))

.build())
.build();
}

@ProcessElement
public
void processElement(final ProcessContext context) {
publisher.publish(context.element());

}

@Teardown
public void teardown() throws Exception {
publisher.shutdown();

}

@Override
public void populateDisplayData(final DisplayData.Builder builder) {

builder.add(DisplayData.item("topic", topicPath));
}
}

```

In small test, this resulted in a publish latency of around 50 – 70 ms instead of 1000 – 1200 with the original `PubsubUnboundedSink`.

I can understand, that the windowing mechanism could lead to better performance and throughput in a scenario with a high number of elements per second. However, it would be nice to enable a "low-latency-mode" using the provided code as an example.

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

Contributor guide

Open the contributing guide

Research direction

Start by locating the PubsubUnboundedSink and tracing its global window, GroupByKey, and publishing DoFn path. Read the Java google-cloud-pubsub Publisher batching API and compare it with the proposed PublishFn and configurable PipelineOptions. Done should mean an agreed low-latency mode or batching design with latency behavior validated against the current sink.

Written by the indexing model from the issue text.

Assessment

Tech stack
google-cloud, java
Domain
distributed-systems, stream-processing
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
30/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.