apache / apache/beam

Dataflow streaming runner fails SDF testWindowedSideInputWithCheckpoints

Open
#18,457 0 comments 0 reactions 0 assignees View on GitHub
bug dataflow P3 runners
Dominant language
Java
Stars
8.7k
Forks
4.7k
Avg merge
2d 2h
Merged PRs (30d)
205

Description

When testing https://github.com/apache/beam/pull/1898

INFO: 2017-06-20T00:04:36.310Z: (22acf084264efa8): Executing operation PAssert$7/GroupGlobally/GroupDummyAndContents/ReadStream****PAssert$7/GroupGlobally/GroupDummyAndContents/MergeBuckets****PAssert$7/GroupGlobally/Values/Values/Map****PAssert$7/GroupGlobally/ParDo(Concat)****PAssert$7/GetPane/Map****PAssert$7/RunChecks****PAssert$7/VerifyAssertions/ParDo(DefaultConclude)
Jun 19, 2017 5:04:59 PM org.apache.beam.runners.dataflow.util.MonitoringUtil$LoggingHandler process
SEVERE: 2017-06-20T00:04:58.641Z: (6a9da772bd8097d0): java.lang.RuntimeException: Exception while fetching side input:
com.google.cloud.dataflow.worker.runners.worker.StateFetcher.fetchSideInput(StateFetcher.java:182)
com.google.cloud.dataflow.worker.runners.worker.StreamingModeExecutionContext.fetchSideInput(StreamingModeExecutionContext.java:211)
com.google.cloud.dataflow.worker.runners.worker.StreamingModeExecutionContext.access$700(StreamingModeExecutionContext.java:70)
com.google.cloud.dataflow.worker.runners.worker.StreamingModeExecutionContext$StepContext.issueSideInputFetch(StreamingModeExecutionContext.java:572)
com.google.cloud.dataflow.worker.runners.worker.StreamingModeExecutionContext$UserStepContext.issueSideInputFetch(StreamingModeExecutionContext.java:624)
com.google.cloud.dataflow.worker.runners.worker.StreamingSideInputFetcher.getReadyWindows(StreamingSideInputFetcher.java:136)
com.google.cloud.dataflow.worker.runners.worker.StreamingSideInputDoFnRunner.startBundle(StreamingSideInputDoFnRunner.java:51)
com.google.cloud.dataflow.worker.runners.worker.SimpleParDoFn.reallyStartBundle(SimpleParDoFn.java:222)
com.google.cloud.dataflow.worker.runners.worker.SimpleParDoFn.startBundle(SimpleParDoFn.java:150)
com.google.cloud.dataflow.worker.runners.worker.ForwardingParDoFn.startBundle(ForwardingParDoFn.java:36)
com.google.cloud.dataflow.worker.util.common.worker.ParDoOperation.start(ParDoOperation.java:40)
com.google.cloud.dataflow.worker.util.common.worker.MapTaskExecutor.execute(MapTaskExecutor.java:69)
com.google.cloud.dataflow.worker.runners.worker.StreamingDataflowWorker.process(StreamingDataflowWorker.java:977)
com.google.cloud.dataflow.worker.runners.worker.StreamingDataflowWorker.access$800(StreamingDataflowWorker.java:127)
com.google.cloud.dataflow.worker.runners.worker.StreamingDataflowWorker$7.run(StreamingDataflowWorker.java:761)
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
java.lang.Thread.run(Thread.java:745)
Caused by: com.google.cloud.dataflow.worker.repackaged.com.google.common.util.concurrent.UncheckedExecutionException: java.util.NoSuchElementException: Empty PCollection accessed as a singleton view.
com.google.cloud.dataflow.worker.repackaged.com.google.common.cache.LocalCache$Segment.get(LocalCache.java:2214)
com.google.cloud.dataflow.worker.repackaged.com.google.common.cache.LocalCache.get(LocalCache.java:4053)
com.google.cloud.dataflow.worker.repackaged.com.google.common.cache.LocalCache$LocalManualCache.get(LocalCache.java:4899)
com.google.cloud.dataflow.worker.runners.worker.StateFetcher.fetchSideInput(StateFetcher.java:173)
com.google.cloud.dataflow.worker.runners.worker.StreamingModeExecutionContext.fetchSideInput(StreamingModeExecutionContext.java:211)
com.google.cloud.dataflow.worker.runners.worker.StreamingModeExecutionContext.access$700(StreamingModeExecutionContext.java:70)
com.google.cloud.dataflow.worker.runners.worker.StreamingModeExecutionContext$StepContext.issueSideInputFetch(StreamingModeExecutionContext.java:572)
com.google.cloud.dataflow.worker.runners.worker.StreamingModeExecutionContext$UserStepContext.issueSideInputFetch(StreamingModeExecutionContext.java:624)
com.google.cloud.dataflow.worker.runners.worker.StreamingSideInputFetcher.getReadyWindows(StreamingSideInputFetcher.java:136)
com.google.cloud.dataflow.worker.runners.worker.StreamingSideInputDoFnRunner.startBundle(StreamingSideInputDoFnRunner.java:51)
com.google.cloud.dataflow.worker.runners.worker.SimpleParDoFn.reallyStartBundle(SimpleParDoFn.java:222)
com.google.cloud.dataflow.worker.runners.worker.SimpleParDoFn.startBundle(SimpleParDoFn.java:150)
com.google.cloud.dataflow.worker.runners.worker.ForwardingParDoFn.startBundle(ForwardingParDoFn.java:36)
com.google.cloud.dataflow.worker.util.common.worker.ParDoOperation.start(ParDoOperation.java:40)
com.google.cloud.dataflow.worker.util.common.worker.MapTaskExecutor.execute(MapTaskExecutor.java:69)
com.google.cloud.dataflow.worker.runners.worker.StreamingDataflowWorker.process(StreamingDataflowWorker.java:977)
com.google.cloud.dataflow.worker.runners.worker.StreamingDataflowWorker.access$800(StreamingDataflowWorker.java:127)
com.google.cloud.dataflow.worker.runners.worker.StreamingDataflowWorker$7.run(StreamingDataflowWorker.java:761)
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
java.lang.Thread.run(Thread.java:745)
Caused by: java.util.NoSuchElementException: Empty PCollection accessed as a singleton view.
org.apache.beam.sdk.values.PCollectionViews$SingletonViewFn.getDefaultValue(PCollectionViews.java:190)
org.apache.beam.sdk.values.PCollectionViews$SingletonViewFn.apply(PCollectionViews.java:215)
org.apache.beam.sdk.values.PCollectionViews$SingletonViewFn.apply(PCollectionViews.java:151)
com.google.cloud.dataflow.worker.runners.worker.StateFetcher$2.call(StateFetcher.java:150)
com.google.cloud.dataflow.worker.runners.worker.StateFetcher$2.call(StateFetcher.java:105)
com.google.cloud.dataflow.worker.repackaged.com.google.common.cache.LocalCache$LocalManualCache$1.load(LocalCache.java:4904)
com.google.cloud.dataflow.worker.repackaged.com.google.common.cache.LocalCache$LoadingValueReference.loadFuture(LocalCache.java:3628)
com.google.cloud.dataflow.worker.repackaged.com.google.common.cache.LocalCache$Segment.loadSync(LocalCache.java:2336)
com.google.cloud.dataflow.worker.repackaged.com.google.common.cache.LocalCache$Segment.lockedGetOrLoad(LocalCache.java:2295)
com.google.cloud.dataflow.worker.repackaged.com.google.common.cache.LocalCache$Segment.get(LocalCache.java:2208)
com.google.cloud.dataflow.worker.repackaged.com.google.common.cache.LocalCache.get(LocalCache.java:4053)
com.google.cloud.dataflow.worker.repackaged.com.google.common.cache.LocalCache$LocalManualCache.get(LocalCache.java:4899)
com.google.cloud.dataflow.worker.runners.worker.StateFetcher.fetchSideInput(StateFetcher.java:173)
com.google.cloud.dataflow.worker.runners.worker.StreamingModeExecutionContext.fetchSideInput(StreamingModeExecutionContext.java:211)
com.google.cloud.dataflow.worker.runners.worker.StreamingModeExecutionContext.access$700(StreamingModeExecutionContext.java:70)
com.google.cloud.dataflow.worker.runners.worker.StreamingModeExecutionContext$StepContext.issueSideInputFetch(StreamingModeExecutionContext.java:572)
com.google.cloud.dataflow.worker.runners.worker.StreamingModeExecutionContext$UserStepContext.issueSideInputFetch(StreamingModeExecutionContext.java:624)
com.google.cloud.dataflow.worker.runners.worker.StreamingSideInputFetcher.getReadyWindows(StreamingSideInputFetcher.java:136)
com.google.cloud.dataflow.worker.runners.worker.StreamingSideInputDoFnRunner.startBundle(StreamingSideInputDoFnRunner.java:51)
com.google.cloud.dataflow.worker.runners.worker.SimpleParDoFn.reallyStartBundle(SimpleParDoFn.java:222)
com.google.cloud.dataflow.worker.runners.worker.SimpleParDoFn.startBundle(SimpleParDoFn.java:150)
com.google.cloud.dataflow.worker.runners.worker.ForwardingParDoFn.startBundle(ForwardingParDoFn.java:36)
com.google.cloud.dataflow.worker.util.common.worker.ParDoOperation.start(ParDoOperation.java:40)
com.google.cloud.dataflow.worker.util.common.worker.MapTaskExecutor.execute(MapTaskExecutor.java:69)
com.google.cloud.dataflow.worker.runners.worker.StreamingDataflowWorker.process(StreamingDataflowWorker.java:977)
com.google.cloud.dataflow.worker.runners.worker.StreamingDataflowWorker.access$800(StreamingDataflowWorker.java:127)
com.google.cloud.dataflow.worker.runners.worker.StreamingDataflowWorker$7.run(StreamingDataflowWorker.java:761)
java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
java.lang.Thread.run(Thread.java:745)

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

Contributor guide

Open the contributing guide

Research direction

Start with the Dataflow streaming runner test named testWindowedSideInputWithCheckpoints and reproduce the failure described in the stack trace. Read StreamingSideInputFetcher, StreamingModeExecutionContext, StateFetcher, and PCollectionViews.SingletonViewFn; done means the test completes without the Empty PCollection singleton-view exception.

Written by the indexing model from the issue text.

Assessment

Tech stack
java
Domain
data-engineering, stream-processing
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.