Path Lines of Code 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/Concatenate.java 40 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 109 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 719 runners/core-java/src/main/java/org/apache/beam/runners/core/InMemoryTimerInternals.java 253 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 31 runners/core-java/src/main/java/org/apache/beam/runners/core/OutputAndTimeBoundedSplittableProcessElementInvoker.java 324 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 103 runners/core-java/src/main/java/org/apache/beam/runners/core/PushbackSideInputDoFnRunner.java 21 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 1122 runners/core-java/src/main/java/org/apache/beam/runners/core/SimplePushbackSideInputDoFnRunner.java 95 runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableParDoViaKeyedWorkItems.java 544 runners/core-java/src/main/java/org/apache/beam/runners/core/SplittableProcessElementInvoker.java 53 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 48 runners/core-java/src/main/java/org/apache/beam/runners/core/StateTags.java 271 runners/core-java/src/main/java/org/apache/beam/runners/core/StatefulDoFnRunner.java 301 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 192 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/construction/SerializablePipelineOptions.java 69 runners/core-java/src/main/java/org/apache/beam/runners/core/construction/package-info.java 4 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/BoundedTrieCell.java 68 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/BoundedTrieData.java 370 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 52 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/DirtyState.java 41 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 34 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/ExecutionStateSampler.java 113 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/ExecutionStateTracker.java 213 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/GcpResourceIdentifiers.java 35 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/HistogramCell.java 74 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 58 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MetricsContainerImpl.java 691 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MetricsContainerStepMap.java 186 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MetricsLogger.java 57 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MetricsMap.java 62 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 182 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoEncodings.java 147 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/MonitoringInfoMetricName.java 79 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/ServiceCallMetric.java 65 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/ShortIdMap.java 57 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 112 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/SimpleStateRegistry.java 53 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/SpecMonitoringInfoValidator.java 51 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/StringSetCell.java 72 runners/core-java/src/main/java/org/apache/beam/runners/core/metrics/StringSetData.java 103 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 42 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 194 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 187 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 360 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 84 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 294 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 294 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 108 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 93 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 76 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 260 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 338 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 120 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 418 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 263 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 190 runners/direct-java/src/main/java/org/apache/beam/runners/direct/StatefulParDoEvaluatorFactory.java 219 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 193 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 138 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 54 runners/direct-java/src/main/java/org/apache/beam/runners/direct/WatermarkCallbackExecutor.java 140 runners/direct-java/src/main/java/org/apache/beam/runners/direct/WatermarkManager.java 1165 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.17/src/main/java/org/apache/beam/runners/flink/translation/types/CoderTypeSerializer.java 119 runners/flink/src/main/java/org/apache/beam/runners/flink/CreateStreamingFlinkView.java 99 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 460 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkBatchTransformTranslators.java 686 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 105 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkExecutionEnvironments.java 388 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 126 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkPipelineOptions.java 262 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkPipelineRunner.java 163 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 202 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkPortablePipelineTranslator.java 45 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkPortableRunnerResult.java 57 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkRunner.java 107 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/FlinkStreamingAggregationsTranslators.java 430 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 909 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkStreamingTransformTranslators.java 1159 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkStreamingTranslationContext.java 102 runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkTransformOverrides.java 44 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/VersionDependentFlinkPipelineOptions.java 12 runners/flink/src/main/java/org/apache/beam/runners/flink/adapter/BeamAdapterCoderUtils.java 72 runners/flink/src/main/java/org/apache/beam/runners/flink/adapter/BeamAdapterUtils.java 115 runners/flink/src/main/java/org/apache/beam/runners/flink/adapter/BeamFlinkDataSetAdapter.java 221 runners/flink/src/main/java/org/apache/beam/runners/flink/adapter/BeamFlinkDataStreamAdapter.java 254 runners/flink/src/main/java/org/apache/beam/runners/flink/adapter/FlinkInput.java 71 runners/flink/src/main/java/org/apache/beam/runners/flink/adapter/FlinkKey.java 54 runners/flink/src/main/java/org/apache/beam/runners/flink/adapter/FlinkOutput.java 61 runners/flink/src/main/java/org/apache/beam/runners/flink/adapter/package-info.java 1 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 57 runners/flink/src/main/java/org/apache/beam/runners/flink/metrics/FlinkMetricContainer.java 30 runners/flink/src/main/java/org/apache/beam/runners/flink/metrics/FlinkMetricContainerBase.java 122 runners/flink/src/main/java/org/apache/beam/runners/flink/metrics/FlinkMetricContainerWithoutAccumulator.java 7 runners/flink/src/main/java/org/apache/beam/runners/flink/metrics/Metrics.java 37 runners/flink/src/main/java/org/apache/beam/runners/flink/metrics/MetricsAccumulator.java 34 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 214 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 90 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/UnversionedTypeSerializerSnapshot.java 55 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 40 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/utils/Locker.java 17 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/SerdeUtils.java 58 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 134 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 1251 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/ExecutableStageDoFnOperator.java 1029 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/KvToFlinkKeyKeySelector.java 25 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/PartialReduceBundleOperator.java 144 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/SdfFlinkKeyKeySelector.java 26 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/SingletonKeyedWorkItem.java 32 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 164 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/WindowDoFnOperator.java 108 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/WorkItemKeySelector.java 27 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 45 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 366 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/io/source/FlinkSource.java 110 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/FlinkSourceReaderBase.java 275 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/FlinkSourceSplit.java 59 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/FlinkSourceSplitEnumerator.java 153 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/LazyFlinkSourceSplitEnumerator.java 139 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/bounded/FlinkBoundedSource.java 37 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/bounded/FlinkBoundedSourceReader.java 118 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/bounded/package-info.java 1 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/impulse/BeamImpulseSource.java 68 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/impulse/package-info.java 1 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/package-info.java 1 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/unbounded/FlinkUnboundedSource.java 42 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/unbounded/FlinkUnboundedSourceReader.java 295 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/io/source/unbounded/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 159 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/stableinput/BufferingDoFnRunner.java 265 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 564 runners/flink/src/main/java/org/apache/beam/runners/flink/translation/wrappers/streaming/state/FlinkStateInternals.java 1587 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 914 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 297 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 381 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 1017 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java 2216 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunnerHooks.java 5 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunnerInfo.java 107 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 284 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/PrimitiveParDoSingleFactory.java 317 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/RedistributeByKeyOverrideFactory.java 118 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 96 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 296 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 76 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/internal/DataflowGroupByKey.java 111 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/internal/IsmFormat.java 505 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/DataflowPipelineDebugOptions.java 156 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowPipelineOptions.java 149 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowPipelineWorkerPoolOptions.java 98 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowProfilingOptions.java 20 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/options/DataflowStreamingPipelineOptions.java 197 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 72 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 59 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 11 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 532 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/CloudObjects.java 101 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 50 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/DataflowTransport.java 64 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/DefaultCoderCloudObjectTranslatorRegistrar.java 112 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 172 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 56 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 44 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/SchemaCoderCloudObjectTranslator.java 128 runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/SerializableCoderCloudObjectTranslator.java 42 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/org/apache/beam/runners/dataflow/worker/ActiveMessageMetadata.java 11 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 92 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 219 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/BatchModeExecutionContext.java 447 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/ByteArrayReader.java 28 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 61 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 12 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 111 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 133 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 345 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/DataflowExecutionStateSampler.java 95 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 19 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowMetricsContainer.java 56 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOperationContext.java 272 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowOutputCounter.java 58 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/DataflowSideInputReadCounter.java 91 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowSystemMetrics.java 48 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 120 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/DataflowWorkUnitClient.java 248 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 94 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/Filepatterns.java 48 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 301 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 311 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 920 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/IsmSideInputReader.java 793 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/LockFreeHistogram.java 136 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 133 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/MetricsToPerStepNamespaceMetricsConverter.java 237 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/OperationalLimits.java 23 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/OutputTooLargeException.java 18 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/PubsubDynamicSink.java 128 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PubsubReader.java 112 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/PubsubSink.java 163 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ReaderCache.java 90 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/RemoveSafeDeltaCounterCell.java 69 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/ShuffleLibrary.java 32 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ShuffleReader.java 15 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/ShuffleSink.java 262 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 394 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 24 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 190 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StackTraceUtil.java 43 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingDataflowWorker.java 867 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 771 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 69 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingSideInputFetcher.java 310 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/StreamingStepMetricsContainer.java 322 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 96 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 108 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 28 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 224 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimeUtils.java 37 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WindmillTimerInternals.java 333 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/WorkItemCancelledException.java 24 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WorkItemStatusClient.java 285 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WorkUnitClient.java 20 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 700 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/WorkerPipelineOptionsFactory.java 55 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 638 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/DataflowCounterUpdateExtractor.java 170 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/counters/NameContext.java 39 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/LengthPrefixUnknownCoders.java 200 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 144 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/graph/Nodes.java 193 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingHandler.java 236 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/logging/DataflowWorkerLoggingInitializer.java 205 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 351 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/JfrzServlet.java 39 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 63 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 101 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/streaming/ActiveWorkState.java 267 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationState.java 110 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationStateCache.java 168 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ComputationWorkExecutor.java 68 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ExecutableWork.java 22 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/KeyCommitTooLargeException.java 30 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/RefreshableWork.java 15 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/ShardedKey.java 19 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/StageInfo.java 99 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/Watermarks.java 38 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/WeightedBoundedQueue.java 44 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/WeightedSemaphore.java 28 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/Work.java 319 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/WorkHeartbeatResponseProcessor.java 39 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/WorkId.java 23 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/WorkIdWithShardingKey.java 10 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/config/ComputationConfig.java 32 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/config/FakeGlobalConfigHandle.java 25 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/config/FixedGlobalConfigHandle.java 21 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/config/StreamingApplianceComputationConfigFetcher.java 105 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/config/StreamingEngineComputationConfigFetcher.java 254 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/config/StreamingGlobalConfig.java 27 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/config/StreamingGlobalConfigHandle.java 10 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/config/StreamingGlobalConfigHandleImpl.java 73 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/FanOutStreamingEngineWorkerHarness.java 375 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/GlobalDataStreamSender.java 36 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/MetricsDataProvider.java 56 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/SingleSourceWorkerHarness.java 241 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/StreamSender.java 4 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/StreamingCounters.java 63 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/StreamingEngineBackends.java 22 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/StreamingWorkerHarness.java 8 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/StreamingWorkerStatusPages.java 233 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/StreamingWorkerStatusReporter.java 387 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/harness/WindmillStreamSender.java 140 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/sideinput/SideInput.java 16 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/sideinput/SideInputCache.java 69 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/sideinput/SideInputState.java 6 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/sideinput/SideInputStateFetcher.java 198 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/streaming/sideinput/SideInputStateFetcherFactory.java 20 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 218 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/JfrInterop.java 53 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/MemoryMonitor.java 540 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 76 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 203 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 14 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/MapTaskExecutor.java 124 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 115 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/util/common/worker/OutputReceiver.java 44 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 315 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 81 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/ApplianceWindmillClient.java 10 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/StreamingEngineWindmillClient.java 16 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/WindmillConnection.java 40 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/WindmillEndpoints.java 163 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/WindmillServerBase.java 108 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/WindmillServerStub.java 17 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/WindmillServiceAddress.java 37 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/appliance/JniWindmillApplianceServer.java 30 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/AbstractWindmillStream.java 330 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/CloseableStream.java 17 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/ResettableThrowingStreamObserver.java 98 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/StreamDebugMetrics.java 150 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/WindmillStream.java 58 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/WindmillStreamPool.java 126 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/WindmillStreamShutdownException.java 6 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java 25 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commits.java 12 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/CompleteCommit.java 33 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitter.java 127 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitter.java 213 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/WorkCommitter.java 12 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/getdata/ApplianceGetDataClient.java 167 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/getdata/GetDataClient.java 22 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/getdata/StreamGetDataClient.java 63 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/getdata/StreamPoolGetDataClient.java 51 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/getdata/ThrottlingGetDataMetricTracker.java 61 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/AppendableInputStream.java 142 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/ChannelzServlet.java 244 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GetWorkResponseChunkAssembler.java 98 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GetWorkTimingInfosTracker.java 142 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcCommitWorkStream.java 368 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcDeadlineClientInterceptor.java 25 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcDirectGetWorkStream.java 292 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcDispatcherClient.java 214 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcGetDataStream.java 399 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcGetDataStreamRequests.java 205 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcGetWorkStream.java 173 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcGetWorkerMetadataStream.java 124 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillServer.java 330 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillStreamFactory.java 306 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/auth/VendoredCredentialsAdapter.java 45 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/auth/VendoredRequestMetadataCallbackAdapter.java 20 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/observers/DirectStreamObserver.java 146 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/observers/ForwardingClientResponseObserver.java 34 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/observers/StreamObserverCancelledException.java 14 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/observers/StreamObserverFactory.java 36 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/observers/TerminatingStreamObserver.java 7 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/ChannelCache.java 136 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/ChannelCachingRemoteStubFactory.java 50 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/ChannelCachingStubFactory.java 7 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/IsolationChannel.java 208 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/WindmillChannelFactory.java 10 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/WindmillChannels.java 140 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/WindmillStubFactory.java 11 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/WindmillStubFactoryFactory.java 6 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/stubs/WindmillStubFactoryFactoryImpl.java 31 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/throttling/StreamingEngineThrottleTimers.java 17 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/throttling/ThrottleTimer.java 30 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/throttling/ThrottledTimeTracker.java 7 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/AbstractWindmillMap.java 4 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/CachingStateTable.java 235 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/ConcatIterables.java 27 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/IdTracker.java 166 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/PagingIterable.java 76 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/RangeCoder.java 51 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/RangeSetCoder.java 24 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/SimpleWindmillState.java 14 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/StateTag.java 54 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/TimestampedValueWithId.java 19 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/ToIterableFunction.java 51 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/ValuesAndContPosition.java 18 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WeightedList.java 30 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillBag.java 147 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillCombiningState.java 127 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillMap.java 357 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillMapViaMultimap.java 120 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillMultimap.java 598 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillOrderedList.java 235 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillSet.java 87 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillState.java 29 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateCache.java 316 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateInternals.java 120 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateReader.java 813 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillStateUtil.java 21 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillValue.java 115 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WindmillWatermarkHold.java 189 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/state/WrappedFuture.java 31 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/WorkItemReceiver.java 16 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/WorkItemScheduler.java 19 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/budget/EvenGetWorkBudgetDistributor.java 32 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/budget/GetWorkBudget.java 53 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/budget/GetWorkBudgetDistributor.java 8 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/budget/GetWorkBudgetDistributors.java 8 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/budget/GetWorkBudgetRefresher.java 82 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/budget/GetWorkBudgetSpender.java 7 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/ComputationWorkExecutorFactory.java 240 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingCommitFinalizer.java 49 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/StreamingWorkScheduler.java 330 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/FailureTracker.java 63 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/HeapDumper.java 7 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/StreamingApplianceFailureTracker.java 37 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/StreamingApplianceStatsReporter.java 6 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/StreamingEngineFailureTracker.java 21 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/processing/failures/WorkFailureProcessor.java 152 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ActiveWorkRefresher.java 145 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/ApplianceHeartbeatSender.java 34 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/FixedStreamHeartbeatSender.java 56 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/HeartbeatSender.java 5 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/Heartbeats.java 42 runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/work/refresh/StreamPoolHeartbeatSender.java 49 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 851 runners/google-cloud-dataflow-java/worker/windmill/src/main/proto/windmill_service.proto 46 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 516 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 105 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 603 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 85 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 184 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 543 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 12 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/data/GrpcDataService.java 142 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 180 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 152 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 160 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 117 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 52 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 441 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/StreamingSideInputHandlerFactory.java 118 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 74 runners/java-fn-execution/src/main/java/org/apache/beam/runners/fnexecution/wire/LengthPrefixUnknownCoders.java 79 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 441 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 228 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 79 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 73 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 102 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/BoundedTrieImpl.java 23 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 250 runners/jet/src/main/java/org/apache/beam/runners/jet/metrics/JetMetricsContainer.java 143 runners/jet/src/main/java/org/apache/beam/runners/jet/metrics/StringSetImpl.java 26 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 167 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 78 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/JobServicePipelineResult.java 177 runners/portability/java/src/main/java/org/apache/beam/runners/portability/PortableMetrics.java 196 runners/portability/java/src/main/java/org/apache/beam/runners/portability/PortableRunner.java 173 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 53 runners/portability/java/src/main/java/org/apache/beam/runners/portability/testing/TestUniversalRunner.java 75 runners/portability/java/src/main/java/org/apache/beam/runners/portability/testing/package-info.java 1 runners/prism/java/src/main/java/org/apache/beam/runners/prism/PrismExecutor.java 120 runners/prism/java/src/main/java/org/apache/beam/runners/prism/PrismLocator.java 183 runners/prism/java/src/main/java/org/apache/beam/runners/prism/PrismPipelineOptions.java 34 runners/prism/java/src/main/java/org/apache/beam/runners/prism/PrismPipelineResult.java 42 runners/prism/java/src/main/java/org/apache/beam/runners/prism/PrismRegistrar.java 24 runners/prism/java/src/main/java/org/apache/beam/runners/prism/PrismRunner.java 88 runners/prism/java/src/main/java/org/apache/beam/runners/prism/TestPrismPipelineOptions.java 3 runners/prism/java/src/main/java/org/apache/beam/runners/prism/TestPrismRunner.java 48 runners/prism/java/src/main/java/org/apache/beam/runners/prism/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 52 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaJobServerDriver.java 68 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaPipelineExceptionContext.java 15 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 118 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaPipelineOptionsValidator.java 26 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaPipelineResult.java 116 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaPipelineRunner.java 53 runners/samza/src/main/java/org/apache/beam/runners/samza/SamzaPortablePipelineOptions.java 11 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 153 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 75 runners/samza/src/main/java/org/apache/beam/runners/samza/adapter/BoundedSourceSystem.java 350 runners/samza/src/main/java/org/apache/beam/runners/samza/adapter/UnboundedSourceSystem.java 420 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 67 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/SamzaGBKMetricOp.java 130 runners/samza/src/main/java/org/apache/beam/runners/samza/metrics/SamzaMetricOp.java 117 runners/samza/src/main/java/org/apache/beam/runners/samza/metrics/SamzaMetricOpFactory.java 29 runners/samza/src/main/java/org/apache/beam/runners/samza/metrics/SamzaMetricsContainer.java 81 runners/samza/src/main/java/org/apache/beam/runners/samza/metrics/SamzaTransformMetricRegistry.java 150 runners/samza/src/main/java/org/apache/beam/runners/samza/metrics/SamzaTransformMetrics.java 97 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/AsyncDoFnRunner.java 124 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/BundleManager.java 14 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/ClassicBundleManager.java 212 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/DoFnOp.java 470 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/DoFnRunnerWithKeyedInternals.java 83 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/FutureCollector.java 11 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/FutureCollectorImpl.java 64 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/GroupByKeyOp.java 196 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/KeyedInternals.java 133 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 187 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/OpEmitter.java 14 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/PortableBundleManager.java 132 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/PortableDoFnOp.java 381 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 425 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/SamzaMetricsBundleProgressHandler.java 77 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/SamzaStateRequestHandlers.java 131 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/SamzaStoreStateInternals.java 936 runners/samza/src/main/java/org/apache/beam/runners/samza/runtime/SamzaTimerInternalsFactory.java 549 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 223 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 9 runners/samza/src/main/java/org/apache/beam/runners/samza/state/SamzaSetState.java 8 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 273 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/ConfigContext.java 51 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/FlattenPCollectionsTranslator.java 73 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/GroupByKeyTranslator.java 219 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 452 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/PortableTranslationContext.java 78 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/RedistributeByKeyTranslator.java 36 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/ReshuffleTranslator.java 90 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 162 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaPortablePipelineTranslator.java 75 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaPortableTranslatorRegistrar.java 5 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 39 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaPublishViewTranslator.java 44 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaTestStreamSystemFactory.java 131 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/SamzaTestStreamTranslator.java 110 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 116 runners/samza/src/main/java/org/apache/beam/runners/samza/translation/StateIdParser.java 42 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 302 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/ConfigUtils.java 17 runners/samza/src/main/java/org/apache/beam/runners/samza/util/DoFnUtils.java 39 runners/samza/src/main/java/org/apache/beam/runners/samza/util/FutureUtils.java 36 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/PipelineJsonRenderer.java 240 runners/samza/src/main/java/org/apache/beam/runners/samza/util/PortableConfigUtils.java 14 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/SamzaPipelineExceptionListener.java 9 runners/samza/src/main/java/org/apache/beam/runners/samza/util/SamzaPipelineTranslatorUtils.java 32 runners/samza/src/main/java/org/apache/beam/runners/samza/util/StateUtils.java 16 runners/samza/src/main/java/org/apache/beam/runners/samza/util/StoreIdGenerator.java 19 runners/samza/src/main/java/org/apache/beam/runners/samza/util/WindowUtils.java 44 runners/samza/src/main/java/org/apache/beam/runners/samza/util/package-info.java 1 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineOptions.java 14 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingPipelineResult.java 92 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingRunner.java 125 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/SparkStructuredStreamingRunnerRegistrar.java 24 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/io/BoundedDatasetFactory.java 253 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/io/package-info.java 1 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/BeamMetricSet.java 25 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/MetricsAccumulator.java 84 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/SparkBeamMetric.java 81 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/SparkBeamMetricSource.java 19 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/WithMetricsSupport.java 55 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/package-info.java 1 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/sink/CodahaleCsvSink.java 44 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/sink/CodahaleGraphiteSink.java 44 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/metrics/sink/package-info.java 1 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/package-info.java 1 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/EvaluationContext.java 69 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/PipelineTranslator.java 343 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/SparkSessionFactory.java 222 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/SparkTransformOverrides.java 34 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/TransformTranslator.java 156 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/Aggregators.java 417 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/CombineGloballyTranslatorBatch.java 73 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/CombineGroupedValuesTranslatorBatch.java 47 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/CombinePerKeyTranslatorBatch.java 105 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DoFnPartitionIteratorFactory.java 131 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DoFnRunnerFactory.java 220 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/DoFnRunnerWithMetrics.java 77 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/FlattenTranslatorBatch.java 42 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/GroupByKeyHelpers.java 54 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/GroupByKeyTranslatorBatch.java 196 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/ImpulseTranslatorBatch.java 24 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/ParDoTranslatorBatch.java 194 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/PipelineTranslatorBatch.java 47 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/ReadSourceTranslatorBatch.java 33 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/ReshuffleTranslatorBatch.java 32 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/WindowAssignTranslatorBatch.java 67 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/functions/CachedSideInputReader.java 121 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/functions/GroupAlsoByWindowViaOutputBufferFn.java 125 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/functions/NoOpStepContext.java 14 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/functions/SideInputValues.java 125 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/functions/SparkSideInputReader.java 103 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/functions/package-info.java 1 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/batch/package-info.java 1 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/CoderHelpers.java 21 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/EncoderFactory.java 84 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/EncoderHelpers.java 425 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/EncoderProvider.java 32 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/helpers/package-info.java 1 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/package-info.java 1 runners/spark/3/src/main/java/org/apache/beam/runners/spark/structuredstreaming/translation/utils/ScalaInterop.java 73 runners/spark/3/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/DependentTransformsVisitor.java 22 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkCommonPipelineOptions.java 60 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 50 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 209 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkPortableStreamingPipelineOptions.java 13 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkRunner.java 309 runners/spark/src/main/java/org/apache/beam/runners/spark/SparkRunnerDebugger.java 87 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 47 runners/spark/src/main/java/org/apache/beam/runners/spark/TestSparkPipelineOptions.java 34 runners/spark/src/main/java/org/apache/beam/runners/spark/TestSparkRunner.java 100 runners/spark/src/main/java/org/apache/beam/runners/spark/coders/CoderHelpers.java 130 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 55 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 258 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/BeamMetricSet.java 24 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 78 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 55 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 44 runners/spark/src/main/java/org/apache/beam/runners/spark/metrics/sink/GraphiteSink.java 44 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 382 runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkStateInternals.java 519 runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/SparkTimerInternals.java 149 runners/spark/src/main/java/org/apache/beam/runners/spark/stateful/StateAndTimers.java 20 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 213 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/GroupByKeyVisitor.java 51 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/GroupCombineFunctions.java 147 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/GroupNonMergingWindowsFunctions.java 253 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/MultiDoFnFunction.java 200 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/SideInputMetadata.java 24 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SingleEmitInputDStream.java 26 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 345 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkCombineFn.java 652 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkContextFactory.java 87 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 251 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/SparkInputDataProcessor.java 268 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkPCollectionView.java 87 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 39 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/SparkStreamingPortablePipelineTranslator.java 296 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 751 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/TranslationUtils.java 287 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/CreateStreamingSparkView.java 99 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/ParDoStateUpdateFn.java 205 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/SparkRunnerStreamingContextFactory.java 82 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/StatefulStreamingParDoEvaluator.java 177 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/StreamingTransformTranslator.java 620 runners/spark/src/main/java/org/apache/beam/runners/spark/translation/streaming/TestDStream.java 102 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 86 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 60 runners/spark/src/main/java/org/apache/beam/runners/spark/util/SideInputReaderFactory.java 15 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/SparkSideInputReader.java 107 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 274 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 207 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 63 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 4 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 89 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 77 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/ByteToElemFunction.java 43 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/ElemToBytesFunction.java 51 runners/twister2/src/main/java/org/apache/beam/runners/twister2/translators/functions/GroupByWindowFunction.java 178 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 122 runners/twister2/src/main/java/org/apache/beam/runners/twister2/utils/package-info.java 1