Skip to content

Commit 81828fd

Browse files
authored
[Dataflow] Added Portable Runner alias to java runners (#38411)
* Added portable runner options to java runner * spotless * Added more tolerance to flaky test * Removed unused experiments * Added experiments back in
1 parent b60082d commit 81828fd

3 files changed

Lines changed: 28 additions & 9 deletions

File tree

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

Lines changed: 7 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1244,8 +1244,8 @@ 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 (!useUnifiedWorker(options)) {
1248+
List<String> experiments = firstNonNull(options.getExperiments(), Collections.emptyList());
12491249
LOG.info(
12501250
"Automatically enabling Dataflow Runner v2 since the pipeline used cross-language"
12511251
+ " transforms or pipeline needed a transform upgrade.");
@@ -1256,7 +1256,9 @@ public DataflowPipelineJob run(Pipeline pipeline) {
12561256
if (useUnifiedWorker(options)) {
12571257
if (hasExperiment(options, "disable_runner_v2")
12581258
|| hasExperiment(options, "disable_runner_v2_until_2023")
1259-
|| hasExperiment(options, "disable_prime_runner_v2")) {
1259+
|| hasExperiment(options, "disable_prime_runner_v2")
1260+
|| hasExperiment(options, "disable_portable_runner")
1261+
|| hasExperiment(options, "enable_streaming_java_runner")) {
12601262
throw new IllegalArgumentException(
12611263
"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.");
12621264
}
@@ -2729,7 +2731,8 @@ static boolean useUnifiedWorker(DataflowPipelineOptions options) {
27292731
return hasExperiment(options, "beam_fn_api")
27302732
|| hasExperiment(options, "use_runner_v2")
27312733
|| hasExperiment(options, "use_unified_worker")
2732-
|| hasExperiment(options, "use_portable_job_submission");
2734+
|| hasExperiment(options, "use_portable_job_submission")
2735+
|| hasExperiment(options, "enable_portable_runner");
27332736
}
27342737

27352738
static void verifyDoFnSupported(

runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowRunnerTest.java

Lines changed: 20 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1783,7 +1783,11 @@ public void testSdkHarnessConfigurationPrime() throws IOException {
17831783
public void testSettingAnyFnApiExperimentEnablesUnifiedWorker() throws Exception {
17841784
for (String experiment :
17851785
ImmutableList.of(
1786-
"beam_fn_api", "use_runner_v2", "use_unified_worker", "use_portable_job_submission")) {
1786+
"beam_fn_api",
1787+
"use_runner_v2",
1788+
"use_unified_worker",
1789+
"use_portable_job_submission",
1790+
"enable_portable_runner")) {
17871791
DataflowPipelineOptions options = buildPipelineOptions();
17881792
ExperimentalOptions.addExperiment(options, experiment);
17891793
Pipeline p = Pipeline.create(options);
@@ -1798,7 +1802,11 @@ public void testSettingAnyFnApiExperimentEnablesUnifiedWorker() throws Exception
17981802

17991803
for (String experiment :
18001804
ImmutableList.of(
1801-
"beam_fn_api", "use_runner_v2", "use_unified_worker", "use_portable_job_submission")) {
1805+
"beam_fn_api",
1806+
"use_runner_v2",
1807+
"use_unified_worker",
1808+
"use_portable_job_submission",
1809+
"enable_portable_runner")) {
18021810
DataflowPipelineOptions options = buildPipelineOptions();
18031811
options.setStreaming(true);
18041812
ExperimentalOptions.addExperiment(options, experiment);
@@ -1822,10 +1830,18 @@ public void testSettingAnyFnApiExperimentEnablesUnifiedWorker() throws Exception
18221830
public void testSettingConflictingEnableAndDisableExperimentsThrowsException() throws Exception {
18231831
for (String experiment :
18241832
ImmutableList.of(
1825-
"beam_fn_api", "use_runner_v2", "use_unified_worker", "use_portable_job_submission")) {
1833+
"beam_fn_api",
1834+
"use_runner_v2",
1835+
"use_unified_worker",
1836+
"use_portable_job_submission",
1837+
"enable_portable_runner")) {
18261838
for (String disabledExperiment :
18271839
ImmutableList.of(
1828-
"disable_runner_v2", "disable_runner_v2_until_2023", "disable_prime_runner_v2")) {
1840+
"disable_runner_v2",
1841+
"disable_runner_v2_until_2023",
1842+
"disable_prime_runner_v2",
1843+
"enable_streaming_java_runner",
1844+
"disable_portable_runner")) {
18291845
DataflowPipelineOptions options = buildPipelineOptions();
18301846
ExperimentalOptions.addExperiment(options, experiment);
18311847
ExperimentalOptions.addExperiment(options, disabledExperiment);

sdks/java/core/src/test/java/org/apache/beam/sdk/util/UnboundedScheduledExecutorServiceTest.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -625,7 +625,7 @@ public void testThreadsAreAddedOnlyAsNeededWithContention() throws Exception {
625625
LOG.info("Created {} threads to execute at most 100 parallel tasks", largestPool);
626626
// Ideally we would never create more than 100, however with contention it is still possible
627627
// some extra threads will be created.
628-
assertTrue(largestPool <= 110);
628+
assertTrue(largestPool <= 120);
629629
executorService.shutdown();
630630
}
631631
}

0 commit comments

Comments
 (0)