This is an automated email from the ASF dual-hosted git repository.

stankiewicz pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new 51b2ad6bc08 Add drain states to PipelineResult (#39020)
51b2ad6bc08 is described below

commit 51b2ad6bc084a65cd5af24d267b9c5c863d1aa52
Author: Lalit Yadav <[email protected]>
AuthorDate: Wed Jul 22 03:24:01 2026 -0500

    Add drain states to PipelineResult (#39020)
    
    * Add drain states to PipelineResult
    
    * Fix draining job lookup for Dataflow updates
---
 CHANGES.md                                         |  1 +
 .../runners/flink/FlinkDetachedRunnerResult.java   |  4 +--
 .../beam/runners/flink/FlinkRunnerResultTest.java  | 12 ++++----
 .../beam/runners/dataflow/DataflowPipelineJob.java |  1 +
 .../beam/runners/dataflow/DataflowRunner.java      |  3 +-
 .../beam/runners/dataflow/util/MonitoringUtil.java |  6 ++--
 .../runners/dataflow/DataflowPipelineJobTest.java  |  2 +-
 .../beam/runners/dataflow/DataflowRunnerTest.java  | 26 ++++++++++++++--
 .../runners/dataflow/util/MonitoringUtilTest.java  |  4 +--
 .../beam/runners/jobsubmission/JobInvocation.java  | 15 +++++++--
 .../runners/jobsubmission/JobInvocationTest.java   | 36 ++++++++++++++++++++++
 .../portability/JobServicePipelineResult.java      |  7 ++---
 .../runners/portability/PortableRunnerTest.java    | 21 +++++++++++++
 .../java/org/apache/beam/sdk/PipelineResult.java   |  6 ++++
 .../apache/beam/sdk/nexmark/NexmarkLauncher.java   |  2 ++
 15 files changed, 123 insertions(+), 23 deletions(-)

diff --git a/CHANGES.md b/CHANGES.md
index 991c341460a..87fdaa8d021 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -75,6 +75,7 @@
 ## Breaking Changes
 
 * (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)).
+* (Java) Added `DRAINING` and `DRAINED` states to `PipelineResult`, including 
runner state mappings and Dataflow update handling 
([#39020](https://github.com/apache/beam/issues/39020)).
 
 ## Deprecations
 
diff --git 
a/runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkDetachedRunnerResult.java
 
b/runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkDetachedRunnerResult.java
index b26e865526d..d7c8912ded0 100644
--- 
a/runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkDetachedRunnerResult.java
+++ 
b/runners/flink/src/main/java/org/apache/beam/runners/flink/FlinkDetachedRunnerResult.java
@@ -120,11 +120,11 @@ public class FlinkDetachedRunnerResult implements 
PipelineResult {
 
   private State getDrainState(CompletableFuture<String> drainFuture) throws 
IOException {
     if (!drainFuture.isDone()) {
-      return State.RUNNING;
+      return State.DRAINING;
     }
     try {
       drainFuture.get();
-      return State.DONE;
+      return State.DRAINED;
     } catch (InterruptedException e) {
       Thread.currentThread().interrupt();
       throw new IOException("Failed to drain Flink job", e);
diff --git 
a/runners/flink/src/test/java/org/apache/beam/runners/flink/FlinkRunnerResultTest.java
 
b/runners/flink/src/test/java/org/apache/beam/runners/flink/FlinkRunnerResultTest.java
index 908d940f5ef..09ceede9d58 100644
--- 
a/runners/flink/src/test/java/org/apache/beam/runners/flink/FlinkRunnerResultTest.java
+++ 
b/runners/flink/src/test/java/org/apache/beam/runners/flink/FlinkRunnerResultTest.java
@@ -66,18 +66,18 @@ public class FlinkRunnerResultTest {
   }
 
   @Test
-  public void testDetachedDrainReturnsRunningThenDone() throws Exception {
+  public void testDetachedDrainReturnsDrainingThenDrained() throws Exception {
     JobClient jobClient = mock(JobClient.class);
     CompletableFuture<String> drainFuture = new CompletableFuture<>();
     when(jobClient.stopWithSavepoint(true, null, SavepointFormatType.DEFAULT))
         .thenReturn(drainFuture);
     FlinkDetachedRunnerResult result = new 
FlinkDetachedRunnerResult(jobClient, 1);
 
-    assertThat(result.drain(), is(PipelineResult.State.RUNNING));
-    assertThat(result.getState(), is(PipelineResult.State.RUNNING));
+    assertThat(result.drain(), is(PipelineResult.State.DRAINING));
+    assertThat(result.getState(), is(PipelineResult.State.DRAINING));
 
     drainFuture.complete("savepoint");
-    assertThat(result.getState(), is(PipelineResult.State.DONE));
+    assertThat(result.getState(), is(PipelineResult.State.DRAINED));
     verify(jobClient).stopWithSavepoint(true, null, 
SavepointFormatType.DEFAULT);
   }
 
@@ -132,11 +132,11 @@ public class FlinkRunnerResultTest {
       result.drain();
       fail("Expected IOException");
     } catch (IOException expected) {
-      assertThat(result.drain(), is(PipelineResult.State.RUNNING));
+      assertThat(result.drain(), is(PipelineResult.State.DRAINING));
     }
 
     retryDrainFuture.complete("savepoint");
-    assertThat(result.getState(), is(PipelineResult.State.DONE));
+    assertThat(result.getState(), is(PipelineResult.State.DRAINED));
     verify(jobClient, times(2)).stopWithSavepoint(true, null, 
SavepointFormatType.DEFAULT);
   }
 }
diff --git 
a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineJob.java
 
b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineJob.java
index 0d7e5eaf68d..278dace6ef4 100644
--- 
a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineJob.java
+++ 
b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowPipelineJob.java
@@ -365,6 +365,7 @@ public class DataflowPipelineJob implements PipelineResult {
     switch (state) {
       case DONE:
       case CANCELLED:
+      case DRAINED:
         LOG.info("Job {} finished with status {}.", getJobId(), state);
         break;
       case UPDATED:
diff --git 
a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java
 
b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java
index 011f60f4fd1..1b03cc10351 100644
--- 
a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java
+++ 
b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/DataflowRunner.java
@@ -2535,8 +2535,9 @@ public class DataflowRunner extends 
PipelineRunner<DataflowPipelineJob> {
         listResult = dataflowClient.listJobs(token);
         token = listResult.getNextPageToken();
         for (Job job : listResult.getJobs()) {
+          State state = MonitoringUtil.toState(job.getCurrentState());
           if (job.getName().equals(jobName)
-              && 
MonitoringUtil.toState(job.getCurrentState()).equals(State.RUNNING)) {
+              && (state.equals(State.RUNNING) || 
state.equals(State.DRAINING))) {
             return job.getId();
           }
         }
diff --git 
a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/MonitoringUtil.java
 
b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/MonitoringUtil.java
index d117cf78612..0e25bbd265a 100644
--- 
a/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/MonitoringUtil.java
+++ 
b/runners/google-cloud-dataflow-java/src/main/java/org/apache/beam/runners/dataflow/util/MonitoringUtil.java
@@ -221,17 +221,19 @@ public class MonitoringUtil {
         return State.CANCELLED;
       case "JOB_STATE_UPDATED":
         return State.UPDATED;
+      case "JOB_STATE_DRAINING":
+        return State.DRAINING;
+      case "JOB_STATE_DRAINED":
+        return State.DRAINED;
 
       case "JOB_STATE_RUNNING":
       case "JOB_STATE_PENDING": // Job has not yet started; closest mapping is 
RUNNING
-      case "JOB_STATE_DRAINING": // Job is still active; the closest mapping 
is RUNNING
       case "JOB_STATE_CANCELLING": // Job is still active; the closest mapping 
is RUNNING
       case "JOB_STATE_PAUSING": // Job is still active; the closest mapping is 
RUNNING
       case "JOB_STATE_RESOURCE_CLEANING_UP": // Job is still active; the 
closest mapping is RUNNING
         return State.RUNNING;
 
       case "JOB_STATE_DONE":
-      case "JOB_STATE_DRAINED": // Job has successfully terminated; closest 
mapping is DONE
         return State.DONE;
       default:
         LOG.warn(
diff --git 
a/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineJobTest.java
 
b/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineJobTest.java
index 4b088eb41a7..be7ad2e6e11 100644
--- 
a/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineJobTest.java
+++ 
b/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowPipelineJobTest.java
@@ -417,7 +417,7 @@ public class DataflowPipelineJobTest {
     DataflowPipelineJob job =
         new DataflowPipelineJob(DataflowClient.create(options), JOB_ID, 
options, null);
 
-    assertEquals(State.RUNNING, job.drain());
+    assertEquals(State.DRAINING, job.drain());
     Job content = new Job();
     content.setProjectId(PROJECT_ID);
     content.setId(JOB_ID);
diff --git 
a/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowRunnerTest.java
 
b/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowRunnerTest.java
index 073b30f928d..6c32dd749d0 100644
--- 
a/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowRunnerTest.java
+++ 
b/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/DataflowRunnerTest.java
@@ -285,6 +285,11 @@ public class DataflowRunnerTest implements Serializable {
   }
 
   static Dataflow buildMockDataflow(Dataflow.Projects.Locations.Jobs mockJobs) 
throws IOException {
+    return buildMockDataflow(mockJobs, "JOB_STATE_RUNNING");
+  }
+
+  static Dataflow buildMockDataflow(Dataflow.Projects.Locations.Jobs mockJobs, 
String currentState)
+      throws IOException {
     Dataflow mockDataflowClient = mock(Dataflow.class);
     Dataflow.Projects mockProjects = mock(Dataflow.Projects.class);
     Dataflow.Projects.Locations mockLocations = 
mock(Dataflow.Projects.Locations.class);
@@ -308,7 +313,7 @@ public class DataflowRunnerTest implements Serializable {
                         new Job()
                             .setName("oldjobname")
                             .setId("oldJobId")
-                            .setCurrentState("JOB_STATE_RUNNING"))));
+                            .setCurrentState(currentState))));
 
     Job resultJob = new Job();
     resultJob.setId("newid");
@@ -375,6 +380,10 @@ public class DataflowRunnerTest implements Serializable {
   }
 
   private DataflowPipelineOptions buildPipelineOptions() throws IOException {
+    return buildPipelineOptions("JOB_STATE_RUNNING");
+  }
+
+  private DataflowPipelineOptions buildPipelineOptions(String currentState) 
throws IOException {
     DataflowPipelineOptions options = 
PipelineOptionsFactory.as(DataflowPipelineOptions.class);
     options.setRunner(DataflowRunner.class);
     options.setProject(PROJECT_ID);
@@ -382,7 +391,7 @@ public class DataflowRunnerTest implements Serializable {
     options.setRegion(REGION_ID);
     // Set FILES_PROPERTY to empty to prevent a default value calculated from 
classpath.
     options.setFilesToStage(new ArrayList<>());
-    options.setDataflowClient(buildMockDataflow(mockJobs));
+    options.setDataflowClient(buildMockDataflow(mockJobs, currentState));
     options.setGcsUtil(mockGcsUtil);
     options.setGcpCredential(new TestCredential());
 
@@ -793,6 +802,19 @@ public class DataflowRunnerTest implements Serializable {
     assertValidJob(jobCaptor.getValue());
   }
 
+  @Test
+  public void testUpdateDrainingJob() throws IOException {
+    DataflowPipelineOptions options = 
buildPipelineOptions("JOB_STATE_DRAINING");
+    options.setUpdate(true);
+    options.setJobName("oldJobName");
+    Pipeline p = buildDataflowPipeline(options);
+    p.run();
+
+    ArgumentCaptor<Job> jobCaptor = ArgumentCaptor.forClass(Job.class);
+    Mockito.verify(mockJobs).create(eq(PROJECT_ID), eq(REGION_ID), 
jobCaptor.capture());
+    assertEquals("oldJobId", jobCaptor.getValue().getReplaceJobId());
+  }
+
   @Test
   public void testUploadGraph() throws IOException {
     DataflowPipelineOptions options = buildPipelineOptions();
diff --git 
a/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/util/MonitoringUtilTest.java
 
b/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/util/MonitoringUtilTest.java
index 5f76b6750ff..35af1be4cdf 100644
--- 
a/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/util/MonitoringUtilTest.java
+++ 
b/runners/google-cloud-dataflow-java/src/test/java/org/apache/beam/runners/dataflow/util/MonitoringUtilTest.java
@@ -100,9 +100,9 @@ public class MonitoringUtilTest {
 
     // Non-trivially mapped cases
     assertEquals(State.STOPPED, MonitoringUtil.toState("JOB_STATE_PAUSED"));
-    assertEquals(State.RUNNING, MonitoringUtil.toState("JOB_STATE_DRAINING"));
+    assertEquals(State.DRAINING, MonitoringUtil.toState("JOB_STATE_DRAINING"));
     assertEquals(State.RUNNING, MonitoringUtil.toState("JOB_STATE_PAUSING"));
-    assertEquals(State.DONE, MonitoringUtil.toState("JOB_STATE_DRAINED"));
+    assertEquals(State.DRAINED, MonitoringUtil.toState("JOB_STATE_DRAINED"));
   }
 
   @Test
diff --git 
a/runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/JobInvocation.java
 
b/runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/JobInvocation.java
index 9da3bf38e57..3bd0a1311b0 100644
--- 
a/runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/JobInvocation.java
+++ 
b/runners/java-job-service/src/main/java/org/apache/beam/runners/jobsubmission/JobInvocation.java
@@ -116,6 +116,12 @@ public class JobInvocation {
                 case RUNNING:
                   setState(JobState.Enum.RUNNING);
                   break;
+                case DRAINING:
+                  setState(JobState.Enum.DRAINING);
+                  break;
+                case DRAINED:
+                  setState(JobState.Enum.DRAINED);
+                  break;
                 case CANCELLED:
                   setState(JobState.Enum.CANCELLED);
                   break;
@@ -169,9 +175,12 @@ public class JobInvocation {
           new FutureCallback<PortablePipelineResult>() {
             @Override
             public void onSuccess(PortablePipelineResult pipelineResult) {
-              // Do not cancel when we are already done.
-              if (pipelineResult != null
-                  && pipelineResult.getState() != PipelineResult.State.DONE) {
+              // Do not cancel when the runner has already successfully 
finished.
+              if (pipelineResult != null) {
+                PipelineResult.State state = pipelineResult.getState();
+                if (state == PipelineResult.State.DONE || state == 
PipelineResult.State.DRAINED) {
+                  return;
+                }
                 try {
                   pipelineResult.cancel();
                   setState(JobState.Enum.CANCELLED);
diff --git 
a/runners/java-job-service/src/test/java/org/apache/beam/runners/jobsubmission/JobInvocationTest.java
 
b/runners/java-job-service/src/test/java/org/apache/beam/runners/jobsubmission/JobInvocationTest.java
index 31380ace43c..7100e57717b 100644
--- 
a/runners/java-job-service/src/test/java/org/apache/beam/runners/jobsubmission/JobInvocationTest.java
+++ 
b/runners/java-job-service/src/test/java/org/apache/beam/runners/jobsubmission/JobInvocationTest.java
@@ -80,6 +80,28 @@ public class JobInvocationTest {
     awaitJobState(jobInvocation, JobApi.JobState.Enum.DONE);
   }
 
+  @Test(timeout = 10_000)
+  public void testStateAfterDrainCompleted() throws Exception {
+    jobInvocation.start();
+    assertThat(jobInvocation.getState(), is(JobApi.JobState.Enum.RUNNING));
+
+    TestPipelineResult pipelineResult = new 
TestPipelineResult(PipelineResult.State.DRAINED);
+    runner.setResult(pipelineResult);
+
+    awaitJobState(jobInvocation, JobApi.JobState.Enum.DRAINED);
+  }
+
+  @Test(timeout = 10_000)
+  public void testStateAfterDrainStarted() throws Exception {
+    jobInvocation.start();
+    assertThat(jobInvocation.getState(), is(JobApi.JobState.Enum.RUNNING));
+
+    TestPipelineResult pipelineResult = new 
TestPipelineResult(PipelineResult.State.DRAINING);
+    runner.setResult(pipelineResult);
+
+    awaitJobState(jobInvocation, JobApi.JobState.Enum.DRAINING);
+  }
+
   @Test(timeout = 10_000)
   public void testStateAfterCompletionWithoutResult() throws Exception {
     jobInvocation.start();
@@ -128,6 +150,20 @@ public class JobInvocationTest {
     assertThat(pipelineResult.cancelLatch.getCount(), is(1L));
   }
 
+  @Test(timeout = 10_000)
+  public void testNoCancellationWhenDrained() throws Exception {
+    jobInvocation.start();
+    assertThat(jobInvocation.getState(), is(JobApi.JobState.Enum.RUNNING));
+
+    TestPipelineResult pipelineResult = new 
TestPipelineResult(PipelineResult.State.DRAINED);
+    runner.setResult(pipelineResult);
+    awaitJobState(jobInvocation, JobApi.JobState.Enum.DRAINED);
+
+    jobInvocation.cancel();
+    assertThat(jobInvocation.getState(), is(JobApi.JobState.Enum.DRAINED));
+    assertThat(pipelineResult.cancelLatch.getCount(), is(1L));
+  }
+
   @Test(timeout = 10_000)
   public void testReturnsMetricsFromJobInvocationAfterSuccess() throws 
Exception {
     JobApi.MetricResults expectedMonitoringInfos = 
JobApi.MetricResults.newBuilder().build();
diff --git 
a/runners/portability/java/src/main/java/org/apache/beam/runners/portability/JobServicePipelineResult.java
 
b/runners/portability/java/src/main/java/org/apache/beam/runners/portability/JobServicePipelineResult.java
index 0f3d3b6d191..7c72a1039c0 100644
--- 
a/runners/portability/java/src/main/java/org/apache/beam/runners/portability/JobServicePipelineResult.java
+++ 
b/runners/portability/java/src/main/java/org/apache/beam/runners/portability/JobServicePipelineResult.java
@@ -160,7 +160,7 @@ class JobServicePipelineResult implements PipelineResult, 
AutoCloseable {
   }
 
   private void propagateErrors() {
-    if (terminalState != State.DONE) {
+    if (terminalState != State.DONE && terminalState != State.DRAINED) {
       JobMessagesRequest messageStreamRequest =
           JobMessagesRequest.newBuilder().setJobIdBytes(jobId).build();
       Iterator<JobMessagesResponse> messageStreamIterator =
@@ -196,10 +196,9 @@ class JobServicePipelineResult implements PipelineResult, 
AutoCloseable {
       case UPDATED:
         return State.UPDATED;
       case DRAINING:
-        // TODO: Determine the correct mappings for the states below.
-        return State.UNKNOWN;
+        return State.DRAINING;
       case DRAINED:
-        return State.UNKNOWN;
+        return State.DRAINED;
       case STARTING:
         return State.RUNNING;
       case CANCELLING:
diff --git 
a/runners/portability/java/src/test/java/org/apache/beam/runners/portability/PortableRunnerTest.java
 
b/runners/portability/java/src/test/java/org/apache/beam/runners/portability/PortableRunnerTest.java
index 68f3f6eae39..068d70b0d25 100644
--- 
a/runners/portability/java/src/test/java/org/apache/beam/runners/portability/PortableRunnerTest.java
+++ 
b/runners/portability/java/src/test/java/org/apache/beam/runners/portability/PortableRunnerTest.java
@@ -109,6 +109,27 @@ public class PortableRunnerTest implements Serializable {
     assertThat(state, is(State.DONE));
   }
 
+  @Test
+  public void mapsDrainingJobState() throws Exception {
+    createJobServer(JobState.Enum.DRAINING, 
JobApi.MetricResults.getDefaultInstance());
+    PortableRunner runner = PortableRunner.create(options, 
ManagedChannelFactory.createInProcess());
+    PipelineResult result = runner.run(p);
+
+    try {
+      assertThat(result.getState(), is(State.DRAINING));
+    } finally {
+      ((AutoCloseable) result).close();
+    }
+  }
+
+  @Test
+  public void mapsDrainedJobState() throws Exception {
+    createJobServer(JobState.Enum.DRAINED, 
JobApi.MetricResults.getDefaultInstance());
+    PortableRunner runner = PortableRunner.create(options, 
ManagedChannelFactory.createInProcess());
+    State state = runner.run(p).waitUntilFinish();
+    assertThat(state, is(State.DRAINED));
+  }
+
   @Test
   public void extractsMetrics() throws Exception {
     JobApi.MetricResults metricResults = generateMetricResults();
diff --git 
a/sdks/java/core/src/main/java/org/apache/beam/sdk/PipelineResult.java 
b/sdks/java/core/src/main/java/org/apache/beam/sdk/PipelineResult.java
index 46cca7833e5..eff031157e1 100644
--- a/sdks/java/core/src/main/java/org/apache/beam/sdk/PipelineResult.java
+++ b/sdks/java/core/src/main/java/org/apache/beam/sdk/PipelineResult.java
@@ -107,6 +107,12 @@ public interface PipelineResult {
     /** The job has been updated. */
     UPDATED(true, true),
 
+    /** The job is draining: no longer accepting new input while finishing 
in-flight work. */
+    DRAINING(false, false),
+
+    /** The job has finished draining. */
+    DRAINED(true, false),
+
     /** The job state reported by a runner cannot be interpreted by the SDK. */
     UNRECOGNIZED(false, false);
 
diff --git 
a/sdks/java/testing/nexmark/src/main/java/org/apache/beam/sdk/nexmark/NexmarkLauncher.java
 
b/sdks/java/testing/nexmark/src/main/java/org/apache/beam/sdk/nexmark/NexmarkLauncher.java
index 7e4e5da0d85..aaf63705ad2 100644
--- 
a/sdks/java/testing/nexmark/src/main/java/org/apache/beam/sdk/nexmark/NexmarkLauncher.java
+++ 
b/sdks/java/testing/nexmark/src/main/java/org/apache/beam/sdk/nexmark/NexmarkLauncher.java
@@ -515,9 +515,11 @@ public class NexmarkLauncher<OptionT extends 
NexmarkOptions> {
         case UNRECOGNIZED:
         case STOPPED:
         case RUNNING:
+        case DRAINING:
           // Keep going.
           break;
         case DONE:
+        case DRAINED:
           // All done.
           running = false;
           break;

Reply via email to