apache / apache/beam

Beam spark runner not working properly with kafka

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

Description

We are running a beam stream processing job on a spark runner, which reads from a kafka topic using kerberos authentication. We are using java-io-kafka v2.4.0 to read from kafka topic in the pipeline. The issue is that the kafkaIO client is continuously creating a new kafka consumer with specified config, doing kerberos login every time. Also, there are spark streaming jobs which get spawned for the unbounded source, every second or so even when there is no data in the kafka topic. Log has these jobs-

INFO SparkContext: Starting job: DStream@SparkUnboumdedSource.java:172

We can see in the logs

INFO MicrobatchSource: No cached reader found for split: [org.apache.beam.sdk.io.kafka.KafkaUnboundedSource@2919a728]. Creating new reader at checkpoint mark...

And then it creates new consumer doing fresh kerberos login, which is creating issues.

We are unsure of what should be correct behavior here and why so many spark streaming jobs are getting created. We tried the beam code with flink runner and did not find this issue there. Can someone point to the correct settings for using unbounded kafka source with spark runner using beam? 

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

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.