Skip to content

Commit 8def3a7

Browse files
committed
Fix draining job lookup for Dataflow updates
1 parent 7231e06 commit 8def3a7

3 files changed

Lines changed: 27 additions & 3 deletions

File tree

CHANGES.md

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,7 @@
7575
## Breaking Changes
7676

7777
* (Python) Removed `google-perftools` from the SDK container images. Users who wish to use `--profiler_agent=tcmalloc` should install google-perftools APT package in their custom container images separately ([#39323](https://github.com/apache/beam/issues/39323)).
78+
* (Java) Added `DRAINING` and `DRAINED` states to `PipelineResult`, including runner state mappings and Dataflow update handling ([#39020](https://github.com/apache/beam/pull/39020)).
7879

7980
## Deprecations
8081

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

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2535,8 +2535,9 @@ private String getJobIdFromName(String jobName) {
25352535
listResult = dataflowClient.listJobs(token);
25362536
token = listResult.getNextPageToken();
25372537
for (Job job : listResult.getJobs()) {
2538+
State state = MonitoringUtil.toState(job.getCurrentState());
25382539
if (job.getName().equals(jobName)
2539-
&& MonitoringUtil.toState(job.getCurrentState()).equals(State.RUNNING)) {
2540+
&& (state.equals(State.RUNNING) || state.equals(State.DRAINING))) {
25402541
return job.getId();
25412542
}
25422543
}

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

Lines changed: 24 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -285,6 +285,11 @@ private static Pipeline buildDataflowPipelineWithLargeGraph(DataflowPipelineOpti
285285
}
286286

287287
static Dataflow buildMockDataflow(Dataflow.Projects.Locations.Jobs mockJobs) throws IOException {
288+
return buildMockDataflow(mockJobs, "JOB_STATE_RUNNING");
289+
}
290+
291+
static Dataflow buildMockDataflow(Dataflow.Projects.Locations.Jobs mockJobs, String currentState)
292+
throws IOException {
288293
Dataflow mockDataflowClient = mock(Dataflow.class);
289294
Dataflow.Projects mockProjects = mock(Dataflow.Projects.class);
290295
Dataflow.Projects.Locations mockLocations = mock(Dataflow.Projects.Locations.class);
@@ -308,7 +313,7 @@ static Dataflow buildMockDataflow(Dataflow.Projects.Locations.Jobs mockJobs) thr
308313
new Job()
309314
.setName("oldjobname")
310315
.setId("oldJobId")
311-
.setCurrentState("JOB_STATE_RUNNING"))));
316+
.setCurrentState(currentState))));
312317

313318
Job resultJob = new Job();
314319
resultJob.setId("newid");
@@ -375,14 +380,18 @@ static GcsUtil buildMockGcsUtil() throws IOException {
375380
}
376381

377382
private DataflowPipelineOptions buildPipelineOptions() throws IOException {
383+
return buildPipelineOptions("JOB_STATE_RUNNING");
384+
}
385+
386+
private DataflowPipelineOptions buildPipelineOptions(String currentState) throws IOException {
378387
DataflowPipelineOptions options = PipelineOptionsFactory.as(DataflowPipelineOptions.class);
379388
options.setRunner(DataflowRunner.class);
380389
options.setProject(PROJECT_ID);
381390
options.setTempLocation(VALID_TEMP_BUCKET);
382391
options.setRegion(REGION_ID);
383392
// Set FILES_PROPERTY to empty to prevent a default value calculated from classpath.
384393
options.setFilesToStage(new ArrayList<>());
385-
options.setDataflowClient(buildMockDataflow(mockJobs));
394+
options.setDataflowClient(buildMockDataflow(mockJobs, currentState));
386395
options.setGcsUtil(mockGcsUtil);
387396
options.setGcpCredential(new TestCredential());
388397

@@ -793,6 +802,19 @@ public void testUpdate() throws IOException {
793802
assertValidJob(jobCaptor.getValue());
794803
}
795804

805+
@Test
806+
public void testUpdateDrainingJob() throws IOException {
807+
DataflowPipelineOptions options = buildPipelineOptions("JOB_STATE_DRAINING");
808+
options.setUpdate(true);
809+
options.setJobName("oldJobName");
810+
Pipeline p = buildDataflowPipeline(options);
811+
p.run();
812+
813+
ArgumentCaptor<Job> jobCaptor = ArgumentCaptor.forClass(Job.class);
814+
Mockito.verify(mockJobs).create(eq(PROJECT_ID), eq(REGION_ID), jobCaptor.capture());
815+
assertEquals("oldJobId", jobCaptor.getValue().getReplaceJobId());
816+
}
817+
796818
@Test
797819
public void testUploadGraph() throws IOException {
798820
DataflowPipelineOptions options = buildPipelineOptions();

0 commit comments

Comments
 (0)