Skip to content

Commit 52b658d

Browse files
committed
Cleanly separate v1 worker and v2 sdk harness container image handling in DataflowRunner
1 parent 4f0bcf6 commit 52b658d

5 files changed

Lines changed: 105 additions & 92 deletions

File tree

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

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -139,10 +139,11 @@ private static byte[] serializeWindowingStrategy(
139139
try {
140140
SdkComponents sdkComponents = SdkComponents.create();
141141

142-
String workerHarnessContainerImageURL =
143-
DataflowRunner.getContainerImageForJob(options.as(DataflowPipelineOptions.class));
142+
String v2SdkHarnessContainerImageURL =
143+
DataflowRunner.getV2SdkHarnessContainerImageForJob(
144+
options.as(DataflowPipelineOptions.class));
144145
RunnerApi.Environment defaultEnvironmentForDataflow =
145-
Environments.createDockerEnvironment(workerHarnessContainerImageURL);
146+
Environments.createDockerEnvironment(v2SdkHarnessContainerImageURL);
146147
sdkComponents.registerEnvironment(defaultEnvironmentForDataflow);
147148

148149
return WindowingStrategyTranslation.toMessageProto(windowingStrategy, sdkComponents)

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

Lines changed: 73 additions & 42 deletions
Original file line numberDiff line numberDiff line change
@@ -518,29 +518,15 @@ static boolean isServiceEndpoint(String endpoint) {
518518
}
519519

520520
static void validateSdkContainerImageOptions(DataflowPipelineWorkerPoolOptions workerOptions) {
521-
// Check against null - empty string value for workerHarnessContainerImage
522-
// must be preserved for legacy dataflowWorkerJar to work.
523-
String sdkContainerOption = workerOptions.getSdkContainerImage();
524-
String workerHarnessOption = workerOptions.getWorkerHarnessContainerImage();
525-
Preconditions.checkArgument(
526-
sdkContainerOption == null
527-
|| workerHarnessOption == null
528-
|| sdkContainerOption.equals(workerHarnessOption),
529-
"Cannot use legacy option workerHarnessContainerImage with sdkContainerImage. Prefer sdkContainerImage.");
530-
531-
// Default to new option, which may be null.
532-
String containerImage = workerOptions.getSdkContainerImage();
533-
if (workerOptions.getWorkerHarnessContainerImage() != null
534-
&& workerOptions.getSdkContainerImage() == null) {
535-
// Set image to old option if old option was set but new option is not set.
521+
if (workerOptions.getSdkContainerImage() != null && workerOptions.getWorkerHarnessContainerImage() != null) {
536522
LOG.warn(
537-
"Prefer --sdkContainerImage over deprecated legacy option --workerHarnessContainerImage.");
538-
containerImage = workerOptions.getWorkerHarnessContainerImage();
523+
"Container specified for both --workerHarnessContainerImage and --sdkContainerImage. "
524+
+ "If you are a Beam of Dataflow developer, this could make sense, "
525+
+ "but otherwise may be a configuration error. "
526+
+ "The value of --workerHarnessContainerImage will be used only if the pipeline runs on Dataflow V1 "
527+
+ "and is *not* supported for end users. "
528+
+ "The value of --sdkContainerImage will be used only if the pipeline runs on Dataflow V2");
539529
}
540-
541-
// Make sure both options have same value.
542-
workerOptions.setSdkContainerImage(containerImage);
543-
workerOptions.setWorkerHarnessContainerImage(containerImage);
544530
}
545531

546532
@VisibleForTesting
@@ -1039,7 +1025,7 @@ protected RunnerApi.Pipeline applySdkEnvironmentOverrides(
10391025
if (containerImage.startsWith("apache/beam")
10401026
&& !updated
10411027
// don't update if the container image is already configured by DataflowRunner
1042-
&& !containerImage.equals(getContainerImageForJob(options))) {
1028+
&& !containerImage.equals(getV2SdkHarnessContainerImageForJob(options))) {
10431029
containerImage =
10441030
DataflowRunnerInfo.getDataflowRunnerInfo().getContainerImageBaseRepository()
10451031
+ containerImage.substring(containerImage.lastIndexOf("/"));
@@ -1290,15 +1276,18 @@ public DataflowPipelineJob run(Pipeline pipeline) {
12901276
+ "related to Google Compute Engine usage and other Google Cloud Services.");
12911277

12921278
DataflowPipelineOptions dataflowOptions = options.as(DataflowPipelineOptions.class);
1293-
String workerHarnessContainerImageURL = DataflowRunner.getContainerImageForJob(dataflowOptions);
1279+
String v1WorkerContainerImageURL =
1280+
DataflowRunner.getV1WorkerContainerImageForJob(dataflowOptions);
1281+
String v2SdkHarnessContainerImageURL =
1282+
DataflowRunner.getV2SdkHarnessContainerImageForJob(dataflowOptions);
12941283

12951284
// This incorrectly puns the worker harness container image (which implements v1beta3 API)
12961285
// with the SDK harness image (which implements Fn API).
12971286
//
12981287
// The same Environment is used in different and contradictory ways, depending on whether
12991288
// it is a v1 or v2 job submission.
13001289
RunnerApi.Environment defaultEnvironmentForDataflow =
1301-
Environments.createDockerEnvironment(workerHarnessContainerImageURL);
1290+
Environments.createDockerEnvironment(v2SdkHarnessContainerImageURL);
13021291

13031292
// The SdkComponents for portable an non-portable job submission must be kept distinct. Both
13041293
// need the default environment.
@@ -1469,7 +1458,7 @@ public DataflowPipelineJob run(Pipeline pipeline) {
14691458
// For runner_v1, only worker_harness_container is set.
14701459
// For runner_v2, both worker_harness_container and sdk_harness_container are set to the same
14711460
// value.
1472-
String containerImage = getContainerImageForJob(options);
1461+
String containerImage = getV1WorkerContainerImageForJob(options);
14731462
for (WorkerPool workerPool : newJob.getEnvironment().getWorkerPools()) {
14741463
workerPool.setWorkerHarnessContainerImage(containerImage);
14751464
}
@@ -2634,59 +2623,101 @@ public Map<PCollection<?>, ReplacementOutput> mapOutputs(
26342623
}
26352624

26362625
@VisibleForTesting
2637-
static String getContainerImageForJob(DataflowPipelineOptions options) {
2626+
static String getV1WorkerContainerImageForJob(DataflowPipelineOptions options) {
2627+
String containerImage = options.getWorkerHarnessContainerImage();
2628+
2629+
if (containerImage == null) {
2630+
// If not set, construct and return default image URL.
2631+
return getDefaultV1WorkerContainerImageUrl(options);
2632+
} else if (containerImage.contains("IMAGE")) {
2633+
// Replace placeholder with default image name
2634+
return containerImage.replace("IMAGE", getDefaultV1WorkerContainerImageNameForJob(options));
2635+
} else {
2636+
return containerImage;
2637+
}
2638+
}
2639+
2640+
static String getV2SdkHarnessContainerImageForJob(DataflowPipelineOptions options) {
26382641
String containerImage = options.getSdkContainerImage();
26392642

26402643
if (containerImage == null) {
26412644
// If not set, construct and return default image URL.
2642-
return getDefaultContainerImageUrl(options);
2645+
return getDefaultV2SdkHarnessContainerImageUrl(options);
26432646
} else if (containerImage.contains("IMAGE")) {
26442647
// Replace placeholder with default image name
2645-
return containerImage.replace("IMAGE", getDefaultContainerImageNameForJob(options));
2648+
return containerImage.replace("IMAGE", getDefaultV2SdkHarnessContainerImageNameForJob());
26462649
} else {
26472650
return containerImage;
26482651
}
26492652
}
26502653

2651-
/** Construct the default Dataflow container full image URL. */
2652-
static String getDefaultContainerImageUrl(DataflowPipelineOptions options) {
2654+
/** Construct the default Dataflow worker container full image URL. */
2655+
static String getDefaultV1WorkerContainerImageUrl(DataflowPipelineOptions options) {
26532656
DataflowRunnerInfo dataflowRunnerInfo = DataflowRunnerInfo.getDataflowRunnerInfo();
26542657
return String.format(
26552658
"%s/%s:%s",
26562659
dataflowRunnerInfo.getContainerImageBaseRepository(),
2657-
getDefaultContainerImageNameForJob(options),
2658-
getDefaultContainerVersion(options));
2660+
getDefaultV1WorkerContainerImageNameForJob(options),
2661+
getDefaultV1WorkerContainerVersion(options));
2662+
}
2663+
2664+
/** Construct the default Java SDK container full image URL. */
2665+
static String getDefaultV2SdkHarnessContainerImageUrl(DataflowPipelineOptions options) {
2666+
DataflowRunnerInfo dataflowRunnerInfo = DataflowRunnerInfo.getDataflowRunnerInfo();
2667+
return String.format(
2668+
"%s/%s:%s",
2669+
dataflowRunnerInfo.getContainerImageBaseRepository(),
2670+
getDefaultV2SdkHarnessContainerImageNameForJob(),
2671+
getDefaultV2SdkHarnessContainerVersion(options));
26592672
}
26602673

26612674
/**
2662-
* Construct the default Dataflow container image name based on pipeline type and Java version.
2675+
* Construct the default Dataflow V1 worker container image name based on pipeline type and Java
2676+
* version.
26632677
*/
2664-
static String getDefaultContainerImageNameForJob(DataflowPipelineOptions options) {
2678+
static String getDefaultV1WorkerContainerImageNameForJob(DataflowPipelineOptions options) {
26652679
Environments.JavaVersion javaVersion = Environments.getJavaVersion();
2666-
if (useUnifiedWorker(options)) {
2667-
return String.format("beam_%s_sdk", javaVersion.name());
2668-
} else if (options.isStreaming()) {
2680+
if (options.isStreaming()) {
26692681
return String.format("beam-%s-streaming", javaVersion.legacyName());
26702682
} else {
26712683
return String.format("beam-%s-batch", javaVersion.legacyName());
26722684
}
26732685
}
26742686

26752687
/**
2676-
* Construct the default Dataflow container image name based on pipeline type and Java version.
2688+
* Construct the default Java SDK container image name based on pipeline type and Java version,
2689+
* for use by Dataflow V2.
2690+
*/
2691+
static String getDefaultV2SdkHarnessContainerImageNameForJob() {
2692+
Environments.JavaVersion javaVersion = Environments.getJavaVersion();
2693+
return String.format("beam_%s_sdk", javaVersion.name());
2694+
}
2695+
2696+
/**
2697+
* Construct the default Dataflow V1 worker container image name based on pipeline type and Java
2698+
* version.
26772699
*/
2678-
static String getDefaultContainerVersion(DataflowPipelineOptions options) {
2700+
static String getDefaultV1WorkerContainerVersion(DataflowPipelineOptions options) {
26792701
DataflowRunnerInfo dataflowRunnerInfo = DataflowRunnerInfo.getDataflowRunnerInfo();
26802702
ReleaseInfo releaseInfo = ReleaseInfo.getReleaseInfo();
26812703
if (releaseInfo.isDevSdkVersion()) {
2682-
if (useUnifiedWorker(options)) {
2683-
return dataflowRunnerInfo.getFnApiDevContainerVersion();
2684-
}
26852704
return dataflowRunnerInfo.getLegacyDevContainerVersion();
26862705
}
26872706
return releaseInfo.getSdkVersion();
26882707
}
26892708

2709+
/**
2710+
* Construct the default Dataflow container image name based on pipeline type and Java version.
2711+
*/
2712+
static String getDefaultV2SdkHarnessContainerVersion(DataflowPipelineOptions options) {
2713+
DataflowRunnerInfo dataflowRunnerInfo = DataflowRunnerInfo.getDataflowRunnerInfo();
2714+
ReleaseInfo releaseInfo = ReleaseInfo.getReleaseInfo();
2715+
if (releaseInfo.isDevSdkVersion()) {
2716+
return dataflowRunnerInfo.getFnApiDevContainerVersion();
2717+
}
2718+
return releaseInfo.getSdkVersion();
2719+
}
2720+
26902721
static boolean useUnifiedWorker(DataflowPipelineOptions options) {
26912722
return hasExperiment(options, "beam_fn_api")
26922723
|| hasExperiment(options, "use_runner_v2")

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

Lines changed: 2 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -104,28 +104,19 @@ public String getAlgorithm() {
104104
void setDiskSizeGb(int value);
105105

106106
/** Container image used as Dataflow worker harness image. */
107-
/** @deprecated Use {@link #getSdkContainerImage} instead. */
108107
@Description(
109-
"Container image used to configure a Dataflow worker. "
110-
+ "Can only be used for official Dataflow container images. "
111-
+ "Prefer using sdkContainerImage instead.")
112-
@Deprecated
108+
"Container image to use for Dataflow V1 worker. Can only be used for official Dataflow container images. ")
113109
@Hidden
114110
String getWorkerHarnessContainerImage();
115111

116-
/** @deprecated Use {@link #setSdkContainerImage} instead. */
117-
@Deprecated
118112
@Hidden
119113
void setWorkerHarnessContainerImage(String value);
120114

121115
/**
122116
* Container image used to configure SDK execution environment on worker. Used for custom
123117
* containers on portable pipelines only.
124118
*/
125-
@Description(
126-
"Container image used to configure the SDK execution environment of "
127-
+ "pipeline code on a worker. For non-portable pipelines, can only be "
128-
+ "used for official Dataflow container images.")
119+
@Description("Container image to use for Beam Java SDK execution environment on Dataflow V2.")
129120
String getSdkContainerImage();
130121

131122
void setSdkContainerImage(String value);

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

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -156,7 +156,8 @@ private SdkComponents createSdkComponents(PipelineOptions options) {
156156
SdkComponents sdkComponents = SdkComponents.create();
157157

158158
String containerImageURL =
159-
DataflowRunner.getContainerImageForJob(options.as(DataflowPipelineOptions.class));
159+
DataflowRunner.getV2SdkHarnessContainerImageForJob(
160+
options.as(DataflowPipelineOptions.class));
160161
RunnerApi.Environment defaultEnvironmentForDataflow =
161162
Environments.createDockerEnvironment(containerImageURL);
162163

@@ -1127,7 +1128,8 @@ public String apply(byte[] input) {
11271128
file2.deleteOnExit();
11281129
SdkComponents sdkComponents = SdkComponents.create();
11291130
sdkComponents.registerEnvironment(
1130-
Environments.createDockerEnvironment(DataflowRunner.getContainerImageForJob(options))
1131+
Environments.createDockerEnvironment(
1132+
DataflowRunner.getV2SdkHarnessContainerImageForJob(options))
11311133
.toBuilder()
11321134
.addAllDependencies(
11331135
Environments.getArtifacts(
@@ -1589,7 +1591,8 @@ public void testSetWorkerHarnessContainerImageInPipelineProto() throws Exception
15891591
Iterables.getOnlyElement(pipelineProto.getComponents().getEnvironmentsMap().values());
15901592

15911593
DockerPayload payload = DockerPayload.parseFrom(defaultEnvironment.getPayload());
1592-
assertEquals(DataflowRunner.getContainerImageForJob(options), payload.getContainerImage());
1594+
assertEquals(
1595+
DataflowRunner.getV2SdkHarnessContainerImageForJob(options), payload.getContainerImage());
15931596
}
15941597

15951598
/**
@@ -1621,7 +1624,8 @@ public void testSetSdkContainerImageInPipelineProto() throws Exception {
16211624
Iterables.getOnlyElement(pipelineProto.getComponents().getEnvironmentsMap().values());
16221625

16231626
DockerPayload payload = DockerPayload.parseFrom(defaultEnvironment.getPayload());
1624-
assertEquals(DataflowRunner.getContainerImageForJob(options), payload.getContainerImage());
1627+
assertEquals(
1628+
DataflowRunner.getV2SdkHarnessContainerImageForJob(options), payload.getContainerImage());
16251629
}
16261630

16271631
@Test

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

Lines changed: 18 additions & 32 deletions
Original file line numberDiff line numberDiff line change
@@ -17,7 +17,6 @@
1717
*/
1818
package org.apache.beam.runners.dataflow;
1919

20-
import static org.apache.beam.runners.dataflow.DataflowRunner.getContainerImageForJob;
2120
import static org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.io.Files.getFileExtension;
2221
import static org.hamcrest.MatcherAssert.assertThat;
2322
import static org.hamcrest.Matchers.containsInAnyOrder;
@@ -644,28 +643,6 @@ public void testZoneAliasWorkerZone() {
644643
assertEquals("us-east1-b", options.getWorkerZone());
645644
}
646645

647-
@Test
648-
public void testAliasForLegacyWorkerHarnessContainerImage() {
649-
DataflowPipelineWorkerPoolOptions options =
650-
PipelineOptionsFactory.as(DataflowPipelineWorkerPoolOptions.class);
651-
String testImage = "image.url:worker";
652-
options.setWorkerHarnessContainerImage(testImage);
653-
DataflowRunner.validateWorkerSettings(options);
654-
assertEquals(testImage, options.getWorkerHarnessContainerImage());
655-
assertEquals(testImage, options.getSdkContainerImage());
656-
}
657-
658-
@Test
659-
public void testAliasForSdkContainerImage() {
660-
DataflowPipelineWorkerPoolOptions options =
661-
PipelineOptionsFactory.as(DataflowPipelineWorkerPoolOptions.class);
662-
String testImage = "image.url:sdk";
663-
options.setSdkContainerImage("image.url:sdk");
664-
DataflowRunner.validateWorkerSettings(options);
665-
assertEquals(testImage, options.getWorkerHarnessContainerImage());
666-
assertEquals(testImage, options.getSdkContainerImage());
667-
}
668-
669646
@Test
670647
public void testRegionRequiredForServiceRunner() throws IOException {
671648
DataflowPipelineOptions options = buildPipelineOptions();
@@ -1736,7 +1713,7 @@ private void verifySdkHarnessConfiguration(DataflowPipelineOptions options) {
17361713

17371714
p.apply(Create.of(Arrays.asList(1, 2, 3)));
17381715

1739-
String defaultSdkContainerImage = DataflowRunner.getContainerImageForJob(options);
1716+
String defaultSdkContainerImage = DataflowRunner.getV2SdkHarnessContainerImageForJob(options);
17401717
SdkComponents sdkComponents = SdkComponents.create();
17411718
RunnerApi.Environment defaultEnvironmentForDataflow =
17421719
Environments.createDockerEnvironment(defaultSdkContainerImage);
@@ -2027,7 +2004,7 @@ public void close() {}
20272004
}
20282005

20292006
@Test
2030-
public void testGetContainerImageForJobFromOption() {
2007+
public void testGetV2SdkHarnessContainerImageForJobFromOption() {
20312008
DataflowPipelineOptions options = PipelineOptionsFactory.as(DataflowPipelineOptions.class);
20322009

20332010
String[] testCases = {
@@ -2042,43 +2019,52 @@ public void testGetContainerImageForJobFromOption() {
20422019
for (String testCase : testCases) {
20432020
// When image option is set, should use that exact image.
20442021
options.setSdkContainerImage(testCase);
2045-
assertThat(getContainerImageForJob(options), equalTo(testCase));
2022+
assertThat(DataflowRunner.getV2SdkHarnessContainerImageForJob(options), equalTo(testCase));
20462023
}
20472024
}
20482025

20492026
@Test
2050-
public void testGetContainerImageForJobFromOptionWithPlaceholder() {
2027+
public void testGetV1WorkerContainerImageForJobFromOptionWithPlaceholder() {
20512028
DataflowPipelineOptions options = PipelineOptionsFactory.as(DataflowPipelineOptions.class);
2052-
options.setSdkContainerImage("gcr.io/IMAGE/foo");
2029+
options.setWorkerHarnessContainerImage("gcr.io/IMAGE/foo");
20532030

20542031
for (Environments.JavaVersion javaVersion : Environments.JavaVersion.values()) {
20552032
System.setProperty("java.specification.version", javaVersion.specification());
20562033
// batch legacy
20572034
options.setExperiments(null);
20582035
options.setStreaming(false);
20592036
assertThat(
2060-
getContainerImageForJob(options),
2037+
DataflowRunner.getV1WorkerContainerImageForJob(options),
20612038
equalTo(String.format("gcr.io/beam-%s-batch/foo", javaVersion.legacyName())));
20622039

20632040
// streaming, legacy
20642041
options.setExperiments(null);
20652042
options.setStreaming(true);
20662043
assertThat(
2067-
getContainerImageForJob(options),
2044+
DataflowRunner.getV1WorkerContainerImageForJob(options),
20682045
equalTo(String.format("gcr.io/beam-%s-streaming/foo", javaVersion.legacyName())));
2046+
}
2047+
}
2048+
2049+
@Test
2050+
public void testGetV2SdkHarnessContainerImageForJobFromOptionWithPlaceholder() {
2051+
DataflowPipelineOptions options = PipelineOptionsFactory.as(DataflowPipelineOptions.class);
2052+
options.setSdkContainerImage("gcr.io/IMAGE/foo");
20692053

2054+
for (Environments.JavaVersion javaVersion : Environments.JavaVersion.values()) {
2055+
System.setProperty("java.specification.version", javaVersion.specification());
20702056
// batch, FnAPI
20712057
options.setExperiments(ImmutableList.of("beam_fn_api"));
20722058
options.setStreaming(false);
20732059
assertThat(
2074-
getContainerImageForJob(options),
2060+
DataflowRunner.getV2SdkHarnessContainerImageForJob(options),
20752061
equalTo(String.format("gcr.io/beam_%s_sdk/foo", javaVersion.name())));
20762062

20772063
// streaming, FnAPI
20782064
options.setExperiments(ImmutableList.of("beam_fn_api"));
20792065
options.setStreaming(true);
20802066
assertThat(
2081-
getContainerImageForJob(options),
2067+
DataflowRunner.getV2SdkHarnessContainerImageForJob(options),
20822068
equalTo(String.format("gcr.io/beam_%s_sdk/foo", javaVersion.name())));
20832069
}
20842070
}

0 commit comments

Comments
 (0)