Skip to content

Commit dbb9826

Browse files
committed
Removed empty if, added disable unit test for python.
1 parent c9445ef commit dbb9826

2 files changed

Lines changed: 24 additions & 3 deletions

File tree

sdks/go/pkg/beam/runners/dataflow/dataflow.go

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -369,9 +369,7 @@ func getJobOptions(ctx context.Context, streaming bool) (*dataflowlib.JobOptions
369369
if !portaSubmission {
370370
experiments = append(experiments, "use_portable_job_submission")
371371
}
372-
if !portableRunnerSet {
373-
// As this option is not documented, we do not set it by default. This behavior will be fixed in later versions.
374-
}
372+
// As portable_runner is not documented, we do not set it by default. This behavior will be fixed in later versions.
375373

376374
// Ensure that streaming specific experiments are set for streaming pipelines
377375
// since runner v2 only supports using streaming engine.

sdks/python/apache_beam/runners/dataflow/dataflow_runner_test.py

Lines changed: 23 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,6 +41,7 @@
4141
from apache_beam.runners.dataflow.dataflow_runner import DataflowRuntimeException
4242
from apache_beam.runners.dataflow.dataflow_runner import _check_and_add_missing_options
4343
from apache_beam.runners.dataflow.dataflow_runner import _check_and_add_missing_streaming_options
44+
from apache_beam.runners.dataflow.dataflow_runner import _is_runner_v2_disabled
4445
from apache_beam.runners.dataflow.internal.clients import dataflow as dataflow_api
4546
from apache_beam.runners.internal import names
4647
from apache_beam.runners.runner import PipelineState
@@ -734,5 +735,27 @@ def test_explicit_streaming_no_unbounded(self):
734735
apiclient.dataflow.Job.TypeValueValuesEnum.JOB_TYPE_STREAMING)
735736

736737

738+
class DataflowRunnerV2DisabledTest(unittest.TestCase):
739+
740+
def test_runner_v2_disabled_experiments_raise(self):
741+
disable_experiments = [
742+
'disable_portable_runner',
743+
'enable_streaming_java_runner',
744+
'disable_runner_v2',
745+
'disable_runner_v2_until_2023',
746+
'disable_runner_v2_until_v2.50',
747+
'disable_prime_runner_v2',
748+
]
749+
for experiment in disable_experiments:
750+
options = PipelineOptions([f'--experiments={experiment}'])
751+
self.assertTrue(
752+
_is_runner_v2_disabled(options),
753+
f'Expected {experiment} to disable runner v2')
754+
with self.assertRaisesRegex(
755+
ValueError,
756+
'Disabling Runner V2 no longer supported'):
757+
DataflowRunner().run_pipeline(None, options)
758+
759+
737760
if __name__ == '__main__':
738761
unittest.main()

0 commit comments

Comments
 (0)