arunpandianp commented on code in PR #38768:
URL: https://github.com/apache/beam/pull/38768#discussion_r3634686426
##########
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();
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() {
+ Set<CompleteCommit> completeCommits = Collections.newSetFromMap(new
ConcurrentHashMap<>());
+ workCommitter = createWorkCommitter(completeCommits::add);
+
+ Work workA = createMockWork(101L);
+ Work workB = createMockWork(102L);
+ Work workC = createMockWork(103L);
+
+ 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);
+
+ // Wait for the server to receive and process the commits
+ fakeWindmillServer.waitForAndGetCommits(3);
+ waitForExpectedSetSize(completeCommits, 3);
+
+ // Verify that FakeWindmillServer received all 3 work requests in
multiKeyCommitsReceived
+ List<Windmill.MultiKeyWorkItemCommitRequest> multiKeyCommits =
+ fakeWindmillServer.getMultiKeyCommitsReceived();
+ assertThat(multiKeyCommits).hasSize(1);
+ assertThat(multiKeyCommits.get(0)).isEqualTo(multiKeyRequest);
+
+ // Verify all three works are completed successfully
+ assertThat(completeCommits)
+ .containsExactly(
+ CompleteCommit.create(
+ "computationId", workA.getShardedKey(), workA.id(),
CommitStatus.OK),
+ CompleteCommit.create(
+ "computationId", workB.getShardedKey(), workB.id(),
CommitStatus.OK),
+ CompleteCommit.create(
+ "computationId", workC.getShardedKey(), workC.id(),
CommitStatus.OK));
+
+ workCommitter.stop();
+ }
+
+ @Test
+ public void testCommit_multiKeyCommitStatusNotOK() {
+ Set<CompleteCommit> completeCommits = Collections.newSetFromMap(new
ConcurrentHashMap<>());
+ workCommitter = createWorkCommitter(completeCommits::add);
+
+ Work workA = createMockWork(101L);
+ Work workB = createMockWork(102L);
+ Work workC = createMockWork(103L);
+
+ 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));
+
+ // Offer NOT_FOUND status for one of the works.
+ fakeWindmillServer.whenCommitWorkStreamCalled().put(workB.id(),
CommitStatus.NOT_FOUND);
+
+ workCommitter.start();
+ workCommitter.commit(commit);
+
+ // Wait for the server to receive and process the commits
+ fakeWindmillServer.waitForAndGetCommits(3);
+ waitForExpectedSetSize(completeCommits, 3);
+
+ // Verify that FakeWindmillServer received the multi-key commit
+ List<Windmill.MultiKeyWorkItemCommitRequest> multiKeyCommits =
+ fakeWindmillServer.getMultiKeyCommitsReceived();
+ assertThat(multiKeyCommits).hasSize(1);
+ assertThat(multiKeyCommits.get(0)).isEqualTo(multiKeyRequest);
+
+ // Verify all three works in the multi-key commit are completed with
NOT_FOUND status
Review Comment:
Added a check for `currentActiveCommitBytes() == 0` that checks commits are
not pending in commiter anymore.
##########
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingEngineWorkCommitterTest.java:
##########
@@ -186,10 +186,14 @@ public void testCommit_sendsCommitsToStreamingEngine() {
waitForExpectedSetSize(completeCommits, 5);
for (Commit commit : commits) {
- WorkItemCommitRequest request =
committed.get(commit.work().getWorkItem().getWorkToken());
+ WorkItemCommitRequest request =
+
committed.get(commit.workBatch().get(0).getWorkItem().getWorkToken());
Review Comment:
done.
##########
runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/commits/StreamingApplianceWorkCommitterTest.java:
##########
@@ -129,9 +129,9 @@ public void testCommit() {
for (Commit commit : commits) {
Windmill.WorkItemCommitRequest request =
- committed.get(commit.work().getWorkItem().getWorkToken());
+
committed.get(commit.workBatch().get(0).getWorkItem().getWorkToken());
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]