Skip to content

Commit 2240a24

Browse files
committed
Set use_staged_dataflow_worker_jar for dev Beam
1 parent 4ae92f5 commit 2240a24

1 file changed

Lines changed: 9 additions & 0 deletions

File tree

it/google-cloud-platform/src/main/java/org/apache/beam/it/gcp/dataflow/DefaultPipelineLauncher.java

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -49,9 +49,11 @@
4949
import org.apache.beam.sdk.metrics.MetricQueryResults;
5050
import org.apache.beam.sdk.metrics.MetricResult;
5151
import org.apache.beam.sdk.metrics.MetricsFilter;
52+
import org.apache.beam.sdk.options.ExperimentalOptions;
5253
import org.apache.beam.sdk.options.PipelineOptions;
5354
import org.apache.beam.sdk.options.PipelineOptionsFactory;
5455
import org.apache.beam.sdk.options.StreamingOptions;
56+
import org.apache.beam.sdk.util.ReleaseInfo;
5557
import org.apache.beam.sdk.util.common.ReflectHelpers;
5658
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Preconditions;
5759
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.base.Strings;
@@ -312,6 +314,13 @@ public LaunchInfo launch(String project, String region, LaunchConfig options) th
312314
optionFromConfig.add(
313315
String.format("--tempLocation=%s", pipelineOptions.getTempLocation()));
314316
}
317+
// for snapshot Beam running on runner v2, need to use staged SDK harness
318+
if (ReleaseInfo.getReleaseInfo().isDevSdkVersion()
319+
&& (ExperimentalOptions.hasExperiment(pipelineOptions, "use_runner_v2")
320+
|| (!Strings.isNullOrEmpty(options.parameters().get("experiments"))
321+
&& options.parameters().get("experiments").contains("use_runner_v2")))) {
322+
optionFromConfig.add("--experiments=use_staged_dataflow_worker_jar");
323+
}
315324

316325
// dataflow runner specific options
317326
PipelineOptions updatedOptions =

0 commit comments

Comments
 (0)