Path Lines of Code runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ArtifactResolver.java 13 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/AvroGenericCoderRegistrar.java 21 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/AvroGenericCoderTranslator.java 24 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/BeamUrns.java 8 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/CoderTranslation.java 130 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/CoderTranslator.java 11 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/CoderTranslatorRegistrar.java 10 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/CoderTranslators.java 159 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/CombineTranslation.java 189 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/CreatePCollectionViewTranslation.java 78 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/DeduplicatedFlattenFactory.java 82 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/DefaultArtifactResolver.java 89 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/DefaultExpansionServiceClientFactory.java 46 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/DisplayDataTranslation.java 52 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/EmptyFlattenAsCreateFactory.java 46 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/Environments.java 425 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ExecutableStageTranslation.java 81 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ExpansionServiceClient.java 5 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ExpansionServiceClientFactory.java 5 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/External.java 326 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ExternalTranslation.java 132 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/FlattenTranslator.java 37 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ForwardingPTransform.java 47 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/GroupByKeyTranslation.java 34 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/GroupIntoBatchesTranslation.java 67 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ImpulseTranslation.java 35 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ModelCoderRegistrar.java 102 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ModelCoders.java 107 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/NativeTransforms.java 23 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/PCollectionTranslation.java 62 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/PCollectionViewTranslation.java 64 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/PTransformMatchers.java 382 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/PTransformReplacements.java 46 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/PTransformTranslation.java 386 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ParDoTranslation.java 684 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/PipelineOptionsTranslation.java 89 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/PipelineTranslation.java 138 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/PrimitiveCreate.java 48 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ReadTranslation.java 166 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/RehydratedComponents.java 121 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ReplacementOutputs.java 69 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/ReshuffleTranslation.java 34 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/RunnerPCollectionView.java 82 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/SdkComponents.java 267 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/SerializablePipelineOptions.java 69 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/SingleInputOutputOverrideFactory.java 18 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/SplittableParDo.java 599 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/SplittableParDoNaiveBounded.java 511 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/SyntheticComponents.java 14 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/TestStreamTranslation.java 161 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/Timer.java 157 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/TransformInputs.java 25 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/TransformPayloadTranslatorRegistrar.java 11 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/TriggerTranslation.java 272 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/UnboundedReadFromBoundedSource.java 402 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/UnconsumedReads.java 47 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/UnsupportedOverrideFactory.java 37 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/WindowIntoTranslation.java 104 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/WindowingStrategyTranslation.java 322 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/WriteFilesTranslation.java 259 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/ExecutableStage.java 123 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/FusedPipeline.java 71 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/GreedyPCollectionFusers.java 233 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/GreedyPipelineFuser.java 292 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/GreedyStageFuser.java 144 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/ImmutableExecutableStage.java 81 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/Networks.java 190 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/OutputDeduplicator.java 265 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/PipelineNode.java 25 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/PipelineValidator.java 260 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/ProtoOverrides.java 50 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/QueryablePipeline.java 328 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/SideInputReference.java 38 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/SplittableParDoExpander.java 263 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/TimerReference.java 18 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/TrivialNativeTransformExpander.java 41 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/UserStateReference.java 35 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/graph/package-info.java 4 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/package-info.java 4 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/renderer/PipelineDotRenderer.java 88 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/renderer/PortablePipelineDotRenderer.java 76 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/renderer/package-info.java 1 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/resources/ClasspathScanningResourcesDetector.java 17 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/resources/PipelineResources.java 88 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/resources/PipelineResourcesDetector.java 9 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/resources/PipelineResourcesOptions.java 45 runners/core-construction-java/src/main/java/org/apache/beam/runners/core/construction/resources/package-info.java 1 runners/core-java/src/main/java/org/apache/beam/runners/core/ActiveWindowSet.java 27 runners/core-java/src/main/java/org/apache/beam/runners/core/DoFnRunner.java 23 runners/core-java/src/main/java/org/apache/beam/runners/core/DoFnRunners.java 139 runners/core-java/src/main/java/org/apache/beam/runners/core/ElementByteSizeObservable.java 6 runners/core-java/src/main/java/org/apache/beam/runners/core/GlobalCombineFnRunner.java 33 runners/core-java/src/main/java/org/apache/beam/runners/core/GlobalCombineFnRunners.java 160 runners/core-java/src/main/java/org/apache/beam/runners/core/GroupAlsoByWindowViaWindowSetNewDoFn.java 107 runners/core-java/src/main/java/org/apache/beam/runners/core/GroupAlsoByWindowsAggregators.java 5 runners/core-java/src/main/java/org/apache/beam/runners/core/GroupByKeyViaGroupByKeyOnly.java 145 runners/core-java/src/main/java/org/apache/beam/runners/core/InMemoryBundleFinalizer.java 27 runners/core-java/src/main/java/org/apache/beam/runners/core/InMemoryMultimapSideInputView.java 65 runners/core-java/src/main/java/org/apache/beam/runners/core/InMemoryStateInternals.java 591 runners/core-java/src/main/java/org/apache/beam/runners/core/InMemoryTimerInternals.java 242 runners/core-java/src/main/java/org/apache/beam/runners/core/KeyedWorkItem.java 8 runners/core-java/src/main/java/org/apache/beam/runners/core/KeyedWorkItemCoder.java 68 runners/core-java/src/main/java/org/apache/beam/runners/core/KeyedWorkItems.java 68 runners/core-java/src/main/java/org/apache/beam/runners/core/LateDataDroppingDoFnRunner.java 106 runners/core-java/src/main/java/org/apache/beam/runners/core/LateDataUtils.java 71 runners/core-java/src/main/java/org/apache/beam/runners/core/MergingActiveWindowSet.java 290 runners/core-java/src/main/java/org/apache/beam/runners/core/MergingStateAccessor.java 7 runners/core-java/src/main/java/org/apache/beam/runners/core/NonEmptyPanes.java 72 runners/core-java/src/main/java/org/apache/beam/runners/core/NonMergingActiveWindowSet.java 50 runners/core-java/src/main/java/org/apache/beam/runners/core/NullSideInputReader.java 30 runners/core-java/src/main/java/org/apache/beam/runners/core/OutputAndTimeBoundedSplittableProcessElementInvoker.java 295 runners/core-java/src/main/java/org/apache/beam/runners/core/OutputWindowedValue.java 19 runners/core-java/src/main/java/org/apache/beam/runners/core/PaneInfoTracker.java 105 runners/core-java/src/main/java/org/apache/beam/runners/core/PeekingReiterator.java 59 runners/core-java/src/main/java/org/apache/beam/runners/core/ProcessFnRunner.java 99 runners/core-java/src/main/java/org/apache/beam/runners/core/PushbackSideInputDoFnRunner.java 20 runners/core-java/src/main/java/org/apache/beam/runners/core/ReadyCheckingSideInputReader.java 6 runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFn.java 37 runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnContextFactory.java 464 runners/core-java/src/main/java/org/apache/beam/runners/core/ReduceFnRunner.java 629 runners/core-java/src/main/java/org/apache/beam/runners/core/SideInputHandler.java 157 runners/core-java/src/main/java/org/apache/beam/runners/core/SideInputReader.java 10 runners/core-java/src/main/java/org/apache/beam/runners/core/SimpleDoFnRunner.java 1033 runners/core-java/src/main/java/org/apache/beam/runners/core/SimplePushbackSideInputDoFnRunner.java 91 runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java 478 runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableProcessElementInvoker.java 50 runners/core-java/src/main/java/org/apache/beam/runners/core/StateAccessor.java 5 runners/core-java/src/main/java/org/apache/beam/runners/core/StateInternals.java 12 runners/core-java/src/main/java/org/apache/beam/runners/core/StateInternalsFactory.java 5 runners/core-java/src/main/java/org/apache/beam/runners/core/StateMerging.java 136 runners/core-java/src/main/java/org/apache/beam/runners/core/StateNamespace.java 7 runners/core-java/src/main/java/org/apache/beam/runners/core/StateNamespaceForTest.java 36 runners/core-java/src/main/java/org/apache/beam/runners/core/StateNamespaces.java 195 runners/core-java/src/main/java/org/apache/beam/runners/core/StateTable.java 60 runners/core-java/src/main/java/org/apache/beam/runners/core/StateTag.java 45 runners/core-java/src/main/java/org/apache/beam/runners/core/StateTags.java 253 runners/core-java/src/main/java/org/apache/beam/runners/core/StatefulDoFnRunner.java 314 runners/core-java/src/main/java/org/apache/beam/runners/core/StepContext.java 9 runners/core-java/src/main/java/org/apache/beam/runners/core/SystemReduceFn.java 94 runners/core-java/src/main/java/org/apache/beam/runners/core/TestInMemoryStateInternals.java 40 runners/core-java/src/main/java/org/apache/beam/runners/core/TimerInternals.java 169 runners/core-java/src/main/java/org/apache/beam/runners/core/TimerInternalsFactory.java 5 runners/core-java/src/main/java/org/apache/beam/runners/core/UnsupportedSideInputReader.java 24 runners/core-java/src/main/java/org/apache/beam/runners/core/WatermarkHold.java 298 runners/core-java/src/main/java/org/apache/beam/runners/core/WindowingInternals.java 26 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/CounterCell.java 63 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/DefaultMetricResults.java 34 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/DirtyState.java 38 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/DistributionCell.java 62 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/DistributionData.java 28 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/ExecutionStateSampler.java 115 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/ExecutionStateTracker.java 176 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/GaugeCell.java 57 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/GaugeData.java 47 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/HistogramCell.java 69 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/LabeledMetrics.java 25 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MetricCell.java 7 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MetricUpdates.java 35 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MetricsContainerImpl.java 426 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MetricsContainerStepMap.java 166 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MetricsLogger.java 60 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MetricsMap.java 51 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MetricsPusher.java 86 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoConstants.java 120 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoEncodings.java 100 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java 80 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/NoOpMetricsSink.java 9 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/ShortIdMap.java 26 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/SimpleExecutionState.java 113 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/SimpleMonitoringInfoBuilder.java 88 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/SimpleStateRegistry.java 48 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/SpecMonitoringInfoValidator.java 46 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/package-info.java 4 runners/core-java/src/main/java/org/apache/beam/runners/core/package-info.java 4 runners/core-java/src/main/java/org/apache/beam/runners/core/serialization/Base64Serializer.java 41 runners/core-java/src/main/java/org/apache/beam/runners/core/serialization/package-info.java 1 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterAllStateMachine.java 81 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterDelayFromFirstElementStateMachine.java 192 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterEachStateMachine.java 88 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterFirstStateMachine.java 89 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterPaneStateMachine.java 87 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterProcessingTimeStateMachine.java 47 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterSynchronizedProcessingTimeStateMachine.java 37 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/AfterWatermarkStateMachine.java 213 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/DefaultTriggerStateMachine.java 47 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/ExecutableTriggerStateMachine.java 98 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/FinishedTriggers.java 7 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/FinishedTriggersBitSet.java 33 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/FinishedTriggersSet.java 38 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/NeverStateMachine.java 30 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/OrFinallyStateMachine.java 72 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/RepeatedlyStateMachine.java 55 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/ReshuffleTriggerStateMachine.java 30 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachine.java 142 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineContextFactory.java 466 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachineRunner.java 132 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/TriggerStateMachines.java 98 runners/core-java/src/main/java/org/apache/beam/runners/core/triggers/package-info.java 4 runners/direct-java/src/main/java/org/apache/beam/runners/direct/AbstractModelEnforcement.java 13 runners/direct-java/src/main/java/org/apache/beam/runners/direct/BoundedReadEvaluatorFactory.java 188 runners/direct-java/src/main/java/org/apache/beam/runners/direct/BundleFactory.java 10 runners/direct-java/src/main/java/org/apache/beam/runners/direct/BundleProcessor.java 6 runners/direct-java/src/main/java/org/apache/beam/runners/direct/Clock.java 6 runners/direct-java/src/main/java/org/apache/beam/runners/direct/CloningBundleFactory.java 69 runners/direct-java/src/main/java/org/apache/beam/runners/direct/CommittedBundle.java 21 runners/direct-java/src/main/java/org/apache/beam/runners/direct/CommittedResult.java 27 runners/direct-java/src/main/java/org/apache/beam/runners/direct/CompletionCallback.java 11 runners/direct-java/src/main/java/org/apache/beam/runners/direct/CopyOnAccessInMemoryStateInternals.java 336 runners/direct-java/src/main/java/org/apache/beam/runners/direct/CreateViewNoopEvaluatorFactory.java 31 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectExecutionContext.java 87 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectGBKIntoKeyedWorkItemsOverrideFactory.java 27 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectGraph.java 85 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectGraphVisitor.java 122 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectGroupByKey.java 91 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectGroupByKeyOverrideFactory.java 30 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectMetrics.java 214 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectOptions.java 43 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectRegistrar.java 24 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectRunner.java 291 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectTestOptions.java 13 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectTimerInternals.java 86 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectTransformExecutor.java 142 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DirectWriteViewVisitor.java 77 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DisplayDataValidator.java 38 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DoFnLifecycleManager.java 75 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DoFnLifecycleManagerRemovingTransformEvaluator.java 61 runners/direct-java/src/main/java/org/apache/beam/runners/direct/DoFnLifecycleManagers.java 20 runners/direct-java/src/main/java/org/apache/beam/runners/direct/EmptyInputProvider.java 18 runners/direct-java/src/main/java/org/apache/beam/runners/direct/EvaluationContext.java 262 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ExecutableGraph.java 11 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ExecutorServiceFactory.java 5 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ExecutorServiceParallelExecutor.java 327 runners/direct-java/src/main/java/org/apache/beam/runners/direct/FlattenEvaluatorFactory.java 52 runners/direct-java/src/main/java/org/apache/beam/runners/direct/GroupAlsoByWindowEvaluatorFactory.java 213 runners/direct-java/src/main/java/org/apache/beam/runners/direct/GroupByKeyOnlyEvaluatorFactory.java 105 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ImmutabilityCheckingBundleFactory.java 94 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ImmutabilityEnforcementFactory.java 121 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ImmutableListBundleFactory.java 134 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ImpulseEvaluatorFactory.java 68 runners/direct-java/src/main/java/org/apache/beam/runners/direct/KeyedPValueTrackingVisitor.java 85 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ModelEnforcement.java 13 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ModelEnforcementFactory.java 5 runners/direct-java/src/main/java/org/apache/beam/runners/direct/MultiStepCombine.java 419 runners/direct-java/src/main/java/org/apache/beam/runners/direct/NanosOffsetClock.java 21 runners/direct-java/src/main/java/org/apache/beam/runners/direct/PCollectionViewWindow.java 35 runners/direct-java/src/main/java/org/apache/beam/runners/direct/PCollectionViewWriter.java 8 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ParDoEvaluator.java 256 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ParDoEvaluatorFactory.java 143 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ParDoMultiOverrideFactory.java 246 runners/direct-java/src/main/java/org/apache/beam/runners/direct/PassthroughTransformEvaluator.java 24 runners/direct-java/src/main/java/org/apache/beam/runners/direct/PipelineExecutor.java 12 runners/direct-java/src/main/java/org/apache/beam/runners/direct/QuiescenceDriver.java 253 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ReadEvaluatorFactory.java 64 runners/direct-java/src/main/java/org/apache/beam/runners/direct/RootInputProvider.java 13 runners/direct-java/src/main/java/org/apache/beam/runners/direct/RootProviderRegistry.java 52 runners/direct-java/src/main/java/org/apache/beam/runners/direct/SideInputContainer.java 233 runners/direct-java/src/main/java/org/apache/beam/runners/direct/SourceShard.java 8 runners/direct-java/src/main/java/org/apache/beam/runners/direct/SplittableProcessElementsEvaluatorFactory.java 187 runners/direct-java/src/main/java/org/apache/beam/runners/direct/StatefulParDoEvaluatorFactory.java 209 runners/direct-java/src/main/java/org/apache/beam/runners/direct/StepAndKey.java 39 runners/direct-java/src/main/java/org/apache/beam/runners/direct/StepTransformResult.java 119 runners/direct-java/src/main/java/org/apache/beam/runners/direct/TestStreamEvaluatorFactory.java 196 runners/direct-java/src/main/java/org/apache/beam/runners/direct/TransformEvaluator.java 6 runners/direct-java/src/main/java/org/apache/beam/runners/direct/TransformEvaluatorFactory.java 12 runners/direct-java/src/main/java/org/apache/beam/runners/direct/TransformEvaluatorRegistry.java 151 runners/direct-java/src/main/java/org/apache/beam/runners/direct/TransformExecutor.java 2 runners/direct-java/src/main/java/org/apache/beam/runners/direct/TransformExecutorFactory.java 9 runners/direct-java/src/main/java/org/apache/beam/runners/direct/TransformExecutorService.java 6 runners/direct-java/src/main/java/org/apache/beam/runners/direct/TransformExecutorServices.java 116 runners/direct-java/src/main/java/org/apache/beam/runners/direct/TransformResult.java 30 runners/direct-java/src/main/java/org/apache/beam/runners/direct/UnboundedReadDeduplicator.java 50 runners/direct-java/src/main/java/org/apache/beam/runners/direct/UnboundedReadEvaluatorFactory.java 274 runners/direct-java/src/main/java/org/apache/beam/runners/direct/UncommittedBundle.java 12 runners/direct-java/src/main/java/org/apache/beam/runners/direct/ViewEvaluatorFactory.java 55 runners/direct-java/src/main/java/org/apache/beam/runners/direct/WatermarkCallbackExecutor.java 134 runners/direct-java/src/main/java/org/apache/beam/runners/direct/WatermarkManager.java 1197 runners/direct-java/src/main/java/org/apache/beam/runners/direct/WindowEvaluatorFactory.java 95 runners/direct-java/src/main/java/org/apache/beam/runners/direct/WriteWithShardingFactory.java 110 runners/direct-java/src/main/java/org/apache/beam/runners/direct/package-info.java 1 runners/extensions-java/metrics/src/main/java/org/apache/beam/runners/extensions/metrics/MetricsGraphiteSink.java 268 runners/extensions-java/metrics/src/main/java/org/apache/beam/runners/extensions/metrics/MetricsHttpSink.java 134 runners/extensions-java/metrics/src/main/java/org/apache/beam/runners/extensions/metrics/package-info.java 4 runners/flink/1.10/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/AbstractStreamOperatorCompat.java 9 runners/flink/1.11/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/AbstractStreamOperatorCompat.java 9 runners/flink/1.12/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/AbstractStreamOperatorCompat.java 21 runners/flink/src/main/java/org/apache/beam/runners/flink/CreateStreamingFlinkView.java 118 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkBatchPipelineTranslator.java 89 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkBatchPortablePipelineTranslator.java 486 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkBatchTransformTranslators.java 719 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkBatchTranslationContext.java 107 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkDetachedRunnerResult.java 32 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkExecutionEnvironments.java 301 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkJobInvoker.java 77 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkJobServerDriver.java 100 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkMiniClusterEntryPoint.java 66 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkPipelineExecutionEnvironment.java 81 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkPipelineOptions.java 187 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkPipelineRunner.java 181 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkPipelineTranslator.java 14 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkPortableClientEntryPoint.java 196 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkPortablePipelineTranslator.java 24 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkPortableRunnerResult.java 30 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkRunner.java 133 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkRunnerRegistrar.java 24 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkRunnerResult.java 50 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkStateBackendFactory.java 5 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkStreamingPipelineTranslator.java 289 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkStreamingPortablePipelineTranslator.java 828 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkStreamingTransformTranslators.java 1291 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkStreamingTranslationContext.java 96 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkTransformOverrides.java 43 runners/flink/src/main/java/org/apache/beam/runners/flink/PipelineTranslationModeOptimizer.java 51 runners/flink/src/main/java/org/apache/beam/runners/flink/TestFlinkRunner.java 53 runners/flink/src/main/java/org/apache/beam/runners/flink/metrics/DoFnRunnerWithMetricsUpdate.java 74 runners/flink/src/main/java/org/apache/beam/runners/flink/metrics/FileReporter.java 53 runners/flink/src/main/java/org/apache/beam/runners/flink/metrics/FlinkMetricContainer.java 142 runners/flink/src/main/java/org/apache/beam/runners/flink/metrics/Metrics.java 36 runners/flink/src/main/java/org/apache/beam/runners/flink/metrics/MetricsAccumulator.java 33 runners/flink/src/main/java/org/apache/beam/runners/flink/metrics/ReaderInvocationUtil.java 44 runners/flink/src/main/java/org/apache/beam/runners/flink/metrics/package-info.java 1 runners/flink/src/main/java/org/apache/beam/runners/flink/package-info.java 1 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/AbstractFlinkCombineRunner.java 153 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkAssignContext.java 35 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkAssignWindows.java 23 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkDoFnFunction.java 199 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkExecutableStageContextFactory.java 32 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkExecutableStageFunction.java 322 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkExecutableStagePruningFunction.java 31 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkExplodeWindowsFunction.java 11 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkIdentityFunction.java 14 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkMergingNonShuffleReduceFunction.java 58 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkMultiOutputPruningFunction.java 34 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkNoOpStepContext.java 14 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkNonMergingReduceFunction.java 76 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkPartialReduceFunction.java 71 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkReduceFunction.java 71 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkSideInputReader.java 86 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkStatefulDoFnFunction.java 195 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/FlinkStreamingSideInputHandlerFactory.java 118 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/HashingFlinkCombineRunner.java 140 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/ImpulseSourceFunction.java 63 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/SideInputInitializer.java 92 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/SingleWindowFlinkCombineRunner.java 85 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/SortingFlinkCombineRunner.java 124 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/functions/package-info.java 1 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/package-info.java 1 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/types/CoderTypeInformation.java 87 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/types/CoderTypeSerializer.java 133 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/types/EncodedValueComparator.java 139 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/types/EncodedValueSerializer.java 84 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/types/EncodedValueTypeInformation.java 61 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/types/KvKeySelector.java 23 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/types/WindowedKvKeySelector.java 33 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/types/package-info.java 1 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/utils/CheckpointStats.java 22 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/utils/CountingPipelineVisitor.java 21 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/utils/FlinkPortableRunnerUtils.java 27 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/utils/LookupPipelineVisitor.java 65 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/utils/Workarounds.java 31 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/utils/package-info.java 1 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/DataInputViewWrapper.java 23 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/DataOutputViewWrapper.java 18 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/ImpulseInputFormat.java 61 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/SourceInputFormat.java 117 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/SourceInputSplit.java 22 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/package-info.java 1 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/DoFnOperator.java 1149 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/ExecutableStageDoFnOperator.java 1012 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/FlinkKeyUtils.java 70 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/KeyedPushedBackElementsHandler.java 70 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/KvToByteBufferKeySelector.java 28 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/NonKeyedPushedBackElementsHandler.java 32 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/PushedBackElementsHandler.java 8 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/SdfByteBufferKeySelector.java 29 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/SingletonKeyedWorkItem.java 28 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/SingletonKeyedWorkItemCoder.java 75 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/SplittableDoFnOperator.java 153 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/WindowDoFnOperator.java 95 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/WorkItemKeySelector.java 29 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/BeamStoppableFunction.java 4 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/DedupingOperator.java 72 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/StreamingImpulseSource.java 44 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/TestStreamSource.java 53 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/UnboundedSourceWrapper.java 368 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/package-info.java 1 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/package-info.java 1 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/BufferedElement.java 5 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/BufferedElements.java 154 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/BufferingDoFnRunner.java 194 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/BufferingElementsHandler.java 7 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/KeyedBufferingElementsHandler.java 67 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/NonKeyedBufferingElementsHandler.java 34 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/package-info.java 4 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/state/FlinkBroadcastStateInternals.java 676 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/state/FlinkStateInternals.java 1324 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/state/package-info.java 1 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/BatchStatefulParDoOverrides.java 222 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/BatchViewOverrides.java 913 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/CreateDataflowView.java 35 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowClient.java 122 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowJobAlreadyExistsException.java 6 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowJobAlreadyUpdatedException.java 6 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowJobException.java 13 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowMetrics.java 252 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPTransformMatchers.java 79 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineJob.java 380 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineRegistrar.java 25 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslator.java 1005 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java 1751 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunnerHooks.java 7 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunnerInfo.java 100 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowServiceException.java 10 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/GroupIntoBatchesOverride.java 182 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/PrimitiveParDoSingleFactory.java 294 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/ReadTranslator.java 62 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/RequiresStableInputParDoOverrides.java 76 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/ReshuffleOverrideFactory.java 61 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/SplittableParDoOverrides.java 50 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/StreamingViewOverrides.java 78 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/TestDataflowPipelineOptions.java 5 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/TestDataflowRunner.java 289 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/TransformTranslator.java 57 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/internal/CustomSources.java 65 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/internal/IsmFormat.java 506 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/internal/package-info.java 1 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/CloudDebuggerOptions.java 28 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowPipelineDebugOptions.java 128 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowPipelineOptions.java 138 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowPipelineWorkerPoolOptions.java 105 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowProfilingOptions.java 22 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowWorkerHarnessOptions.java 18 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowWorkerLoggingOptions.java 85 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DefaultGcpRegionFactory.java 64 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/package-info.java 1 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/package-info.java 1 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/AvroCoderCloudObjectTranslator.java 39 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/CloudKnownType.java 110 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/CloudObject.java 96 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/CloudObjectKinds.java 10 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/CloudObjectTranslator.java 8 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/CloudObjectTranslators.java 477 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/CloudObjects.java 99 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/CoderCloudObjectTranslatorRegistrar.java 12 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/DataflowTemplateJob.java 45 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/DataflowTransport.java 74 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/DefaultCoderCloudObjectTranslatorRegistrar.java 109 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/GcsStager.java 42 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/MonitoringUtil.java 171 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/OutputReference.java 35 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/PackageUtil.java 379 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/PropertyNames.java 54 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/RandomAccessData.java 230 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/RowCoderCloudObjectTranslator.java 52 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/SchemaCoderCloudObjectTranslator.java 95 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/SerializableCoderCloudObjectTranslator.java 41 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/Stager.java 7 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/Structs.java 285 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/TimeUtil.java 98 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/package-info.java 1 runners/google-cloud-dataflow-java/worker/src/main/java/com/google/cloud/dataflow/worker/DataflowRunnerHarness.java 6 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/harness/util/package-info.java 1 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ApplianceShuffleCounters.java 38 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ApplianceShuffleEntryReader.java 41 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ApplianceShuffleReader.java 45 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ApplianceShuffleWriter.java 33 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/AssignWindowsParDoFnFactory.java 96 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/AvroByteReader.java 142 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/AvroByteReaderFactory.java 46 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/AvroByteSink.java 51 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/AvroByteSinkFactory.java 35 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchDataflowWorker.java 310 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchModeExecutionContext.java 419 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchModeUngroupingParDoFn.java 45 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BeamFnMapTaskExecutorFactory.java 669 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ByteStringCoder.java 26 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ChunkingShuffleBatchReader.java 65 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/CombinePhase.java 7 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/CombineValuesFnFactory.java 278 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ConcatReader.java 204 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ConcatReaderFactory.java 101 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ContextActivationObserver.java 14 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ContextActivationObserverRegistry.java 44 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/CounterShortIdCache.java 114 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/CreateIsmShardKeyAndSortKeyDoFnFactory.java 80 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowApiUtils.java 30 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowBatchWorkerHarness.java 126 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowElementExecutionTracker.java 244 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowExecutionContext.java 200 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowExecutionStateKey.java 36 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowExecutionStateRegistry.java 62 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowMapTaskExecutor.java 16 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowMapTaskExecutorFactory.java 28 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowMetricsContainer.java 36 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOperationContext.java 293 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java 64 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowPortabilityPCollectionView.java 125 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowProcessFnRunner.java 93 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowRunnerHarness.java 188 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowSideInputReadCounter.java 91 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowSystemMetrics.java 40 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkExecutor.java 6 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkProgressUpdater.java 119 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java 190 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkerHarnessHelper.java 109 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DefaultParDoFnFactory.java 56 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DeltaCounterCell.java 53 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DeltaDistributionCell.java 50 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DoFnInstanceManager.java 10 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DoFnInstanceManagers.java 81 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DoFnRunnerFactory.java 30 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ExperimentContext.java 49 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/FetchAndFilterStreamingSideInputsOperation.java 216 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/Filepatterns.java 48 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/FnApiWindowMappingFn.java 254 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ForwardingParDoFn.java 29 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/GroupAlsoByWindowFn.java 16 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/GroupAlsoByWindowFnRunner.java 101 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/GroupAlsoByWindowParDoFnFactory.java 317 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/GroupAlsoByWindowsParDoFn.java 169 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/GroupingShuffleReader.java 357 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/GroupingShuffleReaderFactory.java 83 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/GroupingShuffleReaderWithFaultyBytesReadCounter.java 42 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/HotKeyLogger.java 51 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/InMemoryReader.java 138 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/InMemoryReaderFactory.java 43 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutor.java 30 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IntrinsicMapTaskExecutorFactory.java 327 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IsmReader.java 60 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IsmReaderFactory.java 119 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IsmReaderImpl.java 842 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IsmSideInputReader.java 772 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IsmSink.java 202 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IsmSinkFactory.java 79 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/KeyTokenInvalidException.java 18 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/LazilyInitializedSideInputReader.java 35 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/MetricTrackingWindmillServerStub.java 234 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/MetricsContainerRegistry.java 15 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/MetricsEnvironmentContextActivationObserverRegistration.java 30 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/MetricsToCounterUpdateConverter.java 74 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/NoOpSourceOperationExecutor.java 42 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/NoopSideInputReadCounter.java 11 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/OrderedCode.java 388 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PCollectionViewWindow.java 35 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PairWithConstantKeyDoFnFactory.java 71 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ParDoFnFactory.java 19 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PartialGroupByKeyParDoFns.java 305 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PartitioningShuffleReader.java 120 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PartitioningShuffleReaderFactory.java 59 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PubsubReader.java 105 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PubsubSink.java 154 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ReaderCache.java 88 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ReaderFactory.java 24 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ReaderRegistry.java 78 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ReaderUtils.java 31 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ReifyTimestampAndWindowsParDoFnFactory.java 65 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/RunnerHarnessCoderCloudObjectTranslatorRegistrar.java 103 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SdkHarnessRegistries.java 202 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SdkHarnessRegistry.java 24 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ShuffleLibrary.java 32 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ShuffleReader.java 14 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ShuffleSink.java 264 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ShuffleSinkFactory.java 52 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ShuffleWriter.java 7 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SideInputReadCounter.java 6 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SideInputTrackingIsmReader.java 56 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SimpleDoFnRunnerFactory.java 64 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SimpleParDoFn.java 384 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SinkFactory.java 24 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SinkRegistry.java 62 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SizeReportingSinkWrapper.java 44 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SourceOperationExecutor.java 5 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SourceOperationExecutorFactory.java 35 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SourceTranslationUtils.java 118 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/SplittableProcessFnFactory.java 176 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StateFetcher.java 217 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java 2053 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowReshuffleFn.java 36 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowViaWindowSetFn.java 77 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingGroupAlsoByWindowsDoFns.java 34 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingKeyedWorkItemSideInputDoFnRunner.java 122 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingModeExecutionContext.java 650 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingPCollectionViewWriterDoFnFactory.java 48 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingPCollectionViewWriterParDoFn.java 50 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingSideInputDoFnRunner.java 68 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingSideInputFetcher.java 301 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingStepMetricsContainer.java 102 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ToIsmRecordForMultimapDoFnFactory.java 112 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedShuffleReader.java 95 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedShuffleReaderFactory.java 51 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UngroupedWindmillReader.java 101 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/UserParDoFnFactory.java 114 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ValuesDoFnFactory.java 56 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/Weighers.java 11 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillComputationKey.java 19 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillKeyedWorkItem.java 169 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillNamespacePrefix.java 19 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillReaderIteratorBase.java 50 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillSink.java 167 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillStateCache.java 288 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillStateInternals.java 1696 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillStateReader.java 717 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimeUtils.java 34 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternals.java 324 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindowingWindmillReader.java 133 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WorkItemStatusClient.java 254 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WorkUnitClient.java 12 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSourceOperationExecutor.java 74 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WorkerCustomSources.java 677 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WorkerPipelineOptionsFactory.java 52 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WorkerUncaughtExceptionHandler.java 38 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/apiary/Apiary.java 14 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/apiary/FixMultiOutputInfosOnParDoInstructions.java 38 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/Counter.java 126 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/CounterBackedElementByteSizeObserver.java 12 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/CounterFactory.java 632 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/CounterName.java 188 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/CounterSet.java 88 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/CounterUpdateAggregator.java 6 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/CounterUpdateAggregators.java 36 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/DataflowCounterUpdateExtractor.java 170 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/DistributionCounterUpdateAggregator.java 46 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/MeanCounterUpdateAggregator.java 36 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/NameContext.java 17 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/SumCounterUpdateAggregator.java 28 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/BeamFnControlService.java 94 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/BeamFnMapTaskExecutor.java 484 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/DataflowSideInputHandlerFactory.java 126 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/ElementCountMonitoringInfoToCounterUpdateTransformer.java 79 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/ExecutionTimeMonitoringInfoToCounterUpdateTransformer.java 96 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/FnApiMonitoringInfoToCounterUpdateTransformer.java 56 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/MeanByteCountMonitoringInfoToCounterUpdateTransformer.java 84 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/MonitoringInfoToCounterUpdateTransformer.java 8 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/ProcessRemoteBundleOperation.java 125 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/RegisterAndProcessBundleOperation.java 600 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/UserDistributionMonitoringInfoToCounterUpdateTransformer.java 101 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/control/UserMonitoringInfoToCounterUpdateTransformer.java 94 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/data/BeamFnDataGrpcService.java 187 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/data/RemoteGrpcPortReadOperation.java 81 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/data/RemoteGrpcPortWriteOperation.java 201 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/grpc/BeamFnService.java 5 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/logging/BeamFnLoggingService.java 101 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/fn/stream/ServerStreamObserverFactory.java 61 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/CloneAmbiguousFlattensFunction.java 73 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/CreateExecutableStageNodeFunction.java 480 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/CreateRegisterFnOperationFunction.java 182 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/DeduceFlattenLocationsFunction.java 218 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/DeduceNodeLocationsFunction.java 75 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/Edges.java 60 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/InsertFetchAndFilterStreamingSideInputNodes.java 129 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/LengthPrefixUnknownCoders.java 248 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/MapTaskToNetworkFunction.java 118 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/Networks.java 143 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/Nodes.java 276 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/RegisterNodeFunction.java 493 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/RemoveFlattenInstructionsFunction.java 45 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/ReplacePgbkWithPrecombineFunction.java 51 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandler.java 231 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingInitializer.java 180 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingMDC.java 41 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/JulHandlerPrintStreamAdapterFactory.java 350 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/options/StreamingDataflowWorkerOptions.java 144 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/profiler/Profiler.java 8 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/profiler/ScopedProfiler.java 108 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/status/BaseStatusServlet.java 26 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/status/DebugCapture.java 172 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/status/HealthzServlet.java 23 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/status/HeapzServlet.java 71 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/status/LastExceptionDataProvider.java 20 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/status/SdkWorkerStatusServlet.java 66 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/status/StatusDataProvider.java 5 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/status/StatuszServlet.java 50 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/status/ThreadzServlet.java 95 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/status/WorkerStatusPages.java 94 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/status/package-info.java 3 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/BatchGroupAlsoByWindowAndCombineFn.java 231 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/BatchGroupAlsoByWindowFn.java 6 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/BatchGroupAlsoByWindowReshuffleFn.java 34 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/BatchGroupAlsoByWindowViaIteratorsFn.java 194 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/BatchGroupAlsoByWindowViaOutputBufferFn.java 91 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/BatchGroupAlsoByWindowsDoFns.java 30 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/BoundedQueueExecutor.java 50 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/CloudSourceUtils.java 24 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/MemoryMonitor.java 417 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/ScalableBloomFilter.java 155 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/StreamingGroupAlsoByWindowFn.java 6 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/TimerOrElement.java 82 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/ValueInEmptyWindows.java 56 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/WorkerPropertyNames.java 27 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/ForwardingReiterator.java 42 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/TaggedReiteratorList.java 119 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/BatchingShuffleEntryReader.java 88 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ByteArrayShufflePosition.java 70 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/CachingShuffleBatchReader.java 74 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ElementCounter.java 5 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ElementExecutionTracker.java 21 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/FlattenOperation.java 26 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/GroupingShuffleEntryIterator.java 202 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/GroupingShuffleRangeTracker.java 148 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/GroupingTable.java 5 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/GroupingTables.java 354 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/JvmRuntime.java 11 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/KeyGroupedShuffleEntries.java 13 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/MapTaskExecutor.java 132 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/NativeReader.java 74 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/Operation.java 62 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/OperationContext.java 13 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/OutputObjectAndByteCounter.java 114 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/OutputReceiver.java 36 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ParDoFn.java 8 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ParDoOperation.java 50 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ProgressTracker.java 5 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ProgressTrackerGroup.java 28 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ProgressTrackingReiterator.java 23 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ReadOperation.java 302 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/Receiver.java 4 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ReceivingOperation.java 11 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ShuffleBatchReader.java 16 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ShuffleEntry.java 80 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ShuffleEntryReader.java 10 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ShufflePosition.java 2 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ShuffleReadCounter.java 67 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/ShuffleReadCounterFactory.java 13 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/SimplePartialGroupByKeyParDoFn.java 24 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/Sink.java 14 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/WorkExecutor.java 30 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/WorkProgressUpdater.java 209 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/WriteOperation.java 101 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/DirectStreamObserver.java 41 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/ForwardingClientResponseObserver.java 34 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/GrpcWindmillServer.java 1378 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/StreamObserverFactory.java 29 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/WindmillServer.java 29 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/WindmillServerBase.java 103 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/WindmillServerStub.java 135 runners/google-cloud-dataflow-java/worker/windmill/src/main/proto/pubsub.proto 31 runners/google-cloud-dataflow-java/worker/windmill/src/main/proto/windmill.proto 530 runners/google-cloud-dataflow-java/worker/windmill/src/main/proto/windmill_service.proto 41 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/FnService.java 6 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/GrpcContextHeaderAccessorProvider.java 47 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/GrpcFnServer.java 105 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/HeaderAccessor.java 4 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/InProcessServerFactory.java 42 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/ServerFactory.java 191 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/artifact/ArtifactRetrievalService.java 105 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/artifact/ArtifactStagingService.java 517 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/artifact/package-info.java 1 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/BundleCheckpointHandler.java 5 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/BundleCheckpointHandlers.java 106 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/BundleFinalizationHandler.java 4 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/BundleFinalizationHandlers.java 33 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/BundleProgressHandler.java 15 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/BundleSplitHandler.java 15 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/ControlClientPool.java 17 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/DefaultExecutableStageContext.java 21 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/DefaultJobBundleFactory.java 597 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/ExecutableStageContext.java 10 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/FnApiControlClient.java 136 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/FnApiControlClientPoolService.java 89 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/InstructionRequestHandler.java 7 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/JobBundleFactory.java 5 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/MapControlClientPool.java 45 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/OutputReceiverFactory.java 5 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/ProcessBundleDescriptors.java 445 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/ReferenceCountingExecutableStageContextFactory.java 179 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/RemoteBundle.java 17 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/RemoteOutputReceiver.java 16 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/SdkHarnessClient.java 528 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/SingleEnvironmentInstanceJobBundleFactory.java 175 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/StageBundleFactory.java 72 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/TimerReceiverFactory.java 83 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/control/package-info.java 1 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/data/FnDataService.java 11 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/data/GrpcDataService.java 146 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/data/RemoteInputDestination.java 12 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/data/package-info.java 1 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/DockerCommand.java 168 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/DockerContainerEnvironment.java 61 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/DockerEnvironmentFactory.java 211 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/EmbeddedEnvironmentFactory.java 138 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/EnvironmentFactory.java 28 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/ExternalEnvironmentFactory.java 164 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/ProcessEnvironment.java 60 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/ProcessEnvironmentFactory.java 120 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/ProcessManager.java 177 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/RemoteEnvironment.java 22 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/StaticRemoteEnvironment.java 35 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/StaticRemoteEnvironmentFactory.java 40 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/package-info.java 1 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/testing/NeedsDocker.java 2 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/environment/testing/package-info.java 1 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/graph/package-info.java 1 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/logging/GrpcLoggingService.java 74 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/logging/LogWriter.java 5 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/logging/Slf4jLogWriter.java 54 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/logging/package-info.java 1 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/package-info.java 1 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/provisioning/JobInfo.java 22 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/provisioning/StaticGrpcProvisionService.java 53 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/provisioning/package-info.java 1 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/state/GrpcStateService.java 116 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/state/InMemoryBagUserStateFactory.java 82 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/state/StateDelegator.java 10 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/state/StateRequestHandler.java 18 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/state/StateRequestHandlers.java 440 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/state/package-info.java 1 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/status/BeamWorkerStatusGrpcService.java 153 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/status/WorkerStatusClient.java 101 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/status/package-info.java 1 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/translation/BatchSideInputHandlerFactory.java 134 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/translation/PipelineTranslatorUtils.java 156 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/translation/package-info.java 1 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/wire/ByteStringCoder.java 68 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/wire/LengthPrefixUnknownCoders.java 70 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/wire/WireCoders.java 81 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/wire/package-info.java 1 runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/InMemoryJobService.java 422 runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/JobInvocation.java 217 runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/JobInvoker.java 34 runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/JobPreparation.java 20 runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/JobServerDriver.java 208 runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/PortablePipelineJarCreator.java 184 runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/PortablePipelineJarUtils.java 81 runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/PortablePipelineResult.java 8 runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/PortablePipelineRunner.java 6 runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/package-info.java 1 runners/jet/src/main/java/org/apache/beam/runners/jet/DAGBuilder.java 173 runners/jet/src/main/java/org/apache/beam/runners/jet/FailedRunningPipelineResults.java 58 runners/jet/src/main/java/org/apache/beam/runners/jet/JetGraphVisitor.java 72 runners/jet/src/main/java/org/apache/beam/runners/jet/JetPipelineOptions.java 36 runners/jet/src/main/java/org/apache/beam/runners/jet/JetPipelineResult.java 101 runners/jet/src/main/java/org/apache/beam/runners/jet/JetRunner.java 170 runners/jet/src/main/java/org/apache/beam/runners/jet/JetRunnerRegistrar.java 24 runners/jet/src/main/java/org/apache/beam/runners/jet/JetTransformTranslator.java 16 runners/jet/src/main/java/org/apache/beam/runners/jet/JetTransformTranslators.java 370 runners/jet/src/main/java/org/apache/beam/runners/jet/JetTranslationContext.java 27 runners/jet/src/main/java/org/apache/beam/runners/jet/Utils.java 231 runners/jet/src/main/java/org/apache/beam/runners/jet/metrics/AbstractMetric.java 12 runners/jet/src/main/java/org/apache/beam/runners/jet/metrics/CounterImpl.java 29 runners/jet/src/main/java/org/apache/beam/runners/jet/metrics/DistributionImpl.java 25 runners/jet/src/main/java/org/apache/beam/runners/jet/metrics/GaugeImpl.java 18 runners/jet/src/main/java/org/apache/beam/runners/jet/metrics/JetMetricResults.java 162 runners/jet/src/main/java/org/apache/beam/runners/jet/metrics/JetMetricsContainer.java 100 runners/jet/src/main/java/org/apache/beam/runners/jet/metrics/package-info.java 1 runners/jet/src/main/java/org/apache/beam/runners/jet/package-info.java 1 runners/jet/src/main/java/org/apache/beam/runners/jet/processors/AbstractParDoP.java 462 runners/jet/src/main/java/org/apache/beam/runners/jet/processors/AssignWindowP.java 98 runners/jet/src/main/java/org/apache/beam/runners/jet/processors/BoundedSourceP.java 171 runners/jet/src/main/java/org/apache/beam/runners/jet/processors/FlattenP.java 59 runners/jet/src/main/java/org/apache/beam/runners/jet/processors/ImpulseP.java 77 runners/jet/src/main/java/org/apache/beam/runners/jet/processors/ParDoP.java 162 runners/jet/src/main/java/org/apache/beam/runners/jet/processors/StatefulParDoP.java 303 runners/jet/src/main/java/org/apache/beam/runners/jet/processors/UnboundedSourceP.java 229 runners/jet/src/main/java/org/apache/beam/runners/jet/processors/ViewP.java 98 runners/jet/src/main/java/org/apache/beam/runners/jet/processors/WindowGroupP.java 287 runners/jet/src/main/java/org/apache/beam/runners/jet/processors/package-info.java 1 runners/local-java/src/main/java/org/apache/beam/runners/local/Bundle.java 11 runners/local-java/src/main/java/org/apache/beam/runners/local/ExecutionDriver.java 16 runners/local-java/src/main/java/org/apache/beam/runners/local/PipelineMessageReceiver.java 7 runners/local-java/src/main/java/org/apache/beam/runners/local/StructuralKey.java 60 runners/local-java/src/main/java/org/apache/beam/runners/local/package-info.java 1 runners/portability/java/src/main/java/org/apache/beam/runners/portability/CloseableResource.java 54 runners/portability/java/src/main/java/org/apache/beam/runners/portability/ExternalWorkerService.java 68 runners/portability/java/src/main/java/org/apache/beam/runners/portability/JobServicePipelineResult.java 167 runners/portability/java/src/main/java/org/apache/beam/runners/portability/PortableMetrics.java 128 runners/portability/java/src/main/java/org/apache/beam/runners/portability/PortableRunner.java 163 runners/portability/java/src/main/java/org/apache/beam/runners/portability/PortableRunnerRegistrar.java 12 runners/portability/java/src/main/java/org/apache/beam/runners/portability/package-info.java 1 runners/portability/java/src/main/java/org/apache/beam/runners/portability/testing/TestJobService.java 59 runners/portability/java/src/main/java/org/apache/beam/runners/portability/testing/TestPortablePipelineOptions.java 35 runners/portability/java/src/main/java/org/apache/beam/runners/portability/testing/TestPortableRunner.java 56 runners/portability/java/src/main/java/org/apache/beam/runners/portability/testing/TestUniversalRunner.java 77 runners/portability/java/src/main/java/org/apache/beam/runners/portability/testing/package-info.java 1 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaExecutionContext.java 41 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaExecutionEnvironment.java 6 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaJobInvoker.java 53 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaJobServerDriver.java 62 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaPipelineLifeCycleListener.java 12 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaPipelineOptions.java 80 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaPipelineOptionsValidator.java 28 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaPipelineResult.java 121 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaPipelineRunner.java 39 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaPortablePipelineOptions.java 22 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaPortablePipelineResult.java 25 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaRunner.java 138 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaRunnerOverrideConfigs.java 49 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaRunnerRegistrar.java 24 runners/samza/src/main/java/org/apache/beam/runners/samza/TestSamzaRunner.java 61 runners/samza/src/main/java/org/apache/beam/runners/samza/adapter/BoundedSourceSystem.java 355 runners/samza/src/main/java/org/apache/beam/runners/samza/adapter/UnboundedSourceSystem.java 424 runners/samza/src/main/java/org/apache/beam/runners/samza/adapter/package-info.java 1 runners/samza/src/main/java/org/apache/beam/runners/samza/container/BeamContainerRunner.java 55 runners/samza/src/main/java/org/apache/beam/runners/samza/container/BeamJobCoordinatorRunner.java 44 runners/samza/src/main/java/org/apache/beam/runners/samza/container/ContainerCfgLoader.java 39 runners/samza/src/main/java/org/apache/beam/runners/samza/container/ContainerCfgLoaderFactory.java 10 runners/samza/src/main/java/org/apache/beam/runners/samza/container/package-info.java 1 runners/samza/src/main/java/org/apache/beam/runners/samza/metrics/DoFnRunnerWithMetrics.java 69 runners/samza/src/main/java/org/apache/beam/runners/samza/metrics/FnWithMetricsWrapper.java 24 runners/samza/src/main/java/org/apache/beam/runners/samza/metrics/SamzaMetricsContainer.java 74 runners/samza/src/main/java/org/apache/beam/runners/samza/metrics/package-info.java 1 runners/samza/src/main/java/org/apache/beam/runners/samza/package-info.java 1 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/BundleManager.java 212 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/DoFnOp.java 493 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/DoFnRunnerWithKeyedInternals.java 93 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/FutureCollector.java 10 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/GroupByKeyOp.java 190 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/KeyedInternals.java 128 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/KeyedTimerData.java 155 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/KvToKeyedWorkItemOp.java 16 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/Op.java 25 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/OpAdapter.java 148 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/OpEmitter.java 11 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/OpMessage.java 109 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/OutputManagerFactory.java 10 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/SamzaAssignContext.java 32 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/SamzaDoFnInvokerRegistrar.java 11 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/SamzaDoFnRunners.java 256 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/SamzaExecutableStageContextFactory.java 28 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/SamzaStoreStateInternals.java 912 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/SamzaTimerInternalsFactory.java 543 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/SingletonKeyedWorkItem.java 25 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/SplittableParDoProcessKeyedElementsOp.java 215 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/WindowAssignOp.java 29 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/package-info.java 1 runners/samza/src/main/java/org/apache/beam/runners/samza/state/SamzaMapState.java 8 runners/samza/src/main/java/org/apache/beam/runners/samza/state/SamzaSetState.java 7 runners/samza/src/main/java/org/apache/beam/runners/samza/state/package-info.java 1 runners/samza/src/main/java/org/apache/beam/runners/samza/transforms/GroupWithoutRepartition.java 31 runners/samza/src/main/java/org/apache/beam/runners/samza/transforms/UpdatingCombineFn.java 6 runners/samza/src/main/java/org/apache/beam/runners/samza/transforms/package-info.java 1 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/ConfigBuilder.java 243 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/ConfigContext.java 50 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/FlattenPCollectionsTranslator.java 74 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/GroupByKeyTranslator.java 217 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/ImpulseTranslator.java 47 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/PViewToIdMapper.java 48 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/ParDoBoundMultiTranslator.java 315 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/PortableTranslationContext.java 150 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/ReadTranslator.java 61 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaImpulseSystemFactory.java 105 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaPipelineTranslator.java 159 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaPortablePipelineTranslator.java 57 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaPublishView.java 35 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaPublishViewTransformOverride.java 74 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaPublishViewTranslator.java 41 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaTestStreamSystemFactory.java 132 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaTestStreamTranslator.java 66 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaTransformOverrides.java 34 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaTranslatorRegistrar.java 5 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SplittableParDoTranslators.java 115 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/TransformConfigGenerator.java 17 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/TransformTranslator.java 18 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/TranslationContext.java 209 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/WindowAssignTranslator.java 48 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/package-info.java 1 runners/samza/src/main/java/org/apache/beam/runners/samza/util/FutureUtils.java 21 runners/samza/src/main/java/org/apache/beam/runners/samza/util/HashIdGenerator.java 34 runners/samza/src/main/java/org/apache/beam/runners/samza/util/SamzaCoders.java 52 runners/samza/src/main/java/org/apache/beam/runners/samza/util/SamzaPipelineTranslatorUtils.java 64 runners/samza/src/main/java/org/apache/beam/runners/samza/util/package-info.java 1 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptions.java 9 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineResult.java 108 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingRunner.java 130 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingRunnerRegistrar.java 24 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/aggregators/AggregatorsAccumulator.java 41 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/aggregators/NamedAggregators.java 49 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/aggregators/NamedAggregatorsAccumulator.java 35 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/aggregators/package-info.java 1 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/AggregatorMetric.java 15 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/AggregatorMetricSource.java 21 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/CompositeSource.java 22 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/MetricsAccumulator.java 44 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/MetricsContainerStepMapAccumulator.java 37 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/SparkBeamMetric.java 60 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/SparkBeamMetricSource.java 20 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/SparkMetricsContainerStepMap.java 17 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/WithMetricsSupport.java 130 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/package-info.java 1 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/sink/CodahaleCsvSink.java 14 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/sink/CodahaleGraphiteSink.java 14 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/sink/package-info.java 1 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/package-info.java 1 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/PipelineTranslator.java 111 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/SchemaHelpers.java 15 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/SparkTransformOverrides.java 34 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/TransformTranslator.java 9 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/TranslationContext.java 174 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/AggregatorCombiner.java 212 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/CombinePerKeyTranslatorBatch.java 86 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/CreatePCollectionViewTranslatorBatch.java 37 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DatasetSourceBatch.java 117 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DoFnFunction.java 124 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DoFnRunnerWithMetrics.java 76 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/FlattenTranslatorBatch.java 43 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/GroupByKeyTranslatorBatch.java 58 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/ImpulseTranslatorBatch.java 29 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/ParDoTranslatorBatch.java 200 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/PipelineTranslatorBatch.java 51 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/ProcessContext.java 92 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/ReadSourceTranslatorBatch.java 58 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/ReshuffleTranslatorBatch.java 8 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/WindowAssignTranslatorBatch.java 39 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/functions/GroupAlsoByWindowViaOutputBufferFn.java 124 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/functions/NoOpStepContext.java 14 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/functions/SparkSideInputReader.java 150 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/functions/package-info.java 1 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/package-info.java 1 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/CoderHelpers.java 25 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/EncoderHelpers.java 212 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/KVHelpers.java 9 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/MultiOutputCoder.java 55 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/RowHelpers.java 38 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/SideInputBroadcast.java 24 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/WindowingHelpers.java 50 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/package-info.java 1 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/package-info.java 1 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/DatasetSourceStreaming.java 187 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/PipelineTranslatorStreaming.java 36 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/ReadSourceTranslatorStreaming.java 59 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/streaming/package-info.java 1 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/utils/CachedSideInputReader.java 56 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/utils/SideInputStorage.java 62 runners/spark/2/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/utils/package-info.java 1 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkCommonPipelineOptions.java 40 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkContextOptions.java 27 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkJobInvoker.java 75 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkJobServerDriver.java 70 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkNativePipelineVisitor.java 163 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkPipelineOptions.java 56 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkPipelineResult.java 175 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkPipelineRunner.java 227 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkPortableStreamingPipelineOptions.java 15 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkRunner.java 337 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkRunnerDebugger.java 91 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkRunnerRegistrar.java 25 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkTransformOverrides.java 34 runners/spark/src/main/java/org/apache/beam/runners/spark/TestSparkPipelineOptions.java 30 runners/spark/src/main/java/org/apache/beam/runners/spark/TestSparkRunner.java 99 runners/spark/src/main/java/org/apache/beam/runners/spark/aggregators/AggregatorsAccumulator.java 97 runners/spark/src/main/java/org/apache/beam/runners/spark/aggregators/NamedAggregators.java 49 runners/spark/src/main/java/org/apache/beam/runners/spark/aggregators/NamedAggregatorsAccumulator.java 35 runners/spark/src/main/java/org/apache/beam/runners/spark/aggregators/package-info.java 1 runners/spark/src/main/java/org/apache/beam/runners/spark/coders/CoderHelpers.java 121 runners/spark/src/main/java/org/apache/beam/runners/spark/coders/SparkRunnerKryoRegistrator.java 44 runners/spark/src/main/java/org/apache/beam/runners/spark/coders/StatelessJavaSerializer.java 54 runners/spark/src/main/java/org/apache/beam/runners/spark/coders/package-info.java 1 runners/spark/src/main/java/org/apache/beam/runners/spark/io/ConsoleIO.java 31 runners/spark/src/main/java/org/apache/beam/runners/spark/io/CreateStream.java 124 runners/spark/src/main/java/org/apache/beam/runners/spark/io/EmptyCheckpointMark.java 23 runners/spark/src/main/java/org/apache/beam/runners/spark/io/MicrobatchSource.java 259 runners/spark/src/main/java/org/apache/beam/runners/spark/io/SourceDStream.java 150 runners/spark/src/main/java/org/apache/beam/runners/spark/io/SourceRDD.java 265 runners/spark/src/main/java/org/apache/beam/runners/spark/io/SparkUnboundedSource.java 231 runners/spark/src/main/java/org/apache/beam/runners/spark/io/package-info.java 1 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/AggregatorMetric.java 15 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/AggregatorMetricSource.java 21 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/CompositeSource.java 22 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/MetricsAccumulator.java 101 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/MetricsContainerStepMapAccumulator.java 37 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/SparkBeamMetric.java 72 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/SparkBeamMetricSource.java 20 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/SparkMetricsContainerStepMap.java 17 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/WithMetricsSupport.java 130 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/package-info.java 1 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/sink/CsvSink.java 16 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/sink/GraphiteSink.java 16 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/sink/package-info.java 1 runners/spark/src/main/java/org/apache/beam/runners/spark/package-info.java 1 runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkGroupAlsoByWindowViaWindowSet.java 408 runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkStateInternals.java 316 runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkTimerInternals.java 147 runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/StateSpecFunctions.java 139 runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/package-info.java 1 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/BoundedDataset.java 95 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/Dataset.java 8 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/DoFnRunnerWithMetrics.java 76 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/EvaluationContext.java 181 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/GroupCombineFunctions.java 124 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/GroupNonMergingWindowsFunctions.java 190 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/MultiDoFnFunction.java 199 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/ReifyTimestampsAndWindowsFunction.java 15 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkAssignWindowFn.java 42 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkBatchPortablePipelineTranslator.java 340 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkCombineFn.java 656 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkContextFactory.java 66 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkExecutableStageContextFactory.java 29 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkExecutableStageExtractionFunction.java 24 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkExecutableStageFunction.java 258 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkGroupAlsoByWindowViaOutputBufferFn.java 122 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkPCollectionView.java 72 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkPipelineTranslator.java 9 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkPortablePipelineTranslator.java 11 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkProcessContext.java 125 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkStreamingPortablePipelineTranslator.java 297 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkStreamingTranslationContext.java 24 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkTranslationContext.java 70 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/TransformEvaluator.java 7 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/TransformTranslator.java 669 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/TranslationUtils.java 260 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/ValueAndCoderKryoSerializer.java 28 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/ValueAndCoderLazySerializable.java 90 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/package-info.java 1 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/Checkpoint.java 96 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/SparkRunnerStreamingContextFactory.java 76 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/StreamingTransformTranslator.java 525 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/UnboundedDataset.java 47 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/WatermarkSyncedDStream.java 88 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/package-info.java 1 runners/spark/src/main/java/org/apache/beam/runners/spark/util/ByteArray.java 30 runners/spark/src/main/java/org/apache/beam/runners/spark/util/CachedSideInputReader.java 56 runners/spark/src/main/java/org/apache/beam/runners/spark/util/GlobalWatermarkHolder.java 256 runners/spark/src/main/java/org/apache/beam/runners/spark/util/SideInputBroadcast.java 52 runners/spark/src/main/java/org/apache/beam/runners/spark/util/SideInputStorage.java 62 runners/spark/src/main/java/org/apache/beam/runners/spark/util/SparkCommon.java 51 runners/spark/src/main/java/org/apache/beam/runners/spark/util/SparkCompat.java 127 runners/spark/src/main/java/org/apache/beam/runners/spark/util/SparkSideInputReader.java 82 runners/spark/src/main/java/org/apache/beam/runners/spark/util/TimerUtils.java 22 runners/spark/src/main/java/org/apache/beam/runners/spark/util/package-info.java 1 runners/twister2/src/main/java/org/apache/beam/runners/twister2/BeamBatchTSetEnvironment.java 16 runners/twister2/src/main/java/org/apache/beam/runners/twister2/BeamBatchWorker.java 117 runners/twister2/src/main/java/org/apache/beam/runners/twister2/Twister2BatchTranslationContext.java 21 runners/twister2/src/main/java/org/apache/beam/runners/twister2/Twister2PipelineExecutionEnvironment.java 81 runners/twister2/src/main/java/org/apache/beam/runners/twister2/Twister2PipelineOptions.java 41 runners/twister2/src/main/java/org/apache/beam/runners/twister2/Twister2PipelineResult.java 54 runners/twister2/src/main/java/org/apache/beam/runners/twister2/Twister2Runner.java 269 runners/twister2/src/main/java/org/apache/beam/runners/twister2/Twister2RunnerRegistrar.java 24 runners/twister2/src/main/java/org/apache/beam/runners/twister2/Twister2StreamTranslationContext.java 11 runners/twister2/src/main/java/org/apache/beam/runners/twister2/Twister2TestRunner.java 35 runners/twister2/src/main/java/org/apache/beam/runners/twister2/Twister2TranslationContext.java 90 runners/twister2/src/main/java/org/apache/beam/runners/twister2/package-info.java 1 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translation/wrappers/Twister2BoundedSource.java 210 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translation/wrappers/Twister2EmptySource.java 16 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translation/wrappers/package-info.java 1 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/BatchTransformTranslator.java 9 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/StreamTransformTranslator.java 9 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/Twister2BatchPipelineTranslator.java 65 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/Twister2PipelineTranslator.java 7 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/Twister2StreamPipelineTranslator.java 10 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/batch/AssignWindowTranslatorBatch.java 28 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/batch/FlattenTranslatorBatch.java 48 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/batch/GroupByKeyTranslatorBatch.java 53 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/batch/ImpulseTranslatorBatch.java 20 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/batch/PCollectionViewTranslatorBatch.java 60 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/batch/ParDoMultiOutputTranslatorBatch.java 98 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/batch/ReadSourceTranslatorBatch.java 27 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/batch/package-info.java 1 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/AssignWindowsFunction.java 87 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/ByteToWindowFunction.java 74 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/ByteToWindowFunctionPrimitive.java 76 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/DoFnFunction.java 293 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/GroupByWindowFunction.java 186 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/ImpulseSource.java 20 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/MapToTupleFunction.java 75 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/OutputTagFilter.java 32 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/Twister2SinkFunction.java 23 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/internal/SystemReduceFnBuffering.java 82 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/internal/package-info.java 1 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/package-info.java 1 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/package-info.java 1 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/streaming/ReadSourceTranslatorStream.java 12 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/streaming/package-info.java 1 runners/twister2/src/main/java/org/apache/beam/runners/twister2/utils/NoOpStepContext.java 15 runners/twister2/src/main/java/org/apache/beam/runners/twister2/utils/TranslationUtils.java 31 runners/twister2/src/main/java/org/apache/beam/runners/twister2/utils/Twister2AssignContext.java 27 runners/twister2/src/main/java/org/apache/beam/runners/twister2/utils/Twister2SideInputReader.java 86 runners/twister2/src/main/java/org/apache/beam/runners/twister2/utils/package-info.java 1