Skip to content

Commit 615fcc5

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

5 files changed

Lines changed: 109 additions & 100 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: 77 additions & 50 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,21 +1277,19 @@ 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

1295-
// This incorrectly puns the worker harness container image (which implements v1beta3 API)
1296-
// with the SDK harness image (which implements Fn API).
1297-
//
1298-
// The same Environment is used in different and contradictory ways, depending on whether
1299-
// it is a v1 or v2 job submission.
1300-
RunnerApi.Environment defaultEnvironmentForDataflow =
1301-
Environments.createDockerEnvironment(workerHarnessContainerImageURL);
1285+
RunnerApi.Environment defaultEnvironmentForDataflowV2 =
1286+
Environments.createDockerEnvironment(v2SdkHarnessContainerImageURL);
13021287

13031288
// The SdkComponents for portable an non-portable job submission must be kept distinct. Both
13041289
// need the default environment.
13051290
SdkComponents portableComponents = SdkComponents.create();
13061291
portableComponents.registerEnvironment(
1307-
defaultEnvironmentForDataflow
1292+
defaultEnvironmentForDataflowV2
13081293
.toBuilder()
13091294
.addAllDependencies(getDefaultArtifacts())
13101295
.addAllCapabilities(Environments.getJavaCapabilities())
@@ -1343,7 +1328,7 @@ public DataflowPipelineJob run(Pipeline pipeline) {
13431328
// Capture the SdkComponents for look up during step translations
13441329
SdkComponents dataflowV1Components = SdkComponents.create();
13451330
dataflowV1Components.registerEnvironment(
1346-
defaultEnvironmentForDataflow
1331+
defaultEnvironmentForDataflowV2
13471332
.toBuilder()
13481333
.addAllDependencies(getDefaultArtifacts())
13491334
.addAllCapabilities(Environments.getJavaCapabilities())
@@ -1469,7 +1454,7 @@ public DataflowPipelineJob run(Pipeline pipeline) {
14691454
// For runner_v1, only worker_harness_container is set.
14701455
// For runner_v2, both worker_harness_container and sdk_harness_container are set to the same
14711456
// value.
1472-
String containerImage = getContainerImageForJob(options);
1457+
String containerImage = getV1WorkerContainerImageForJob(options);
14731458
for (WorkerPool workerPool : newJob.getEnvironment().getWorkerPools()) {
14741459
workerPool.setWorkerHarnessContainerImage(containerImage);
14751460
}
@@ -2634,59 +2619,101 @@ public Map<PCollection<?>, ReplacementOutput> mapOutputs(
26342619
}
26352620

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

26402639
if (containerImage == null) {
26412640
// If not set, construct and return default image URL.
2642-
return getDefaultContainerImageUrl(options);
2641+
return getDefaultV2SdkHarnessContainerImageUrl(options);
26432642
} else if (containerImage.contains("IMAGE")) {
26442643
// Replace placeholder with default image name
2645-
return containerImage.replace("IMAGE", getDefaultContainerImageNameForJob(options));
2644+
return containerImage.replace("IMAGE", getDefaultV2SdkHarnessContainerImageNameForJob());
26462645
} else {
26472646
return containerImage;
26482647
}
26492648
}
26502649

2651-
/** Construct the default Dataflow container full image URL. */
2652-
static String getDefaultContainerImageUrl(DataflowPipelineOptions options) {
2650+
/** Construct the default Dataflow worker container full image URL. */
2651+
static String getDefaultV1WorkerContainerImageUrl(DataflowPipelineOptions options) {
26532652
DataflowRunnerInfo dataflowRunnerInfo = DataflowRunnerInfo.getDataflowRunnerInfo();
26542653
return String.format(
26552654
"%s/%s:%s",
26562655
dataflowRunnerInfo.getContainerImageBaseRepository(),
2657-
getDefaultContainerImageNameForJob(options),
2658-
getDefaultContainerVersion(options));
2656+
getDefaultV1WorkerContainerImageNameForJob(options),
2657+
getDefaultV1WorkerContainerVersion(options));
2658+
}
2659+
2660+
/** Construct the default Java SDK container full image URL. */
2661+
static String getDefaultV2SdkHarnessContainerImageUrl(DataflowPipelineOptions options) {
2662+
DataflowRunnerInfo dataflowRunnerInfo = DataflowRunnerInfo.getDataflowRunnerInfo();
2663+
return String.format(
2664+
"%s/%s:%s",
2665+
dataflowRunnerInfo.getContainerImageBaseRepository(),
2666+
getDefaultV2SdkHarnessContainerImageNameForJob(),
2667+
getDefaultV2SdkHarnessContainerVersion(options));
26592668
}
26602669

26612670
/**
2662-
* Construct the default Dataflow container image name based on pipeline type and Java version.
2671+
* Construct the default Dataflow V1 worker container image name based on pipeline type and Java
2672+
* version.
26632673
*/
2664-
static String getDefaultContainerImageNameForJob(DataflowPipelineOptions options) {
2674+
static String getDefaultV1WorkerContainerImageNameForJob(DataflowPipelineOptions options) {
26652675
Environments.JavaVersion javaVersion = Environments.getJavaVersion();
2666-
if (useUnifiedWorker(options)) {
2667-
return String.format("beam_%s_sdk", javaVersion.name());
2668-
} else if (options.isStreaming()) {
2676+
if (options.isStreaming()) {
26692677
return String.format("beam-%s-streaming", javaVersion.legacyName());
26702678
} else {
26712679
return String.format("beam-%s-batch", javaVersion.legacyName());
26722680
}
26732681
}
26742682

26752683
/**
2676-
* Construct the default Dataflow container image name based on pipeline type and Java version.
2684+
* Construct the default Java SDK container image name based on pipeline type and Java version,
2685+
* for use by Dataflow V2.
2686+
*/
2687+
static String getDefaultV2SdkHarnessContainerImageNameForJob() {
2688+
Environments.JavaVersion javaVersion = Environments.getJavaVersion();
2689+
return String.format("beam_%s_sdk", javaVersion.name());
2690+
}
2691+
2692+
/**
2693+
* Construct the default Dataflow V1 worker container image name based on pipeline type and Java
2694+
* version.
26772695
*/
2678-
static String getDefaultContainerVersion(DataflowPipelineOptions options) {
2696+
static String getDefaultV1WorkerContainerVersion(DataflowPipelineOptions options) {
26792697
DataflowRunnerInfo dataflowRunnerInfo = DataflowRunnerInfo.getDataflowRunnerInfo();
26802698
ReleaseInfo releaseInfo = ReleaseInfo.getReleaseInfo();
26812699
if (releaseInfo.isDevSdkVersion()) {
2682-
if (useUnifiedWorker(options)) {
2683-
return dataflowRunnerInfo.getFnApiDevContainerVersion();
2684-
}
26852700
return dataflowRunnerInfo.getLegacyDevContainerVersion();
26862701
}
26872702
return releaseInfo.getSdkVersion();
26882703
}
26892704

2705+
/**
2706+
* Construct the default Dataflow container image name based on pipeline type and Java version.
2707+
*/
2708+
static String getDefaultV2SdkHarnessContainerVersion(DataflowPipelineOptions options) {
2709+
DataflowRunnerInfo dataflowRunnerInfo = DataflowRunnerInfo.getDataflowRunnerInfo();
2710+
ReleaseInfo releaseInfo = ReleaseInfo.getReleaseInfo();
2711+
if (releaseInfo.isDevSdkVersion()) {
2712+
return dataflowRunnerInfo.getFnApiDevContainerVersion();
2713+
}
2714+
return releaseInfo.getSdkVersion();
2715+
}
2716+
26902717
static boolean useUnifiedWorker(DataflowPipelineOptions options) {
26912718
return hasExperiment(options, "beam_fn_api")
26922719
|| 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

0 commit comments

Comments
 (0)