Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

[Python] Managed Transforms API #31495

Merged
merged 31 commits into from
Oct 15, 2024

remove managed dependency

68dce26
Select commit
Loading
Failed to load commit list.
Sign in for the full log view
Merged

[Python] Managed Transforms API #31495

remove managed dependency
68dce26
Select commit
Loading
Failed to load commit list.
GitHub Actions / Test Results succeeded Aug 6, 2024 in 0s

All 373 tests pass, 1 skipped in 8m 29s

  153 files  +121    153 suites  +121   8m 29s ⏱️ - 30m 44s
  374 tests +239    373 ✅ +239  1 💤 ±0  0 ❌ ±0 
1 122 runs  +527  1 119 ✅ +534  3 💤  - 7  0 ❌ ±0 

Results for commit 68dce26. ± Comparison against earlier commit efd4c88.

Annotations

Check notice on line 0 in .github

See this annotation in the file changed.

@github-actions github-actions / Test Results

1 skipped test found

There is 1 skipped test, see "Raw output" for the name of the skipped test.
Raw output
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldNotPrepareFilesToStagewhenFlinkMasterIsSetToCollection[UseDataStreamForBatch = true]

Check notice on line 0 in .github

See this annotation in the file changed.

@github-actions github-actions / Test Results

374 tests found

