Skip to content

Commit 0b5d216

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

5 files changed

Lines changed: 106 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: 74 additions & 42 deletions
Original file line numberDiff line numberDiff line change
@@ -518,29 +518,16 @@ 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
522+
&& workerOptions.getWorkerHarnessContainerImage() != null) {
536523
LOG.warn(
537-
"Prefer --sdkContainerImage over deprecated legacy option --workerHarnessContainerImage.");
538-
containerImage = workerOptions.getWorkerHarnessContainerImage();
524+
"Container specified for both --workerHarnessContainerImage and --sdkContainerImage. "
525+
+ "If you are a Beam of Dataflow developer, this could make sense, "
526+
+ "but otherwise may be a configuration error. "
527+
+ "The value of --workerHarnessContainerImage will be used only if the pipeline runs on Dataflow V1 "
528+
+ "and is *not* supported for end users. "
529+
+ "The value of --sdkContainerImage will be used only if the pipeline runs on Dataflow V2");
539530
}
540-
541-
// Make sure both options have same value.
542-
workerOptions.setSdkContainerImage(containerImage);
543-
workerOptions.setWorkerHarnessContainerImage(containerImage);
544531
}
545532

546533
@VisibleForTesting
@@ -1039,7 +1026,7 @@ protected RunnerApi.Pipeline applySdkEnvironmentOverrides(
10391026
if (containerImage.startsWith("apache/beam")
10401027
&& !updated
10411028
// don't update if the container image is already configured by DataflowRunner
1042-
&& !containerImage.equals(getContainerImageForJob(options))) {
1029+
&& !containerImage.equals(getV2SdkHarnessContainerImageForJob(options))) {
10431030
containerImage =
10441031
DataflowRunnerInfo.getDataflowRunnerInfo().getContainerImageBaseRepository()
10451032
+ containerImage.substring(containerImage.lastIndexOf("/"));
@@ -1290,15 +1277,18 @@ public DataflowPipelineJob run(Pipeline pipeline) {
12901277
+ "related to Google Compute Engine usage and other Google Cloud Services.");
12911278

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

12951285
// This incorrectly puns the worker harness container image (which implements v1beta3 API)
12961286
// with the SDK harness image (which implements Fn API).
12971287
//
12981288
// The same Environment is used in different and contradictory ways, depending on whether
12991289
// it is a v1 or v2 job submission.
13001290
RunnerApi.Environment defaultEnvironmentForDataflow =
1301-
Environments.createDockerEnvironment(workerHarnessContainerImageURL);
1291+
Environments.createDockerEnvironment(v2SdkHarnessContainerImageURL);
13021292

13031293
// The SdkComponents for portable an non-portable job submission must be kept distinct. Both
13041294
// need the default environment.
@@ -1469,7 +1459,7 @@ public DataflowPipelineJob run(Pipeline pipeline) {
14691459
// For runner_v1, only worker_harness_container is set.
14701460
// For runner_v2, both worker_harness_container and sdk_harness_container are set to the same
14711461
// value.
1472-
String containerImage = getContainerImageForJob(options);
1462+
String containerImage = getV1WorkerContainerImageForJob(options);
14731463
for (WorkerPool workerPool : newJob.getEnvironment().getWorkerPools()) {
14741464
workerPool.setWorkerHarnessContainerImage(containerImage);
14751465
}
@@ -2634,59 +2624,101 @@ public Map<PCollection<?>, ReplacementOutput> mapOutputs(
26342624
}
26352625

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

26402644
if (containerImage == null) {
26412645
// If not set, construct and return default image URL.
2642-
return getDefaultContainerImageUrl(options);
2646+
return getDefaultV2SdkHarnessContainerImageUrl(options);
26432647
} else if (containerImage.contains("IMAGE")) {
26442648
// Replace placeholder with default image name
2645-
return containerImage.replace("IMAGE", getDefaultContainerImageNameForJob(options));
2649+
return containerImage.replace("IMAGE", getDefaultV2SdkHarnessContainerImageNameForJob());
26462650
} else {
26472651
return containerImage;
26482652
}
26492653
}
26502654

2651-
/** Construct the default Dataflow container full image URL. */
2652-
static String getDefaultContainerImageUrl(DataflowPipelineOptions options) {
2655+
/** Construct the default Dataflow worker container full image URL. */
2656+
static String getDefaultV1WorkerContainerImageUrl(DataflowPipelineOptions options) {
26532657
DataflowRunnerInfo dataflowRunnerInfo = DataflowRunnerInfo.getDataflowRunnerInfo();
26542658
return String.format(
26552659
"%s/%s:%s",
26562660
dataflowRunnerInfo.getContainerImageBaseRepository(),
2657-
getDefaultContainerImageNameForJob(options),
2658-
getDefaultContainerVersion(options));
2661+
getDefaultV1WorkerContainerImageNameForJob(options),
2662+
getDefaultV1WorkerContainerVersion(options));
2663+
}
2664+
2665+
/** Construct the default Java SDK container full image URL. */
2666+
static String getDefaultV2SdkHarnessContainerImageUrl(DataflowPipelineOptions options) {
2667+
DataflowRunnerInfo dataflowRunnerInfo = DataflowRunnerInfo.getDataflowRunnerInfo();
2668+
return String.format(
2669+
"%s/%s:%s",
2670+
dataflowRunnerInfo.getContainerImageBaseRepository(),
2671+
getDefaultV2SdkHarnessContainerImageNameForJob(),
2672+
getDefaultV2SdkHarnessContainerVersion(options));
26592673
}
26602674

26612675
/**
2662-
* Construct the default Dataflow container image name based on pipeline type and Java version.
2676+
* Construct the default Dataflow V1 worker container image name based on pipeline type and Java
2677+
* version.
26632678
*/
2664-
static String getDefaultContainerImageNameForJob(DataflowPipelineOptions options) {
2679+
static String getDefaultV1WorkerContainerImageNameForJob(DataflowPipelineOptions options) {
26652680
Environments.JavaVersion javaVersion = Environments.getJavaVersion();
2666-
if (useUnifiedWorker(options)) {
2667-
return String.format("beam_%s_sdk", javaVersion.name());
2668-
} else if (options.isStreaming()) {
2681+
if (options.isStreaming()) {
26692682
return String.format("beam-%s-streaming", javaVersion.legacyName());
26702683
} else {
26712684
return String.format("beam-%s-batch", javaVersion.legacyName());
26722685
}
26732686
}
26742687

26752688
/**
2676-
* Construct the default Dataflow container image name based on pipeline type and Java version.
2689+
* Construct the default Java SDK container image name based on pipeline type and Java version,
2690+
* for use by Dataflow V2.
2691+
*/
2692+
static String getDefaultV2SdkHarnessContainerImageNameForJob() {
2693+
Environments.JavaVersion javaVersion = Environments.getJavaVersion();
2694+
return String.format("beam_%s_sdk", javaVersion.name());
2695+
}
2696+
2697+
/**
2698+
* Construct the default Dataflow V1 worker container image name based on pipeline type and Java
2699+
* version.
26772700
*/
2678-
static String getDefaultContainerVersion(DataflowPipelineOptions options) {
2701+
static String getDefaultV1WorkerContainerVersion(DataflowPipelineOptions options) {
26792702
DataflowRunnerInfo dataflowRunnerInfo = DataflowRunnerInfo.getDataflowRunnerInfo();
26802703
ReleaseInfo releaseInfo = ReleaseInfo.getReleaseInfo();
26812704
if (releaseInfo.isDevSdkVersion()) {
2682-
if (useUnifiedWorker(options)) {
2683-
return dataflowRunnerInfo.getFnApiDevContainerVersion();
2684-
}
26852705
return dataflowRunnerInfo.getLegacyDevContainerVersion();
26862706
}
26872707
return releaseInfo.getSdkVersion();
26882708
}
26892709

2710+
/**
2711+
* Construct the default Dataflow container image name based on pipeline type and Java version.
2712+
*/
2713+
static String getDefaultV2SdkHarnessContainerVersion(DataflowPipelineOptions options) {
2714+
DataflowRunnerInfo dataflowRunnerInfo = DataflowRunnerInfo.getDataflowRunnerInfo();
2715+
ReleaseInfo releaseInfo = ReleaseInfo.getReleaseInfo();
2716+
if (releaseInfo.isDevSdkVersion()) {
2717+
return dataflowRunnerInfo.getFnApiDevContainerVersion();
2718+
}
2719+
return releaseInfo.getSdkVersion();
2720+
}
2721+
26902722
static boolean useUnifiedWorker(DataflowPipelineOptions options) {
26912723
return hasExperiment(options, "beam_fn_api")
26922724
|| 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)