Skip to content

Commit 51b2ad6

Browse files
authored
Add drain states to PipelineResult (#39020)
* Add drain states to PipelineResult * Fix draining job lookup for Dataflow updates
1 parent 55e1ecb commit 51b2ad6

15 files changed

Lines changed: 123 additions & 23 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/issues/39020)).
7879

7980
## Deprecations
8081

runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkDetachedRunnerResult.java

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -120,11 +120,11 @@ public synchronized State drain() throws IOException {
120120

121121
private State getDrainState(CompletableFuture<String> drainFuture) throws IOException {
122122
if (!drainFuture.isDone()) {
123-
return State.RUNNING;
123+
return State.DRAINING;
124124
}
125125
try {
126126
drainFuture.get();
127-
return State.DONE;
127+
return State.DRAINED;
128128
} catch (InterruptedException e) {
129129
Thread.currentThread().interrupt();
130130
throw new IOException("Failed to drain Flink job", e);

runners/flink/src/test/java/org/apache/beam/runners/flink/FlinkRunnerResultTest.java

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -66,18 +66,18 @@ public void testDrainDoneResultDoesNotThrowAnException() throws Exception {
6666
}
6767

6868
@Test
69-
public void testDetachedDrainReturnsRunningThenDone() throws Exception {
69+
public void testDetachedDrainReturnsDrainingThenDrained() throws Exception {
7070
JobClient jobClient = mock(JobClient.class);
7171
CompletableFuture<String> drainFuture = new CompletableFuture<>();
7272
when(jobClient.stopWithSavepoint(true, null, SavepointFormatType.DEFAULT))
7373
.thenReturn(drainFuture);
7474
FlinkDetachedRunnerResult result = new FlinkDetachedRunnerResult(jobClient, 1);
7575

76-
assertThat(result.drain(), is(PipelineResult.State.RUNNING));
77-
assertThat(result.getState(), is(PipelineResult.State.RUNNING));
76+
assertThat(result.drain(), is(PipelineResult.State.DRAINING));
77+
assertThat(result.getState(), is(PipelineResult.State.DRAINING));
7878

7979
drainFuture.complete("savepoint");
80-
assertThat(result.getState(), is(PipelineResult.State.DONE));
80+
assertThat(result.getState(), is(PipelineResult.State.DRAINED));
8181
verify(jobClient).stopWithSavepoint(true, null, SavepointFormatType.DEFAULT);
8282
}
8383

@@ -132,11 +132,11 @@ public void testDetachedDrainRetriesAfterFailure() throws Exception {
132132
result.drain();
133133
fail("Expected IOException");
134134
} catch (IOException expected) {
135-
assertThat(result.drain(), is(PipelineResult.State.RUNNING));
135+
assertThat(result.drain(), is(PipelineResult.State.DRAINING));
136136
}
137137

138138
retryDrainFuture.complete("savepoint");
139-
assertThat(result.getState(), is(PipelineResult.State.DONE));
139+
assertThat(result.getState(), is(PipelineResult.State.DRAINED));
140140
verify(jobClient, times(2)).stopWithSavepoint(true, null, SavepointFormatType.DEFAULT);
141141
}
142142
}

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

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -365,6 +365,7 @@ private void logTerminalState(State state) {
365365
switch (state) {
366366
case DONE:
367367
case CANCELLED:
368+
case DRAINED:
368369
LOG.info("Job {} finished with status {}.", getJobId(), state);
369370
break;
370371
case UPDATED:

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/main/java/org/apache/beam/runners/dataflow/util/MonitoringUtil.java

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -221,17 +221,19 @@ public static State toState(@Nullable String stateName) {
221221
return State.CANCELLED;
222222
case "JOB_STATE_UPDATED":
223223
return State.UPDATED;
224+
case "JOB_STATE_DRAINING":
225+
return State.DRAINING;
226+
case "JOB_STATE_DRAINED":
227+
return State.DRAINED;
224228

225229
case "JOB_STATE_RUNNING":
226230
case "JOB_STATE_PENDING": // Job has not yet started; closest mapping is RUNNING
227-
case "JOB_STATE_DRAINING": // Job is still active; the closest mapping is RUNNING
228231
case "JOB_STATE_CANCELLING": // Job is still active; the closest mapping is RUNNING
229232
case "JOB_STATE_PAUSING": // Job is still active; the closest mapping is RUNNING
230233
case "JOB_STATE_RESOURCE_CLEANING_UP": // Job is still active; the closest mapping is RUNNING
231234
return State.RUNNING;
232235

233236
case "JOB_STATE_DONE":
234-
case "JOB_STATE_DRAINED": // Job has successfully terminated; closest mapping is DONE
235237
return State.DONE;
236238
default:
237239
LOG.warn(

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -417,7 +417,7 @@ public void testDrainUnterminatedJobThatSucceeds() throws IOException {
417417
DataflowPipelineJob job =
418418
new DataflowPipelineJob(DataflowClient.create(options), JOB_ID, options, null);
419419

420-
assertEquals(State.RUNNING, job.drain());
420+
assertEquals(State.DRAINING, job.drain());
421421
Job content = new Job();
422422
content.setProjectId(PROJECT_ID);
423423
content.setId(JOB_ID);

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();

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -100,9 +100,9 @@ public void testToStateNormal() {
100100

101101
// Non-trivially mapped cases
102102
assertEquals(State.STOPPED, MonitoringUtil.toState("JOB_STATE_PAUSED"));
103-
assertEquals(State.RUNNING, MonitoringUtil.toState("JOB_STATE_DRAINING"));
103+
assertEquals(State.DRAINING, MonitoringUtil.toState("JOB_STATE_DRAINING"));
104104
assertEquals(State.RUNNING, MonitoringUtil.toState("JOB_STATE_PAUSING"));
105-
assertEquals(State.DONE, MonitoringUtil.toState("JOB_STATE_DRAINED"));
105+
assertEquals(State.DRAINED, MonitoringUtil.toState("JOB_STATE_DRAINED"));
106106
}
107107

108108
@Test

runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/JobInvocation.java

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -116,6 +116,12 @@ public void onSuccess(PortablePipelineResult pipelineResult) {
116116
case RUNNING:
117117
setState(JobState.Enum.RUNNING);
118118
break;
119+
case DRAINING:
120+
setState(JobState.Enum.DRAINING);
121+
break;
122+
case DRAINED:
123+
setState(JobState.Enum.DRAINED);
124+
break;
119125
case CANCELLED:
120126
setState(JobState.Enum.CANCELLED);
121127
break;
@@ -169,9 +175,12 @@ public synchronized void cancel() {
169175
new FutureCallback<PortablePipelineResult>() {
170176
@Override
171177
public void onSuccess(PortablePipelineResult pipelineResult) {
172-
// Do not cancel when we are already done.
173-
if (pipelineResult != null
174-
&& pipelineResult.getState() != PipelineResult.State.DONE) {
178+
// Do not cancel when the runner has already successfully finished.
179+
if (pipelineResult != null) {
180+
PipelineResult.State state = pipelineResult.getState();
181+
if (state == PipelineResult.State.DONE || state == PipelineResult.State.DRAINED) {
182+
return;
183+
}
175184
try {
176185
pipelineResult.cancel();
177186
setState(JobState.Enum.CANCELLED);

0 commit comments

Comments
 (0)