Skip to content

Commit 6e70264

Browse files
committed
Use ExperimentalOptions.addExperiment() instead of get/set pattern in Java
Replaces getExperiments()/modify/setExperiments() boilerplate with ExperimentalOptions.addExperiment() which handles null-init and deduplication. Resolves #19347
1 parent 9ad7856 commit 6e70264

3 files changed

Lines changed: 20 additions & 62 deletions

File tree

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

Lines changed: 3 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -69,6 +69,7 @@
6969
import org.apache.beam.sdk.coders.IterableCoder;
7070
import org.apache.beam.sdk.coders.KvCoder;
7171
import org.apache.beam.sdk.extensions.gcp.options.GcpOptions;
72+
import org.apache.beam.sdk.options.ExperimentalOptions;
7273
import org.apache.beam.sdk.options.PipelineOptions;
7374
import org.apache.beam.sdk.options.StreamingOptions;
7475
import org.apache.beam.sdk.runners.AppliedPTransform;
@@ -414,19 +415,8 @@ public Job translate(List<DataflowPackage> packages) {
414415
// back end as well. If streaming engine is not enabled make sure the experiments are also
415416
// not enabled.
416417
if (options.isEnableStreamingEngine()) {
417-
List<String> experiments = options.getExperiments();
418-
if (experiments == null) {
419-
experiments = new ArrayList<String>();
420-
} else {
421-
experiments = new ArrayList<String>(experiments);
422-
}
423-
if (!experiments.contains(GcpOptions.STREAMING_ENGINE_EXPERIMENT)) {
424-
experiments.add(GcpOptions.STREAMING_ENGINE_EXPERIMENT);
425-
}
426-
if (!experiments.contains(GcpOptions.WINDMILL_SERVICE_EXPERIMENT)) {
427-
experiments.add(GcpOptions.WINDMILL_SERVICE_EXPERIMENT);
428-
}
429-
options.setExperiments(experiments);
418+
ExperimentalOptions.addExperiment(options, GcpOptions.STREAMING_ENGINE_EXPERIMENT);
419+
ExperimentalOptions.addExperiment(options, GcpOptions.WINDMILL_SERVICE_EXPERIMENT);
430420
} else {
431421
List<String> experiments = options.getExperiments();
432422
if (experiments != null) {

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

Lines changed: 11 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -1244,13 +1244,11 @@ public DataflowPipelineJob run(Pipeline pipeline) {
12441244
// Multi-language pipelines and pipelines that include upgrades should automatically be upgraded
12451245
// to Runner v2.
12461246
if (DataflowRunner.isMultiLanguagePipeline(pipeline) || includesTransformUpgrades(pipeline)) {
1247-
List<String> experiments = firstNonNull(options.getExperiments(), Collections.emptyList());
1248-
if (!experiments.contains("use_runner_v2")) {
1247+
if (!firstNonNull(options.getExperiments(), Collections.emptyList()).contains("use_runner_v2")) {
12491248
LOG.info(
12501249
"Automatically enabling Dataflow Runner v2 since the pipeline used cross-language"
12511250
+ " transforms or pipeline needed a transform upgrade.");
1252-
options.setExperiments(
1253-
ImmutableList.<String>builder().addAll(experiments).add("use_runner_v2").build());
1251+
ExperimentalOptions.addExperiment(options, "use_runner_v2");
12541252
}
12551253
}
12561254
if (useUnifiedWorker(options)) {
@@ -1260,21 +1258,10 @@ public DataflowPipelineJob run(Pipeline pipeline) {
12601258
throw new IllegalArgumentException(
12611259
"Runner V2 both disabled and enabled: at least one of ['beam_fn_api', 'use_unified_worker', 'use_runner_v2', 'use_portable_job_submission'] is set and also one of ['disable_runner_v2', 'disable_runner_v2_until_2023', 'disable_prime_runner_v2'] is set.");
12621260
}
1263-
List<String> experiments =
1264-
new ArrayList<>(options.getExperiments()); // non-null if useUnifiedWorker is true
1265-
if (!experiments.contains("use_runner_v2")) {
1266-
experiments.add("use_runner_v2");
1267-
}
1268-
if (!experiments.contains("use_unified_worker")) {
1269-
experiments.add("use_unified_worker");
1270-
}
1271-
if (!experiments.contains("beam_fn_api")) {
1272-
experiments.add("beam_fn_api");
1273-
}
1274-
if (!experiments.contains("use_portable_job_submission")) {
1275-
experiments.add("use_portable_job_submission");
1276-
}
1277-
options.setExperiments(ImmutableList.copyOf(experiments));
1261+
ExperimentalOptions.addExperiment(options, "use_runner_v2");
1262+
ExperimentalOptions.addExperiment(options, "use_unified_worker");
1263+
ExperimentalOptions.addExperiment(options, "beam_fn_api");
1264+
ExperimentalOptions.addExperiment(options, "use_portable_job_submission");
12781265
// Ensure that logging via the FnApi is enabled
12791266
options.as(SdkHarnessOptions.class).setEnableLogViaFnApi(true);
12801267
}
@@ -1301,14 +1288,9 @@ public DataflowPipelineJob run(Pipeline pipeline) {
13011288
options.setStreaming(true);
13021289

13031290
{
1304-
List<String> experiments =
1305-
options.getExperiments() == null
1306-
? new ArrayList<>()
1307-
: new ArrayList<>(options.getExperiments());
13081291
// Experiment marking that the harness supports tag encoding v2
13091292
// Backend will enable tag encoding v2 only if the harness supports it.
1310-
experiments.add("streaming_engine_state_tag_encoding_v2_supported");
1311-
options.setExperiments(ImmutableList.copyOf(experiments));
1293+
ExperimentalOptions.addExperiment(options, "streaming_engine_state_tag_encoding_v2_supported");
13121294
}
13131295

13141296
if (useUnifiedWorker(options)) {
@@ -1413,15 +1395,7 @@ public DataflowPipelineJob run(Pipeline pipeline) {
14131395
pipeline, dataflowV1PipelineProto, dataflowV1Components, this, packages);
14141396

14151397
if (!isNullOrEmpty(dataflowOptions.getDataflowWorkerJar()) && !useUnifiedWorker(options)) {
1416-
List<String> experiments =
1417-
firstNonNull(dataflowOptions.getExperiments(), Collections.emptyList());
1418-
if (!experiments.contains("use_staged_dataflow_worker_jar")) {
1419-
dataflowOptions.setExperiments(
1420-
ImmutableList.<String>builder()
1421-
.addAll(experiments)
1422-
.add("use_staged_dataflow_worker_jar")
1423-
.build());
1424-
}
1398+
ExperimentalOptions.addExperiment(options, "use_staged_dataflow_worker_jar");
14251399
}
14261400

14271401
Job newJob = jobSpecification.getJob();
@@ -1480,11 +1454,7 @@ public DataflowPipelineJob run(Pipeline pipeline) {
14801454
.collect(Collectors.toList());
14811455

14821456
if (minCpuFlags.isEmpty()) {
1483-
dataflowOptions.setExperiments(
1484-
ImmutableList.<String>builder()
1485-
.addAll(experiments)
1486-
.add("min_cpu_platform=" + dataflowOptions.getMinCpuPlatform())
1487-
.build());
1457+
ExperimentalOptions.addExperiment(dataflowOptions, "min_cpu_platform=" + dataflowOptions.getMinCpuPlatform());
14881458
} else {
14891459
LOG.warn(
14901460
"Flag min_cpu_platform is defined in both top level PipelineOption, "
@@ -1521,19 +1491,16 @@ public DataflowPipelineJob run(Pipeline pipeline) {
15211491
byte[] jobGraphBytes = DataflowPipelineTranslator.jobToString(newJob).getBytes(UTF_8);
15221492
int jobGraphByteSize = jobGraphBytes.length;
15231493
if (jobGraphByteSize >= CREATE_JOB_REQUEST_LIMIT_BYTES
1524-
&& !hasExperiment(options, "upload_graph")
15251494
&& !useUnifiedWorker(options)) {
1526-
List<String> experiments = firstNonNull(options.getExperiments(), Collections.emptyList());
1527-
options.setExperiments(
1528-
ImmutableList.<String>builder().addAll(experiments).add("upload_graph").build());
1495+
ExperimentalOptions.addExperiment(options, "upload_graph");
15291496
LOG.info(
15301497
"The job graph size ({} in bytes) is larger than {}. Automatically add "
15311498
+ "the upload_graph option to experiments.",
15321499
jobGraphByteSize,
15331500
CREATE_JOB_REQUEST_LIMIT_BYTES);
15341501
}
15351502

1536-
if (hasExperiment(options, "upload_graph") && useUnifiedWorker(options)) {
1503+
if (useUnifiedWorker(options)) {
15371504
ArrayList<String> experiments = new ArrayList<>(options.getExperiments());
15381505
while (experiments.remove("upload_graph")) {}
15391506
options.setExperiments(experiments);

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

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -55,6 +55,7 @@
5555
import org.apache.beam.runners.dataflow.worker.windmill.client.grpc.stubs.WindmillStubFactoryFactory;
5656
import org.apache.beam.runners.dataflow.worker.windmill.work.WorkItemReceiver;
5757
import org.apache.beam.sdk.extensions.gcp.options.GcpOptions;
58+
import org.apache.beam.sdk.options.ExperimentalOptions;
5859
import org.apache.beam.sdk.options.PipelineOptionsFactory;
5960
import org.apache.beam.sdk.util.BackOff;
6061
import org.apache.beam.sdk.util.BackOffUtils;
@@ -113,13 +114,13 @@ private static DataflowWorkerHarnessOptions testOptions(
113114
options.setProject("project");
114115
options.setJobId("job");
115116
options.setWorkerId("worker");
116-
List<String> experiments =
117-
options.getExperiments() == null ? new ArrayList<>() : options.getExperiments();
117+
118118
if (enableStreamingEngine) {
119-
experiments.add(GcpOptions.STREAMING_ENGINE_EXPERIMENT);
119+
ExperimentalOptions.addExperiment(options, GcpOptions.STREAMING_ENGINE_EXPERIMENT);
120+
}
121+
for (String experiment : additionalExperiments) {
122+
ExperimentalOptions.addExperiment(options, experiment);
120123
}
121-
experiments.addAll(additionalExperiments);
122-
options.setExperiments(experiments);
123124

124125
options.setWindmillServiceStreamingRpcBatchLimit(Integer.MAX_VALUE);
125126
options.setWindmillServiceStreamingRpcHealthCheckPeriodMs(NO_HEALTH_CHECK);

0 commit comments

Comments
 (0)