apache / apache/beam

[Bug]: BundleProcessorCache calls shutdown on timed out DoFns inline, blocking unrelated processing

Open
#26,599 0 comments 0 reactions 0 assignees View on GitHub
awaiting triage bug java P2
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

Open the contributing 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

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.