apache / apache/beam

[Bug]: Potential performance regression in KafkaIO and schema registry

Open
#26,262 1 comment 0 reactions 0 assignees View on GitHub
bug java kafka P2
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
1d 20h
Merged PRs (30d)
196

Description

### What needs to happen?

From [email thread](https://lists.apache.org/thread/97dkw2rc4tf1f0gfmnfgdfgdg7nvbwmw):

"I am trying to understand the effect of schema registry on our pipeline's performance. In order to do sowe created a very simple pipeline that reads from kafka, runs a simple transformation of adding new field and writes of kafka. the messages are in avro format

I ran this pipeline with 3 different options on same configuration : 1 kafka partition, 1 task manager, 1 slot, 1 parallelism:

* when i used apicurio as the schema registry i was able to process only 2000 messages per second
* when i used confluent schema registry i was able to process 7000 messages per second
* when I did not use any schema registry and used plain avro deserializer/serializer i was able to process 30K messages per second.

```
KafkaIO.read()
.withBootstrapServers(bootstrapServers)
.withTopic(topic)
.withConsumerConfigUpdates(Map.ofEntries(
Map.entry("schema.registry.url", registryURL),
Map.entry(ConsumerConfig.GROUP_ID_CONFIG, consumerGroup+ UUID.randomUUID()),
))
.withKeyDeserializer(StringDeserializer.class)
.withValueDeserializerAndCoder((Class) io.confluent.kafka.serializers.KafkaAvroDeserializer.class, AvroCoder.of(avroClass));
```

I have made the suggested change and used `ConfluentSchemaRegistryDeserializerProvider`
the results are slightly better.. average of 8000 msg/sec "

We need to investigate and find out the cause of this performance issue.

### Issue Priority

Priority: 2 (default / most normal work should be filed as P2)

### Issue Components

- [ ] Component: Python SDK
- [X] Component: Java SDK
- [ ] Component: Go SDK
- [ ] Component: Typescript SDK
- [ ] Component: IO connector
- [ ] Component: Beam examples
- [ ] Component: Beam playground
- [ ] Component: Beam katas
- [ ] Component: Website
- [ ] Component: Spark Runner
- [ ] Component: Flink Runner
- [ ] Component: Samza Runner
- [ ] Component: Twister2 Runner
- [ ] Component: Hazelcast Jet Runner
- [ ] Component: Google Cloud Dataflow Runner

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.