There are 374 tests, see "Raw output" for the full list of tests.
Raw output
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldAcceptExplicitlySetIdleSourcesFlagWithCheckpointing[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldAcceptExplicitlySetIdleSourcesFlagWithCheckpointing[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldAcceptExplicitlySetIdleSourcesFlagWithoutCheckpointing[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldAcceptExplicitlySetIdleSourcesFlagWithoutCheckpointing[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldAllowPortOmissionForRemoteEnvironmentBatch[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldAllowPortOmissionForRemoteEnvironmentBatch[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldAllowPortOmissionForRemoteEnvironmentStreaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldAllowPortOmissionForRemoteEnvironmentStreaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldAutoSetIdleSourcesFlagWithCheckpointing[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldAutoSetIdleSourcesFlagWithCheckpointing[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldAutoSetIdleSourcesFlagWithoutCheckpointing[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldAutoSetIdleSourcesFlagWithoutCheckpointing[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldCreateFileSystemStateBackend[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldCreateFileSystemStateBackend[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldCreateRocksDbStateBackend[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldCreateRocksDbStateBackend[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldDetectMalformedPortBatch[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldDetectMalformedPortBatch[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldDetectMalformedPortStreaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldDetectMalformedPortStreaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldFailOnNoStoragePathProvided[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldFailOnNoStoragePathProvided[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldFailOnUnknownStateBackend[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldFailOnUnknownStateBackend[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldFallbackToDefaultParallelismBatch[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldFallbackToDefaultParallelismBatch[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldFallbackToDefaultParallelismStreaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldFallbackToDefaultParallelismStreaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldInferParallelismFromEnvironmentBatch[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldInferParallelismFromEnvironmentBatch[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldInferParallelismFromEnvironmentStreaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldInferParallelismFromEnvironmentStreaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldParsePortForRemoteEnvironmentBatch[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldParsePortForRemoteEnvironmentBatch[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldParsePortForRemoteEnvironmentStreaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldParsePortForRemoteEnvironmentStreaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldRemoveHttpProtocolFromHostBatch[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldRemoveHttpProtocolFromHostBatch[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldRemoveHttpProtocolFromHostStreaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldRemoveHttpProtocolFromHostStreaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSetMaxParallelismStreaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSetMaxParallelismStreaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSetParallelismBatch[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSetParallelismBatch[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSetParallelismStreaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSetParallelismStreaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSetSavepointRestoreForRemoteStreaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSetSavepointRestoreForRemoteStreaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSetWebUIOptions[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSetWebUIOptions[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSupportIPv4Batch[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSupportIPv4Batch[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSupportIPv4Streaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSupportIPv4Streaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSupportIPv6Batch[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSupportIPv6Batch[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSupportIPv6Streaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldSupportIPv6Streaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldTreatAutoAndEmptyHostTheSameBatch[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldTreatAutoAndEmptyHostTheSameBatch[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldTreatAutoAndEmptyHostTheSameStreaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ shouldTreatAutoAndEmptyHostTheSameStreaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ useDefaultParallelismFromContextBatch[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ useDefaultParallelismFromContextBatch[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ useDefaultParallelismFromContextStreaming[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkExecutionEnvironmentsTest ‑ useDefaultParallelismFromContextStreaming[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkJobServerDriverTest ‑ testConfigurationDefaults
org.apache.beam.runners.flink.FlinkJobServerDriverTest ‑ testConfigurationFromArgs
org.apache.beam.runners.flink.FlinkJobServerDriverTest ‑ testConfigurationFromConfig
org.apache.beam.runners.flink.FlinkJobServerDriverTest ‑ testJobServerDriver
org.apache.beam.runners.flink.FlinkJobServerDriverTest ‑ testJobServerDriverWithoutExpansionService
org.apache.beam.runners.flink.FlinkJobServerDriverTest ‑ testLegacyMasterUrlParameter
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldFailWhenFileDoesNotExistAndFlinkMasterIsSetExplicitly[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldFailWhenFileDoesNotExistAndFlinkMasterIsSetExplicitly[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldLogWarningWhenCheckpointingIsDisabled[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldLogWarningWhenCheckpointingIsDisabled[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldNotPrepareFilesToStageWhenFlinkMasterIsSetToAuto[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldNotPrepareFilesToStageWhenFlinkMasterIsSetToAuto[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldNotPrepareFilesToStageWhenFlinkMasterIsSetToLocal[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldNotPrepareFilesToStageWhenFlinkMasterIsSetToLocal[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldNotPrepareFilesToStagewhenFlinkMasterIsSetToCollection[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldNotPrepareFilesToStagewhenFlinkMasterIsSetToCollection[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldPrepareFilesToStageWhenFlinkMasterIsSetExplicitly[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldPrepareFilesToStageWhenFlinkMasterIsSetExplicitly[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldProvideParallelismToTransformOverrides[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldProvideParallelismToTransformOverrides[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldRecognizeAndTranslateStreamingPipeline[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldRecognizeAndTranslateStreamingPipeline[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldUseDefaultTempLocationIfNoneSet[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldUseDefaultTempLocationIfNoneSet[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldUsePreparedFilesOnRemoteEnvironment[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldUsePreparedFilesOnRemoteEnvironment[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldUseStreamingTransformOverridesWithUnboundedSources[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldUseStreamingTransformOverridesWithUnboundedSources[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldUseTransformOverrides[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ shouldUseTransformOverrides[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ testTranslationModeNoOverrideWithoutUnboundedSources[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ testTranslationModeNoOverrideWithoutUnboundedSources[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ testTranslationModeOverrideWithUnboundedSources[UseDataStreamForBatch = false]
org.apache.beam.runners.flink.FlinkPipelineExecutionEnvironmentTest ‑ testTranslationModeOverrideWithUnboundedSources[UseDataStreamForBatch = true]
org.apache.beam.runners.flink.FlinkPipelineOptionsTest ‑ parDoBaseClassPipelineOptionsNullTest
org.apache.beam.runners.flink.FlinkPipelineOptionsTest ‑ parDoBaseClassPipelineOptionsSerializationTest
org.apache.beam.runners.flink.FlinkPipelineOptionsTest ‑ testDefaults
org.apache.beam.runners.flink.FlinkRequiresStableInputTest ‑ testParDoRequiresStableInput
org.apache.beam.runners.flink.FlinkRequiresStableInputTest ‑ testParDoRequiresStableInputPortable
org.apache.beam.runners.flink.FlinkRequiresStableInputTest ‑ testParDoRequiresStableInputStateful
org.apache.beam.runners.flink.FlinkRequiresStableInputTest ‑ testParDoRequiresStableInputStatefulPortable
org.apache.beam.runners.flink.FlinkRunnerRegistrarTest ‑ testClassName
org.apache.beam.runners.flink.FlinkRunnerRegistrarTest ‑ testFullName
org.apache.beam.runners.flink.FlinkRunnerResultTest ‑ testCancelDoesNotThrowAnException
org.apache.beam.runners.flink.FlinkRunnerResultTest ‑ testPipelineResultReturnsDone
org.apache.beam.runners.flink.FlinkRunnerResultTest ‑ testWaitUntilFinishReturnsDone
org.apache.beam.runners.flink.FlinkRunnerTest ‑ testEnsureStdoutStdErrIsRestored
org.apache.beam.runners.flink.FlinkSavepointTest ‑ testSavepointRestoreLegacy
org.apache.beam.runners.flink.FlinkSavepointTest ‑ testSavepointRestorePortable
org.apache.beam.runners.flink.FlinkStreamingPipelineTranslatorTest ‑ testAutoBalanceShardKeyCacheIsNotSerialized
org.apache.beam.runners.flink.FlinkStreamingPipelineTranslatorTest ‑ testAutoBalanceShardKeyCacheIsStable
org.apache.beam.runners.flink.FlinkStreamingPipelineTranslatorTest ‑ testAutoBalanceShardKeyCacheMaxSize
org.apache.beam.runners.flink.FlinkStreamingPipelineTranslatorTest ‑ testAutoBalanceShardKeyResolvesMaxParallelism
org.apache.beam.runners.flink.FlinkStreamingPipelineTranslatorTest ‑ testStatefulParDoAfterCombineChaining
org.apache.beam.runners.flink.FlinkStreamingPipelineTranslatorTest ‑ testStatefulParDoAfterGroupByKeyChaining
org.apache.beam.runners.flink.FlinkStreamingTransformTranslatorsTest ‑ readSourceTranslatorBoundedWithMaxParallelism
org.apache.beam.runners.flink.FlinkStreamingTransformTranslatorsTest ‑ readSourceTranslatorBoundedWithoutMaxParallelism
org.apache.beam.runners.flink.FlinkStreamingTransformTranslatorsTest ‑ readSourceTranslatorUnboundedWithMaxParallelism
org.apache.beam.runners.flink.FlinkStreamingTransformTranslatorsTest ‑ readSourceTranslatorUnboundedWithoutMaxParallelism
org.apache.beam.runners.flink.FlinkSubmissionTest ‑ testDetachedSubmissionBatch
org.apache.beam.runners.flink.FlinkSubmissionTest ‑ testDetachedSubmissionBatchUseDataStream
org.apache.beam.runners.flink.FlinkSubmissionTest ‑ testDetachedSubmissionStreaming
org.apache.beam.runners.flink.FlinkSubmissionTest ‑ testSubmissionBatch
org.apache.beam.runners.flink.FlinkSubmissionTest ‑ testSubmissionBatchUseDataStream
org.apache.beam.runners.flink.FlinkSubmissionTest ‑ testSubmissionStreaming
org.apache.beam.runners.flink.FlinkTransformOverridesTest ‑ testRunnerDeterminedSharding
org.apache.beam.runners.flink.PipelineTranslationModeOptimizerTest ‑ testBoundedCollectionProducingTransform
org.apache.beam.runners.flink.PipelineTranslationModeOptimizerTest ‑ testUnboundedCollectionProducingTransform
org.apache.beam.runners.flink.PortableExecutionTest ‑ testExecution[streaming: false]
org.apache.beam.runners.flink.PortableExecutionTest ‑ testExecution[streaming: true]
org.apache.beam.runners.flink.PortableStateExecutionTest ‑ testExecution[streaming: false]
org.apache.beam.runners.flink.PortableStateExecutionTest ‑ testExecution[streaming: true]
org.apache.beam.runners.flink.PortableTimersExecutionTest ‑ testTimerExecution[streaming: false]
org.apache.beam.runners.flink.PortableTimersExecutionTest ‑ testTimerExecution[streaming: true]
org.apache.beam.runners.flink.ReadSourcePortableTest ‑ testExecution[streaming: false]
org.apache.beam.runners.flink.ReadSourcePortableTest ‑ testExecution[streaming: true]
org.apache.beam.runners.flink.ReadSourceStreamingTest ‑ testBatch
org.apache.beam.runners.flink.ReadSourceStreamingTest ‑ testStreaming
org.apache.beam.runners.flink.ReadSourceTest ‑ testJobCollectionExecution
org.apache.beam.runners.flink.ReadSourceTest ‑ testJobWithObjectReuse
org.apache.beam.runners.flink.ReadSourceTest ‑ testJobWithoutObjectReuse
org.apache.beam.runners.flink.adapter.BeamFlinkDataSetAdapterTest ‑ testApplyCompositeTransform
org.apache.beam.runners.flink.adapter.BeamFlinkDataSetAdapterTest ‑ testApplyGroupingTransform
org.apache.beam.runners.flink.adapter.BeamFlinkDataSetAdapterTest ‑ testApplyMultiInputTransform
org.apache.beam.runners.flink.adapter.BeamFlinkDataSetAdapterTest ‑ testApplyMultiOutputTransform
org.apache.beam.runners.flink.adapter.BeamFlinkDataSetAdapterTest ‑ testApplySimpleTransform
org.apache.beam.runners.flink.adapter.BeamFlinkDataSetAdapterTest ‑ testCustomCoder
org.apache.beam.runners.flink.adapter.BeamFlinkDataStreamAdapterTest ‑ testApplyCompositeTransform
org.apache.beam.runners.flink.adapter.BeamFlinkDataStreamAdapterTest ‑ testApplyGroupingTransform
org.apache.beam.runners.flink.adapter.BeamFlinkDataStreamAdapterTest ‑ testApplyMultiInputTransform
org.apache.beam.runners.flink.adapter.BeamFlinkDataStreamAdapterTest ‑ testApplyMultiOutputTransform
org.apache.beam.runners.flink.adapter.BeamFlinkDataStreamAdapterTest ‑ testApplyPreservesInputTimestamps
org.apache.beam.runners.flink.adapter.BeamFlinkDataStreamAdapterTest ‑ testApplyPreservesOutputTimestamps
org.apache.beam.runners.flink.adapter.BeamFlinkDataStreamAdapterTest ‑ testApplySimpleTransform
org.apache.beam.runners.flink.batch.NonMergingGroupByKeyTest ‑ testDisabledReIterationThrowsAnException
org.apache.beam.runners.flink.batch.NonMergingGroupByKeyTest ‑ testEnabledReIterationDoesNotThrowAnException
org.apache.beam.runners.flink.batch.ReshuffleTest ‑ testEqualDistributionOnReshuffleAcrossMultipleStages
org.apache.beam.runners.flink.metrics.FlinkMetricContainerTest ‑ testCounter
org.apache.beam.runners.flink.metrics.FlinkMetricContainerTest ‑ testDistribution
org.apache.beam.runners.flink.metrics.FlinkMetricContainerTest ‑ testDropUnexpectedMonitoringInfoTypes
org.apache.beam.runners.flink.metrics.FlinkMetricContainerTest ‑ testGauge
org.apache.beam.runners.flink.metrics.FlinkMetricContainerTest ‑ testMetricNameGeneration
org.apache.beam.runners.flink.metrics.FlinkMetricContainerTest ‑ testMonitoringInfoUpdate
org.apache.beam.runners.flink.streaming.BoundedSourceRestoreTest ‑ testRestore[0]
org.apache.beam.runners.flink.streaming.BoundedSourceRestoreTest ‑ testRestore[1]
org.apache.beam.runners.flink.streaming.BoundedSourceRestoreTest ‑ testRestore[2]
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testBag
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testBagIsEmpty
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testBagWithBadCoderEquality
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testCombiningIsEmpty
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testCombiningValue
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testMap
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testMapReadable
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testMergeBagIntoNewNamespace
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testMergeBagIntoSource
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testMergeCombiningValueIntoNewNamespace
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testMergeCombiningValueIntoSource
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testMergeCombiningWithContextValueIntoNewNamespace
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testMergeCombiningWithContextValueIntoSource
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testMergeSetIntoNewNamespace
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testMergeSetIntoSource
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testSet
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testSetIsEmpty
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testSetReadable
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testValue
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testWatermarkEarliestState
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testWatermarkEndOfWindowState
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testWatermarkLatestState
org.apache.beam.runners.flink.streaming.FlinkBroadcastStateInternalsTest ‑ testWatermarkStateIsEmpty
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testBag
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testBagIsEmpty
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testBagWithBadCoderEquality
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testCombiningIsEmpty
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testCombiningValue
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testGlobalWindowWatermarkHoldClear
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testMap
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testMapReadable
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testMergeBagIntoNewNamespace
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testMergeBagIntoSource
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testMergeCombiningValueIntoNewNamespace
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testMergeCombiningValueIntoSource
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testMergeCombiningWithContextValueIntoNewNamespace
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testMergeCombiningWithContextValueIntoSource
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testMergeSetIntoNewNamespace
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testMergeSetIntoSource
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testSet
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testSetIsEmpty
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testSetReadable
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testValue
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testWatermarkEarliestState
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testWatermarkEndOfWindowState
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testWatermarkHoldsPersistence
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testWatermarkLatestState
org.apache.beam.runners.flink.streaming.FlinkStateInternalsTest ‑ testWatermarkStateIsEmpty
org.apache.beam.runners.flink.streaming.GroupByNullKeyTest ‑ testProgram
org.apache.beam.runners.flink.streaming.GroupByWithNullValuesTest ‑ testGroupByWithNullValues
org.apache.beam.runners.flink.streaming.TopWikipediaSessionsTest ‑ testProgram
org.apache.beam.runners.flink.translation.functions.FlinkDoFnFunctionTest ‑ testAccumulatorRegistrationOnOperatorClose
org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunctionTest ‑ expectedInputsAreSent[0]
org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunctionTest ‑ expectedInputsAreSent[1]
org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunctionTest ‑ outputsAreTaggedCorrectly[0]
org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunctionTest ‑ outputsAreTaggedCorrectly[1]
org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunctionTest ‑ sdkErrorsSurfaceOnClose[0]
org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunctionTest ‑ sdkErrorsSurfaceOnClose[1]
org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunctionTest ‑ testAccumulatorRegistrationOnOperatorClose[0]
org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunctionTest ‑ testAccumulatorRegistrationOnOperatorClose[1]
org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunctionTest ‑ testStageBundleClosed[0]
org.apache.beam.runners.flink.translation.functions.FlinkExecutableStageFunctionTest ‑ testStageBundleClosed[1]
org.apache.beam.runners.flink.translation.functions.FlinkStatefulDoFnFunctionTest ‑ testAccumulatorRegistrationOnOperatorClose
org.apache.beam.runners.flink.translation.functions.ImpulseSourceFunctionTest ‑ testImpulseInitial
org.apache.beam.runners.flink.translation.functions.ImpulseSourceFunctionTest ‑ testImpulseRestored
org.apache.beam.runners.flink.translation.functions.ImpulseSourceFunctionTest ‑ testInstanceOfSourceFunction
org.apache.beam.runners.flink.translation.functions.ImpulseSourceFunctionTest ‑ testKeepAlive
org.apache.beam.runners.flink.translation.functions.ImpulseSourceFunctionTest ‑ testKeepAliveDuringInterrupt
org.apache.beam.runners.flink.translation.types.CoderTypeSerializerTest ‑ shouldWriteAndReadSnapshotForAnonymousClassCoder
org.apache.beam.runners.flink.translation.types.CoderTypeSerializerTest ‑ shouldWriteAndReadSnapshotForConcreteClassCoder
org.apache.beam.runners.flink.translation.types.UnversionedTypeSerializerSnapshotTest ‑ testSerialization
org.apache.beam.runners.flink.translation.wrappers.SourceInputFormatTest ‑ testAccumulatorRegistrationOnOperatorClose
org.apache.beam.runners.flink.translation.wrappers.streaming.DedupingOperatorTest ‑ testDeduping
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ keyedParDoPushbackDataCheckpointing
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ keyedParDoSideInputCheckpointing
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ nonKeyedParDoPushbackDataCheckpointing
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ nonKeyedParDoSideInputCheckpointing
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testAccumulatorRegistrationOnOperatorClose
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testBundle
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testBundleKeyed
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testBundleProcessingExceptionIsFatalDuringCheckpointing
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testCheckpointBufferingWithMultipleBundles
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testExactlyOnceBuffering
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testExactlyOnceBufferingFlushDuringDrain
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testExactlyOnceBufferingKeyed
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testFailOnRequiresStableInputAndDisabledCheckpointing
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testGCForGlobalWindow
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testKeyedParDoSideInputs
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testLateDroppingForStatefulFn
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testMultiOutputOutput
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testNormalParDoSideInputs
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testRemoveCachedClassReferences
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testSingleOutput
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testStateGCForStatefulFn
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testStateRestore
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testTimersRestore
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testWatermarkContract
org.apache.beam.runners.flink.translation.wrappers.streaming.DoFnOperatorTest ‑ testWatermarkUpdateAfterWatermarkHoldRelease
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ expectedInputsAreSent
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ outputsAreTaggedCorrectly
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ sdkErrorsSurfaceOnClose
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ testCacheTokenHandling
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ testEnsureDeferredStateCleanupTimerFiring
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ testEnsureDeferredStateCleanupTimerFiringWithCheckpointing
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ testEnsureStateCleanupOnFinalWatermark
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ testEnsureStateCleanupWithKeyedInput
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ testEnsureStateCleanupWithKeyedInputCleanupTimer
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ testEnsureStateCleanupWithKeyedInputStateCleaner
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ testSerialization
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ testStableInputApplied
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ testStageBundleClosed
org.apache.beam.runners.flink.translation.wrappers.streaming.ExecutableStageDoFnOperatorTest ‑ testWatermarkHandling
org.apache.beam.runners.flink.translation.wrappers.streaming.FlinkKeyUtilsTest ‑ testCoderContext
org.apache.beam.runners.flink.translation.wrappers.streaming.FlinkKeyUtilsTest ‑ testEncodeDecode
org.apache.beam.runners.flink.translation.wrappers.streaming.FlinkKeyUtilsTest ‑ testFromEncodedKey
org.apache.beam.runners.flink.translation.wrappers.streaming.FlinkKeyUtilsTest ‑ testNullKey
org.apache.beam.runners.flink.translation.wrappers.streaming.WindowDoFnOperatorTest ‑ testRestore
org.apache.beam.runners.flink.translation.wrappers.streaming.WindowDoFnOperatorTest ‑ testTimerCleanupOfPendingTimerList
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$BasicTest ‑ testAccumulatorRegistrationOnOperatorClose
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$BasicTest ‑ testSequentialReadingFromBoundedSource
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$BasicTest ‑ testSerialization
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$BasicTest ‑ testSourceWithNoReaderDoesNotShutdown
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$BasicTest ‑ testSourceWithReadersDoesNotShutdown
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$IntegrationTests ‑ testPollingOfIdleReaders
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testNullCheckpoint[numTasks = 1; numSplits=1]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testNullCheckpoint[numTasks = 1; numSplits=2]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testNullCheckpoint[numTasks = 1; numSplits=4]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testNullCheckpoint[numTasks = 2; numSplits=1]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testNullCheckpoint[numTasks = 2; numSplits=2]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testNullCheckpoint[numTasks = 2; numSplits=4]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testNullCheckpoint[numTasks = 4; numSplits=1]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testNullCheckpoint[numTasks = 4; numSplits=2]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testNullCheckpoint[numTasks = 4; numSplits=4]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testRestore[numTasks = 1; numSplits=1]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testRestore[numTasks = 1; numSplits=2]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testRestore[numTasks = 1; numSplits=4]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testRestore[numTasks = 2; numSplits=1]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testRestore[numTasks = 2; numSplits=2]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testRestore[numTasks = 2; numSplits=4]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testRestore[numTasks = 4; numSplits=1]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testRestore[numTasks = 4; numSplits=2]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testRestore[numTasks = 4; numSplits=4]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testValueEmission[numTasks = 1; numSplits=1]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testValueEmission[numTasks = 1; numSplits=2]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testValueEmission[numTasks = 1; numSplits=4]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testValueEmission[numTasks = 2; numSplits=1]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testValueEmission[numTasks = 2; numSplits=2]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testValueEmission[numTasks = 2; numSplits=4]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testValueEmission[numTasks = 4; numSplits=1]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testValueEmission[numTasks = 4; numSplits=2]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testValueEmission[numTasks = 4; numSplits=4]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testWatermarkEmission[numTasks = 1; numSplits=1]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testWatermarkEmission[numTasks = 1; numSplits=2]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testWatermarkEmission[numTasks = 1; numSplits=4]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testWatermarkEmission[numTasks = 2; numSplits=1]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testWatermarkEmission[numTasks = 2; numSplits=2]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testWatermarkEmission[numTasks = 2; numSplits=4]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testWatermarkEmission[numTasks = 4; numSplits=1]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testWatermarkEmission[numTasks = 4; numSplits=2]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.UnboundedSourceWrapperTest$ParameterizedUnboundedSourceWrapperTest ‑ testWatermarkEmission[numTasks = 4; numSplits=4]
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.FlinkSourceSplitEnumeratorTest ‑ testAddSplitsBack
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.FlinkSourceSplitEnumeratorTest ‑ testAssignSplitsWithBoundedSource
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.FlinkSourceSplitEnumeratorTest ‑ testAssignSplitsWithUnboundedSource
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testExceptionInExecutorThread
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testIsAvailableOnNoMoreSplitsNotification
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testIsAvailableOnSplitsAssignment
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testIsAvailableWithIdleTimeout
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testIsAvailableWithoutIdleTimeout
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testMetricsContainer
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testNumBytesInMetrics
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testPollBasic
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testPollEmitsMaxWatermark
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testPollFromEmptySplit
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testPollWithIdleTimeout
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testPollWithTimestampExtractor
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testPollWithoutIdleTimeout
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.bounded.FlinkBoundedSourceReaderTest ‑ testSnapshotStateAndRestore
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testCheckMarksFinalized
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testExceptionInExecutorThread
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testIsAvailableAlwaysWakenUp
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testIsAvailableOnNoMoreSplitsNotification
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testIsAvailableOnSplitChangeWhenNoDataAvailableForAliveReaders
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testIsAvailableWithIdleTimeout
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testIsAvailableWithoutIdleTimeout
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testMetricsContainer
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testNumBytesInMetrics
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testPendingBytesMetric
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testPollBasic
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testPollFromEmptySplit
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testPollWithTimestampExtractor
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testSnapshotStateAndRestore
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testWatermark
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testWatermarkOnEmptySource
org.apache.beam.runners.flink.translation.wrappers.streaming.io.source.unbounded.FlinkUnboundedSourceReaderTest ‑ testWatermarkOnNoSplits
org.apache.beam.runners.flink.translation.wrappers.streaming.stableinput.BufferedElementsTest ‑ testCoder
org.apache.beam.runners.flink.translation.wrappers.streaming.stableinput.BufferingDoFnRunnerTest ‑ testRejectConcurrentCheckpointingBoundaries
org.apache.beam.runners.flink.translation.wrappers.streaming.stableinput.BufferingDoFnRunnerTest ‑ testRestoreWithConcurrentCheckpoints
org.apache.beam.runners.flink.translation.wrappers.streaming.stableinput.BufferingDoFnRunnerTest ‑ testRestoreWithConcurrentCheckpointsFromPendingCheckpoint
org.apache.beam.runners.flink.translation.wrappers.streaming.stableinput.BufferingDoFnRunnerTest ‑ testRestoreWithConcurrentCheckpointsFromPendingCheckpoints
org.apache.beam.runners.flink.translation.wrappers.streaming.stableinput.BufferingDoFnRunnerTest ‑ testRestoreWithoutConcurrentCheckpoints
org.apache.beam.runners.flink.translation.wrappers.streaming.stableinput.BufferingDoFnRunnerTest ‑ testRestoreWithoutConcurrentCheckpointsWithPendingCheckpoint
org.apache.beam.runners.flink.translation.wrappers.streaming.stableinput.BufferingDoFnRunnerTest ‑ testRestoreWithoutConcurrentCheckpointsWithPendingCheckpointFromConcurrentCheckpointing