Skip to content

Commit c01230e

Browse files
committed
Aded portable runner to python and go runners
1 parent efe4e94 commit c01230e

4 files changed

Lines changed: 52 additions & 5 deletions

File tree

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

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -335,7 +335,7 @@ func getJobOptions(ctx context.Context, streaming bool) (*dataflowlib.JobOptions
335335
experiments := jobopts.GetExperiments()
336336
// Ensure that we enable the same set of experiments across all SDKs
337337
// for runner v2.
338-
var fnApiSet, v2set, uwSet, portaSubmission, seSet, wsSet bool
338+
var fnApiSet, v2set, uwSet, portableRunnerSet, portaSubmission, seSet, wsSet bool
339339
for _, e := range experiments {
340340
if strings.Contains(e, "beam_fn_api") {
341341
fnApiSet = true
@@ -349,7 +349,10 @@ func getJobOptions(ctx context.Context, streaming bool) (*dataflowlib.JobOptions
349349
if strings.Contains(e, "use_portable_job_submission") {
350350
portaSubmission = true
351351
}
352-
if strings.Contains(e, "disable_runner_v2") || strings.Contains(e, "disable_runner_v2_until_2023") || strings.Contains(e, "disable_prime_runner_v2") {
352+
if strings.Contains(e, "enable_portable_runner") {
353+
portableRunnerSet = true
354+
}
355+
if strings.Contains(e, "disable_runner_v2") || strings.Contains(e, "disable_runner_v2_until_2023") || strings.Contains(e, "disable_prime_runner_v2") || strings.Contains(e, "disable_portable_runner") || strings.Contains(e, "enable_streaming_java_runner") {
353356
return nil, errors.New("detected one of the following experiments: disable_runner_v2 | disable_runner_v2_until_2023 | disable_prime_runner_v2. Disabling runner v2 is no longer supported as of Beam version 2.45.0+")
354357
}
355358
}
@@ -366,6 +369,9 @@ func getJobOptions(ctx context.Context, streaming bool) (*dataflowlib.JobOptions
366369
if !portaSubmission {
367370
experiments = append(experiments, "use_portable_job_submission")
368371
}
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+
}
369375

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

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

Lines changed: 36 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -191,12 +191,12 @@ func TestGetJobOptions_NoExperimentsSet(t *testing.T) {
191191
if err != nil {
192192
t.Fatalf("getJobOptions() returned error %q, want %q", err, "nil")
193193
}
194-
if got, want := len(opts.Experiments), 4; got != want {
194+
if got, want := len(opts.Experiments), 5; got != want {
195195
t.Fatalf("len(getJobOptions().Experiments) = %d, want %d", got, want)
196196
}
197197
sort.Strings(opts.Experiments)
198198
expectedExperiments := []string{"beam_fn_api", "use_portable_job_submission", "use_unified_worker", "use_runner_v2"}
199-
for i := 0; i < 2; i++ {
199+
for i := 0; i < 5; i++ {
200200
if got, want := opts.Experiments[i], expectedExperiments[i]; got != want {
201201
t.Errorf("getJobOptions().Experiments[%d] = %q, want %q", i, got, want)
202202
}
@@ -244,6 +244,40 @@ func TestGetJobOptions_DisableRunnerV2ExperimentsSet(t *testing.T) {
244244
}
245245
}
246246

247+
func TestGetJobOptions_DisablePortableRunnerExperimentsSet(t *testing.T) {
248+
resetGlobals()
249+
*stagingLocation = "gs://testStagingLocation"
250+
*gcpopts.Project = "testProject"
251+
*gcpopts.Region = "testRegion"
252+
*jobopts.Experiments = "disable_portable_runner"
253+
254+
opts, err := getJobOptions(context.Background(), false)
255+
256+
if err == nil {
257+
t.Error("getJobOptions() returned error nil, want an error")
258+
}
259+
if opts != nil {
260+
t.Errorf("getJobOptions() returned JobOptions when it should not have, got %#v, want nil", opts)
261+
}
262+
}
263+
264+
func TestGetJobOptions_EnableStreamingJavaRunnerExperimentsSet(t *testing.T) {
265+
resetGlobals()
266+
*stagingLocation = "gs://testStagingLocation"
267+
*gcpopts.Project = "testProject"
268+
*gcpopts.Region = "testRegion"
269+
*jobopts.Experiments = "enable_streaming_java_runner"
270+
271+
opts, err := getJobOptions(context.Background(), false)
272+
273+
if err == nil {
274+
t.Error("getJobOptions() returned error nil, want an error")
275+
}
276+
if opts != nil {
277+
t.Errorf("getJobOptions() returned JobOptions when it should not have, got %#v, want nil", opts)
278+
}
279+
}
280+
247281
func TestGetJobOptions_NoStagingLocation(t *testing.T) {
248282
resetGlobals()
249283
*stagingLocation = ""

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)