Skip to content

Commit 86063e1

Browse files
committed
Copy experiments list to mutable ArrayList before addExperiment() calls
1 parent f46213e commit 86063e1

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
@@ -1241,15 +1241,21 @@ private static boolean includesTransformUpgrades(Pipeline pipeline) {
12411241
@SuppressWarnings("Slf4jFormatShouldBeConst")
12421242
@Override
12431243
public DataflowPipelineJob run(Pipeline pipeline) {
1244+
// Ensure the experiments list is mutable before any experiments are added.
1245+
if (options.getExperiments() != null) {
1246+
options.setExperiments(new ArrayList<>(options.getExperiments()));
1247+
}
12441248
// Multi-language pipelines and pipelines that include upgrades should automatically be upgraded
12451249
// to Dataflow Portable Runner.
12461250
if (DataflowRunner.isMultiLanguagePipeline(pipeline) || includesTransformUpgrades(pipeline)) {
1247-
if (!firstNonNull(options.getExperiments(), Collections.emptyList())
1248-
.contains("use_runner_v2")) {
1249-
LOG.info(
1250-
"Automatically enabling Dataflow Portable Runner since the pipeline used cross-language"
1251-
+ " transforms or pipeline needed a transform upgrade.");
1252-
ExperimentalOptions.addExperiment(options, "use_runner_v2");
1251+
if (!useUnifiedWorker(options)) {
1252+
if (!firstNonNull(options.getExperiments(), Collections.emptyList())
1253+
.contains("use_runner_v2")) {
1254+
LOG.info(
1255+
"Automatically enabling Dataflow Portable Runner since the pipeline used cross-language"
1256+
+ " transforms or pipeline needed a transform upgrade.");
1257+
ExperimentalOptions.addExperiment(options, "use_runner_v2");
1258+
}
12531259
}
12541260
}
12551261
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)