FlinkRuner: Pipeline using KafkaIO seems not be able to terminate
- 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
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