Skip to content

Commit 6e89f04

Browse files
committed
Added Portable Runner Passing for Python Runners
1 parent efe4e94 commit 6e89f04

2 files changed

Lines changed: 8 additions & 1 deletion

File tree

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -591,6 +591,7 @@ def _add_runner_v2_missing_options(options):
591591
debug_options.add_experiment('use_unified_worker')
592592
debug_options.add_experiment('use_runner_v2')
593593
debug_options.add_experiment('use_portable_job_submission')
594+
debug_options.add_experiment('enable_portable_runner')
594595

595596

596597
def _check_and_add_missing_options(options):
@@ -662,6 +663,8 @@ def _is_runner_v2_disabled(options):
662663
"""Returns true if runner v2 is disabled."""
663664
debug_options = options.view_as(DebugOptions)
664665
return (
666+
debug_options.lookup_experiment('disable_portable_runner') or
667+
debug_options.lookup_experiment('enable_streaming_java_runner') or
665668
debug_options.lookup_experiment('disable_runner_v2') or
666669
debug_options.lookup_experiment('disable_runner_v2_until_2023') or
667670
debug_options.lookup_experiment('disable_runner_v2_until_v2.50') or

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

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -624,7 +624,8 @@ def test_batch_is_runner_v2(self):
624624
for expected in ['beam_fn_api',
625625
'use_unified_worker',
626626
'use_runner_v2',
627-
'use_portable_job_submission']:
627+
'use_portable_job_submission',
628+
'enable_portable_runner']:
628629
self.assertTrue(
629630
options.view_as(DebugOptions).lookup_experiment(expected, False),
630631
expected)
@@ -636,6 +637,7 @@ def test_streaming_is_runner_v2(self):
636637
for expected in ['beam_fn_api',
637638
'use_unified_worker',
638639
'use_runner_v2',
640+
'enable_portable_runner',
639641
'use_portable_job_submission',
640642
'enable_windmill_service',
641643
'enable_streaming_engine']:
@@ -653,6 +655,7 @@ def test_dataflow_service_options_enable_prime_sets_runner_v2(self):
653655
for expected in ['beam_fn_api',
654656
'use_unified_worker',
655657
'use_runner_v2',
658+
'enable_portable_runner',
656659
'use_portable_job_submission']:
657660
self.assertTrue(
658661
options.view_as(DebugOptions).lookup_experiment(expected, False),
@@ -669,6 +672,7 @@ def test_dataflow_service_options_enable_prime_sets_runner_v2(self):
669672
'use_unified_worker',
670673
'use_runner_v2',
671674
'use_portable_job_submission',
675+
'enable_portable_runner',
672676
'enable_windmill_service',
673677
'enable_streaming_engine']:
674678
self.assertTrue(

0 commit comments

Comments
 (0)