[Bug]: BundleProcessorCache calls shutdown on timed out DoFns inline, blocking unrelated processing
- Dominant language
- Java
- Stars
- 8.7k
- Forks
- 4.7k
- Avg merge
- 1d 20h
- Merged PRs (30d)
- 196
Description
### What happened?
Stack traces showing how shutdown blocks other processing:
Lots of new bundles unable to get a processor:
```
java.base@11.0.17.0.101/jdk.internal.misc.Unsafe.park(Native Method)
java.base@11.0.17.0.101/java.util.concurrent.locks.LockSupport.park(LockSupport.java:194)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.util.concurrent.AbstractFuture.get(AbstractFuture.java:502)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.util.concurrent.AbstractFuture$TrustedFuture.get(AbstractFuture.java:83)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.util.concurrent.Uninterruptibles.getUninterruptibly(Uninterruptibles.java:196)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$LoadingValueReference.waitForValue(LocalCache.java:3581)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$Segment.waitForLoadingValue(LocalCache.java:2174)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$Segment.lockedGetOrLoad(LocalCache.java:2161)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$Segment.get(LocalCache.java:2044)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache.get(LocalCache.java:3952)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache.getOrLoad(LocalCache.java:3974)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$LocalLoadingCache.get(LocalCache.java:4958)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$LocalLoadingCache.getUnchecked(LocalCache.java:4964)
app//org.apache.beam.fn.harness.control.ProcessBundleHandler$BundleProcessorCache.get(ProcessBundleHandler.java:965)
app//org.apache.beam.fn.harness.control.ProcessBundleHandler.processBundle(ProcessBundleHandler.java:507)
app//org.apache.beam.fn.harness.FnHarness$$Lambda$201/0x0000000800386040.apply(Unknown Source)
app//org.apache.beam.fn.harness.control.BeamFnControlClient.delegateOnInstructionRequestType(BeamFnControlClient.java:151)
app//org.apache.beam.fn.harness.control.BeamFnControlClient$InboundObserver.lambda$onNext$0(BeamFnControlClient.java:116)
app//org.apache.beam.fn.harness.control.BeamFnControlClient$InboundObserver$$Lambda$699/0x0000000800766840.run(Unknown Source)
java.base@11.0.17.0.101/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
java.base@11.0.17.0.101/java.util.concurrent.FutureTask.run(FutureTask.java:264)
app//org.apache.beam.sdk.util.UnboundedScheduledExecutorService$ScheduledFutureTask.run(UnboundedScheduledExecutorService.java:163)
java.base@11.0.17.0.101/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
java.base@11.0.17.0.101/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
java.base@11.0.17.0.101/java.lang.Thread.run(Thread.java:829)
```
Because shutdown is happening for a (possibly) different Dofn and is slow:
```
java.base@11.0.17.0.101/jdk.internal.misc.Unsafe.park(Native Method)
java.base@11.0.17.0.101/java.util.concurrent.locks.LockSupport.parkNanos(LockSupport.java:234)
java.base@11.0.17.0.101/java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.awaitNanos(AbstractQueuedSynchronizer.java:2123)
java.base@11.0.17.0.101/java.util.concurrent.ThreadPoolExecutor.awaitTermination(ThreadPoolExecutor.java:1454)
app//org.apache.beam.sdk.util.UnboundedScheduledExecutorService.awaitTermination(UnboundedScheduledExecutorService.java:363)
app//com.google.api.gax.rpc.Watchdog.awaitTermination(Watchdog.java:158)
app//com.google.api.gax.core.BackgroundResourceAggregation.awaitTermination(BackgroundResourceAggregation.java:82)
app//com.google.cloud.bigquery.storage.v1.stub.GrpcBigQueryWriteStub.awaitTermination(GrpcBigQueryWriteStub.java:330)
app//com.google.cloud.bigquery.storage.v1.BigQueryWriteClient.awaitTermination(BigQueryWriteClient.java:896)
app//org.apache.beam.sdk.io.gcp.bigquery.BigQueryServicesImpl$DatasetServiceImpl.close(BigQueryServicesImpl.java:1440)
app//org.apache.beam.sdk.io.gcp.bigquery.StorageApiWriteUnshardedRecords$WriteRecordsDoFn.teardown(StorageApiWriteUnshardedRecords.java:784)
app//org.apache.beam.sdk.io.gcp.bigquery.StorageApiWriteUnshardedRecords$WriteRecordsDoFn$DoFnInvoker.invokeTeardown(Unknown Source)
app//org.apache.beam.fn.harness.FnApiDoFnRunner.tearDown(FnApiDoFnRunner.java:1778)
app//org.apache.beam.fn.harness.FnApiDoFnRunner$$Lambda$782/0x000000080091ac40.run(Unknown Source)
app//org.apache.beam.fn.harness.control.ProcessBundleHandler$BundleProcessor.shutdown(ProcessBundleHandler.java:1179)
app//org.apache.beam.fn.harness.control.ProcessBundleHandler$BundleProcessorCache.lambda$new$0(ProcessBundleHandler.java:929)
app//org.apache.beam.fn.harness.control.ProcessBundleHandler$BundleProcessorCache$$Lambda$989/0x0000000800cc1840.accept(Unknown Source)
java.base@11.0.17.0.101/java.util.concurrent.ConcurrentLinkedQueue.forEachFrom(ConcurrentLinkedQueue.java:1037)
java.base@11.0.17.0.101/java.util.concurrent.ConcurrentLinkedQueue.forEach(ConcurrentLinkedQueue.java:1054)
app//org.apache.beam.fn.harness.control.ProcessBundleHandler$BundleProcessorCache.lambda$new$1(ProcessBundleHandler.java:927)
app//org.apache.beam.fn.harness.control.ProcessBundleHandler$BundleProcessorCache$$Lambda$197/0x000000080029a440.onRemoval(Unknown Source)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache.processPendingNotifications(LocalCache.java:1809)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$Segment.runUnlockedCleanup(LocalCache.java:3462)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$Segment.postWriteCleanup(LocalCache.java:3438)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$Segment.lockedGetOrLoad(LocalCache.java:2145)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$Segment.get(LocalCache.java:2044)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache.get(LocalCache.java:3952)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache.getOrLoad(LocalCache.java:3974)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$LocalLoadingCache.get(LocalCache.java:4958)
app//org.apache.beam.vendor.guava.v26_0_jre.com.google.common.cache.LocalCache$LocalLoadingCache.getUnchecked(LocalCache.java:4964)
app//org.apache.beam.fn.harness.control.ProcessBundleHandler$BundleProcessorCache.get(ProcessBundleHandler.java:965)
app//org.apache.beam.fn.harness.control.ProcessBundleHandler.processBundle(ProcessBundleHandler.java:507)
app//org.apache.beam.fn.harness.FnHarness$$Lambda$201/0x0000000800386040.apply(Unknown Source)
app//org.apache.beam.fn.harness.control.BeamFnControlClient.delegateOnInstructionRequestType(BeamFnControlClient.java:151)
app//org.apache.beam.fn.harness.control.BeamFnControlClient$InboundObserver.lambda$onNext$0(BeamFnControlClient.java:116)
app//org.apache.beam.fn.harness.control.BeamFnControlClient$InboundObserver$$Lambda$699/0x0000000800766840.run(Unknown Source)
java.base@11.0.17.0.101/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)
java.base@11.0.17.0.101/java.util.concurrent.FutureTask.run(FutureTask.java:264)
app//org.apache.beam.sdk.util.UnboundedScheduledExecutorService$ScheduledFutureTask.run(UnboundedScheduledExecutorService.java:163)
java.base@11.0.17.0.101/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)
java.base@11.0.17.0.101/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)
java.base@11.0.17.0.101/java.lang.Thread.run(Thread.java:829)
```
### Issue Priority
Priority: 2 (default / most bugs 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
Research direction
Start in ProcessBundleHandler.java, especially BundleProcessorCache and BundleProcessor.shutdown, using the stack traces to trace cache removal and timed-out DoFn teardown. The fix should prevent a slow shutdown from blocking unrelated bundle processing; verify the behavior with the existing Java SDK test suite, since no specific test is named.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- java
- Domain
- backend, distributed-systems
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 35/100