apache / apache/beam

FlinkRuner: Pipeline using KafkaIO seems not be able to terminate

Open
#20,914 4 comments 0 reactions 0 assignees View on GitHub
bug io java kafka P3
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
2d 5h
Merged PRs (30d)
204

Description

2021-03-25 14:21:27,210 WARN org.apache.beam.sdk.io.kafka.KafkaUnboundedReader [] - Reader-0: Unexpected
java.lang.InterruptedException: null
at java.util.concurrent.SynchronousQueue.poll(Unknown Source) ~[?:?]
at org.apache.beam.sdk.io.kafka.KafkaUnboundedReader.nextBatch(KafkaUnboundedReader.java:584) ~[blob_p-e4f6919ea552b3197dcb3d58dab934634011ea1d-f375010a962934be4febeb9924152473:?]
at org.apache.beam.sdk.io.kafka.KafkaUnboundedReader.advance(KafkaUnboundedReader.java:214) ~[blob_p-e4f6919ea552b3197dcb3d58dab934634011ea1d-f375010a962934be4febeb9924152473:?]
at org.apache.beam.sdk.io.Read$UnboundedSourceAsSDFWrapperFn$UnboundedSourceAsSDFRestrictionTracker.tryClaim(Read.java:841) ~[blob_p-cfeb2021150481a2d2069a38f7abd261d89645c3-b56abda9593e69f19cd5f833293fbd4f:?]
at org.apache.beam.sdk.io.Read$UnboundedSourceAsSDFWrapperFn$UnboundedSourceAsSDFRestrictionTracker.tryClaim(Read.java:781) ~[blob_p-cfeb2021150481a2d2069a38f7abd261d89645c3-b56abda9593e69f19cd5f833293fbd4f:?]
at org.apache.beam.sdk.fn.splittabledofn.RestrictionTrackers$RestrictionTrackerObserver.tryClaim(RestrictionTrackers.java:59) ~[blob_p-bcbef6ab6822495d6ebec31ea6f945a2703e27a2-beed43efff1368b0a55e8950c1c78418:?]
at org.apache.beam.sdk.io.Read$UnboundedSourceAsSDFWrapperFn.processElement(Read.java:537) ~[blob_p-cfeb2021150481a2d2069a38f7abd261d89645c3-b56abda9593e69f19cd5f833293fbd4f:?]
at org.apache.beam.sdk.io.Read$UnboundedSourceAsSDFWrapperFn$DoFnInvoker.invokeProcessElement(Unknown Source) ~[?:?]
at org.apache.beam.runners.core.OutputAndTimeBoundedSplittableProcessElementInvoker.invokeProcessElement(OutputAndTimeBoundedSplittableProcessElementInvoker.java:123) ~[blob_p-a8a186bf74efa331ed1d0a699183e46f7fe5e71f-7d0e42800d9b187a9ca3e2acd5574aa3:?]
at org.apache.beam.runners.core.SplittableParDoViaKeyedWorkItems$ProcessFn.processElement(SplittableParDoViaKeyedWorkItems.java:523) ~[blob_p-a8a186bf74efa331ed1d0a699183e46f7fe5e71f-7d0e42800d9b187a9ca3e2acd5574aa3:?]
at org.apache.beam.runners.core.SplittableParDoViaKeyedWorkItems$ProcessFn$DoFnInvoker.invokeProcessElement(Unknown Source) ~[?:?]
at org.apache.beam.runners.core.SimpleDoFnRunner.invokeProcessElement(SimpleDoFnRunner.java:232) ~[blob_p-a8a186bf74efa331ed1d0a699183e46f7fe5e71f-7d0e42800d9b187a9ca3e2acd5574aa3:?]
at org.apache.beam.runners.core.SimpleDoFnRunner.processElement(SimpleDoFnRunner.java:188) ~[blob_p-a8a186bf74efa331ed1d0a699183e46f7fe5e71f-7d0e42800d9b187a9ca3e2acd5574aa3:?]
at org.apache.beam.runners.flink.metrics.DoFnRunnerWithMetricsUpdate.processElement(DoFnRunnerWithMetricsUpdate.java:62) ~[blob_p-f2532293ec647e3493b3c93016324a1bd4a24416-94f2748503eb0c25216d37d822b6e6cf:?]

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

Contributor guide

Open the contributing guide

Research direction

Start with the reported stack trace and inspect KafkaUnboundedReader.java, especially nextBatch and advance, along with the Flink runner path named in the issue. Reproduce the KafkaIO pipeline termination problem and determine the expected behavior when the pipeline is stopped; done means the pipeline terminates without the reported unexpected interruption.

Written by the indexing model from the issue text.

Assessment

Tech stack
java, kafka
Domain
backend, stream-processing
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.