Skip to content

Commit 25ea087

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 d841b8d commit 25ea087

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 Dataflow Portable Runner.
12461246
if (DataflowRunner.isMultiLanguagePipeline(pipeline) || includesTransformUpgrades(pipeline)) {
1247-
if (!useUnifiedWorker(options)) {
1248-
List<String> experiments = firstNonNull(options.getExperiments(), Collections.emptyList());
1247+
if (!firstNonNull(options.getExperiments(), Collections.emptyList()).contains("use_runner_v2")) {
12491248
LOG.info(
12501249
"Automatically enabling Dataflow Portable Runner 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)) {
@@ -1262,21 +1260,10 @@ public DataflowPipelineJob run(Pipeline pipeline) {
12621260
throw new IllegalArgumentException(
12631261
"Dataflow Portable Runner both disabled and enabled: at least one of ['enable_portable_runner', 'beam_fn_api', 'use_unified_worker', 'use_runner_v2', 'use_portable_job_submission'] is set and also one of ['enable_streaming_java_runner', 'disable_portable_runner', 'disable_runner_v2', 'disable_runner_v2_until_2023', 'disable_prime_runner_v2'] is set.");
12641262
}
1265-
List<String> experiments =
1266-
new ArrayList<>(options.getExperiments()); // non-null if useUnifiedWorker is true
1267-
if (!experiments.contains("use_runner_v2")) {
1268-
experiments.add("use_runner_v2");
1269-
}
1270-
if (!experiments.contains("use_unified_worker")) {
1271-
experiments.add("use_unified_worker");
1272-
}
1273-
if (!experiments.contains("beam_fn_api")) {
1274-
experiments.add("beam_fn_api");
1275-
}
1276-
if (!experiments.contains("use_portable_job_submission")) {
1277-
experiments.add("use_portable_job_submission");
1278-
}
1279-
options.setExperiments(ImmutableList.copyOf(experiments));
1263+
ExperimentalOptions.addExperiment(options, "use_runner_v2");
1264+
ExperimentalOptions.addExperiment(options, "use_unified_worker");
1265+
ExperimentalOptions.addExperiment(options, "beam_fn_api");
1266+
ExperimentalOptions.addExperiment(options, "use_portable_job_submission");
12801267
// Ensure that logging via the FnApi is enabled
12811268
options.as(SdkHarnessOptions.class).setEnableLogViaFnApi(true);
12821269
}
@@ -1303,14 +1290,9 @@ public DataflowPipelineJob run(Pipeline pipeline) {
13031290
options.setStreaming(true);
13041291

13051292
{
1306-
List<String> experiments =
1307-
options.getExperiments() == null
1308-
? new ArrayList<>()
1309-
: new ArrayList<>(options.getExperiments());
13101293
// Experiment marking that the harness supports tag encoding v2
13111294
// Backend will enable tag encoding v2 only if the harness supports it.
1312-
experiments.add("streaming_engine_state_tag_encoding_v2_supported");
1313-
options.setExperiments(ImmutableList.copyOf(experiments));
1295+
ExperimentalOptions.addExperiment(options, "streaming_engine_state_tag_encoding_v2_supported");
13141296
}
13151297

13161298
if (useUnifiedWorker(options)) {
@@ -1421,15 +1403,7 @@ public DataflowPipelineJob run(Pipeline pipeline) {
14211403
pipeline, dataflowV1PipelineProto, dataflowV1Components, this, packages);
14221404

14231405
if (!isNullOrEmpty(dataflowOptions.getDataflowWorkerJar()) && !useUnifiedWorker(options)) {
1424-
List<String> experiments =
1425-
firstNonNull(dataflowOptions.getExperiments(), Collections.emptyList());
1426-
if (!experiments.contains("use_staged_dataflow_worker_jar")) {
1427-
dataflowOptions.setExperiments(
1428-
ImmutableList.<String>builder()
1429-
.addAll(experiments)
1430-
.add("use_staged_dataflow_worker_jar")
1431-
.build());
1432-
}
1406+
ExperimentalOptions.addExperiment(options, "use_staged_dataflow_worker_jar");
14331407
}
14341408

14351409
Job newJob = jobSpecification.getJob();
@@ -1488,11 +1462,7 @@ public DataflowPipelineJob run(Pipeline pipeline) {
14881462
.collect(Collectors.toList());
14891463

14901464
if (minCpuFlags.isEmpty()) {
1491-
dataflowOptions.setExperiments(
1492-
ImmutableList.<String>builder()
1493-
.addAll(experiments)
1494-
.add("min_cpu_platform=" + dataflowOptions.getMinCpuPlatform())
1495-
.build());
1465+
ExperimentalOptions.addExperiment(dataflowOptions, "min_cpu_platform=" + dataflowOptions.getMinCpuPlatform());
14961466
} else {
14971467
LOG.warn(
14981468
"Flag min_cpu_platform is defined in both top level PipelineOption, "
@@ -1529,19 +1499,16 @@ public DataflowPipelineJob run(Pipeline pipeline) {
15291499
byte[] jobGraphBytes = DataflowPipelineTranslator.jobToString(newJob).getBytes(UTF_8);
15301500
int jobGraphByteSize = jobGraphBytes.length;
15311501
if (jobGraphByteSize >= CREATE_JOB_REQUEST_LIMIT_BYTES
1532-
&& !hasExperiment(options, "upload_graph")
15331502
&& !useUnifiedWorker(options)) {
1534-
List<String> experiments = firstNonNull(options.getExperiments(), Collections.emptyList());
1535-
options.setExperiments(
1536-
ImmutableList.<String>builder().addAll(experiments).add("upload_graph").build());
1503+
ExperimentalOptions.addExperiment(options, "upload_graph");
15371504
LOG.info(
15381505
"The job graph size ({} in bytes) is larger than {}. Automatically add "
15391506
+ "the upload_graph option to experiments.",
15401507
jobGraphByteSize,
15411508
CREATE_JOB_REQUEST_LIMIT_BYTES);
15421509
}
15431510

1544-
if (hasExperiment(options, "upload_graph") && useUnifiedWorker(options)) {
1511+
if (useUnifiedWorker(options)) {
15451512
ArrayList<String> experiments = new ArrayList<>(options.getExperiments());
15461513
while (experiments.remove("upload_graph")) {}
15471514
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)