arunpandianp commented on code in PR #38768:
URL: https://github.com/apache/beam/pull/38768#discussion_r3634725420


##########
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcCommitWorkStream.java:
##########
@@ -512,7 +528,34 @@ public boolean commitWorkItem(
         return false;
       }
 
-      PendingRequest request = new PendingRequest(computation, commitRequest, 
onDone);
+      PendingRequest request =
+          new PendingRequest(
+              computation,
+              commitRequest.getShardingKey(),
+              commitRequest.toByteString(),
+              StreamingCommitRequestChunk.CommitType.COMMIT_TYPE_SINGLE_KEY,
+              onDone);
+      add(idGenerator.incrementAndGet(), request);
+      return true;
+    }
+
+    @Override
+    public boolean commitMultiKeyWorkItem(
+        String computation,
+        Windmill.MultiKeyWorkItemCommitRequest commitRequest,
+        Consumer<CommitStatus> onDone) {
+      if (!canAccept(commitRequest.getSerializedSize() + 
computation.length())) {
+        return false;
+      }
+      Preconditions.checkArgument(commitRequest.getRequestsCount() > 0);

Review Comment:
   done.



##########
runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/Commit.java:
##########
@@ -32,20 +35,43 @@ public abstract class Commit {
   public static Commit create(
       WorkItemCommitRequest request, ComputationState computationState, Work 
work) {
     Preconditions.checkArgument(request.getSerializedSize() > 0);
-    return new AutoValue_Commit(request, computationState, work);
+    return new AutoValue_Commit(
+        Optional.of(request), computationState, Optional.empty(), 
ImmutableList.of(work));
+  }
+
+  public static Commit createMultiKey(
+      Windmill.MultiKeyWorkItemCommitRequest multiKeyRequest,
+      ComputationState computationState,
+      ImmutableList<Work> workBatch) {
+    Preconditions.checkArgument(!workBatch.isEmpty());
+    return new AutoValue_Commit(
+        Optional.empty(), computationState, Optional.of(multiKeyRequest), 
workBatch);
   }
 
   public final String computationId() {
     return computationState().getComputationId();
   }
 
-  public abstract WorkItemCommitRequest request();
+  public abstract Optional<WorkItemCommitRequest> singleKeyRequest();
 
   public abstract ComputationState computationState();
 
-  public abstract Work work();
+  public abstract Optional<Windmill.MultiKeyWorkItemCommitRequest> 
multiKeyRequest();
+
+  public abstract ImmutableList<Work> workBatch();
+
+  public final boolean isFailed() {
+    for (Work w : workBatch()) {
+      if (w.isFailed()) {
+        return true;
+      }
+    }
+    return false;
+  }
 
   public final int getSize() {

Review Comment:
   done.



##########
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java:
##########
@@ -282,11 +296,16 @@ public void 
testCommit_handlesCompleteCommits_commitStatusNotOK() {
     waitForExpectedSetSize(completeCommits, commits.size());
 
     for (Commit commit : commits) {
-      WorkItemCommitRequest request = 
committed.get(commit.work().getWorkItem().getWorkToken());
+      WorkItemCommitRequest request =

Review Comment:
   done



##########
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitterTest.java:
##########
@@ -141,12 +141,13 @@ public void testCommit() {
                 (CompleteCommit completeCommit, Commit commit) ->
                     
completeCommit.computationId().equals(commit.computationId())
                         && completeCommit.status() == Windmill.CommitStatus.OK
-                        && completeCommit.workId().equals(commit.work().id())
+                        && 
completeCommit.workId().equals(commit.workBatch().get(0).id())

Review Comment:
   done.



##########
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java:
##########
@@ -409,10 +436,14 @@ public void testMultipleCommitSendersSingleStream() {
     waitForExpectedSetSize(completeCommits, commits.size());
 
     for (Commit commit : commits) {
-      WorkItemCommitRequest request = 
committed.get(commit.work().getWorkItem().getWorkToken());
+      WorkItemCommitRequest request =

Review Comment:
   done.



##########
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java:
##########
@@ -474,4 +505,201 @@ public void testStop_drainsCommitQueue_concurrentCommit()
 
     waitForExpectedSetSize(completeCommits, sentCommits.intValue());
   }
+
+  @Test
+  public void testCommit_multiKeyCommitFailedWork() {
+    Set<CompleteCommit> completeCommits = Collections.newSetFromMap(new 
ConcurrentHashMap<>());
+    workCommitter = createWorkCommitter(completeCommits::add);
+
+    Work workA = createMockWork(101L);
+    Work workB = createMockWork(102L);
+    Work workC = createMockWork(103L);
+
+    // Mark non-primary key B as failed
+    workB.setFailed();
+
+    Windmill.MultiKeyWorkItemCommitRequest multiKeyRequest =
+        Windmill.MultiKeyWorkItemCommitRequest.newBuilder()
+            .addRequests(
+                Windmill.WorkItemCommitRequest.newBuilder()
+                    .setKey(workA.getWorkItem().getKey())
+                    .setShardingKey(workA.getWorkItem().getShardingKey())
+                    .setWorkToken(workA.getWorkItem().getWorkToken())
+                    .setCacheToken(workA.getWorkItem().getCacheToken())
+                    .build())
+            .addRequests(
+                Windmill.WorkItemCommitRequest.newBuilder()
+                    .setKey(workB.getWorkItem().getKey())
+                    .setShardingKey(workB.getWorkItem().getShardingKey())
+                    .setWorkToken(workB.getWorkItem().getWorkToken())
+                    .setCacheToken(workB.getWorkItem().getCacheToken())
+                    .build())
+            .addRequests(
+                Windmill.WorkItemCommitRequest.newBuilder()
+                    .setKey(workC.getWorkItem().getKey())
+                    .setShardingKey(workC.getWorkItem().getShardingKey())
+                    .setWorkToken(workC.getWorkItem().getWorkToken())
+                    .setCacheToken(workC.getWorkItem().getCacheToken())
+                    .build())
+            .build();
+
+    Commit commit =
+        Commit.createMultiKey(
+            multiKeyRequest,
+            createComputationState("computationId"),
+            ImmutableList.of(workA, workB, workC));
+
+    workCommitter.start();
+    workCommitter.commit(commit);
+
+    // The entire batch must be aborted immediately without making network 
calls
+    waitForExpectedSetSize(completeCommits, 3);
+
+    // Verify all three works are aborted individually
+    assertThat(completeCommits)
+        .containsExactly(
+            CompleteCommit.create(
+                "computationId", workA.getShardedKey(), workA.id(), 
CommitStatus.ABORTED),
+            CompleteCommit.create(
+                "computationId", workB.getShardedKey(), workB.id(), 
CommitStatus.ABORTED),
+            CompleteCommit.create(
+                "computationId", workC.getShardedKey(), workC.id(), 
CommitStatus.ABORTED));
+
+    workCommitter.stop();
+  }
+
+  @Test
+  public void testCommit_multiKeyCommitSuccess() {

Review Comment:
   done.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to