Skip to content

Commit 13549dc

Browse files
committed
Copy experiments list to mutable ArrayList before addExperiment() calls
1 parent f3e822c commit 13549dc

3 files changed

Lines changed: 20 additions & 6 deletions

File tree

runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineTranslator.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -403,6 +403,10 @@ public Translator(Pipeline pipeline, DataflowRunner runner, SdkComponents sdkCom
403403
* @return a Job definition filled in with the type of job, the environment, and the job steps.
404404
*/
405405
public Job translate(List<DataflowPackage> packages) {
406+
// Ensure the experiments list is mutable before any experiments are added.
407+
if (options.getExperiments() != null) {
408+
options.setExperiments(new ArrayList<>(options.getExperiments()));
409+
}
406410
job.setName(options.getJobName().toLowerCase());
407411

408412
Environment environment = new Environment();

runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java

Lines changed: 12 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1243,15 +1243,21 @@ private static boolean includesTransformUpgrades(Pipeline pipeline) {
12431243
@SuppressWarnings("Slf4jFormatShouldBeConst")
12441244
@Override
12451245
public DataflowPipelineJob run(Pipeline pipeline) {
1246+
// Ensure the experiments list is mutable before any experiments are added.
1247+
if (options.getExperiments() != null) {
1248+
options.setExperiments(new ArrayList<>(options.getExperiments()));
1249+
}
12461250
// Multi-language pipelines and pipelines that include upgrades should automatically be upgraded
12471251
// to Dataflow Portable Runner.
12481252
if (DataflowRunner.isMultiLanguagePipeline(pipeline) || includesTransformUpgrades(pipeline)) {
1249-
if (!firstNonNull(options.getExperiments(), Collections.emptyList())
1250-
.contains("use_runner_v2")) {
1251-
LOG.info(
1252-
"Automatically enabling Dataflow Portable Runner since the pipeline used cross-language"
1253-
+ " transforms or pipeline needed a transform upgrade.");
1254-
ExperimentalOptions.addExperiment(options, "use_runner_v2");
1253+
if (!useUnifiedWorker(options)) {
1254+
if (!firstNonNull(options.getExperiments(), Collections.emptyList())
1255+
.contains("use_runner_v2")) {
1256+
LOG.info(
1257+
"Automatically enabling Dataflow Portable Runner since the pipeline used cross-language"
1258+
+ " transforms or pipeline needed a transform upgrade.");
1259+
ExperimentalOptions.addExperiment(options, "use_runner_v2");
1260+
}
12551261
}
12561262
}
12571263
if (useUnifiedWorker(options)) {

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcWindmillServer.java

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -111,6 +111,10 @@ private static DataflowWorkerHarnessOptions testOptions(
111111
boolean enableStreamingEngine, List<String> additionalExperiments) {
112112
DataflowWorkerHarnessOptions options =
113113
PipelineOptionsFactory.create().as(DataflowWorkerHarnessOptions.class);
114+
// Ensure the experiments list is mutable before any experiments are added.
115+
if (options.getExperiments() != null) {
116+
options.setExperiments(new ArrayList<>(options.getExperiments()));
117+
}
114118
options.setProject("project");
115119
options.setJobId("job");
116120
options.setWorkerId("worker");

0 commit comments

Comments
 (0)