Skip to content

Commit 9cfd2a7

Browse files
authored
[Dataflow Streaming] Fix grpc commit stream test (#35552)
Test was expecting ordered requests from the stream. Stream uses a hashmap internally to keep track of requests and ordering is not guaranteed.
1 parent de1aca3 commit 9cfd2a7

1 file changed

Lines changed: 62 additions & 24 deletions

File tree

  • runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc

runners/google-cloud-dataflow-java/worker/src/test/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/GrpcCommitWorkStreamTest.java

Lines changed: 62 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,9 @@
2626
import static org.junit.Assert.assertTrue;
2727

2828
import java.io.IOException;
29+
import java.util.HashMap;
2930
import java.util.HashSet;
31+
import java.util.Map;
3032
import java.util.Set;
3133
import java.util.concurrent.CountDownLatch;
3234
import java.util.concurrent.ExecutionException;
@@ -212,12 +214,21 @@ public void testCommitWorkItem_retryOnNewStream() throws Exception {
212214
}
213215
Windmill.StreamingCommitWorkRequest request = streamInfo.requests.take();
214216
assertEquals(5, request.getCommitChunkCount());
215-
for (int i = 0; i < 5; ++i) {
216-
assertEquals(i + 1, request.getCommitChunk(i).getRequestId());
217-
Windmill.WorkItemCommitRequest parsedRequest =
218-
Windmill.WorkItemCommitRequest.parseFrom(
219-
request.getCommitChunk(i).getSerializedWorkItemCommit());
220-
assertEquals(parsedRequest.getWorkToken(), i);
217+
{
218+
// Check if request ids and work tokens match.
219+
Map<Long, Long> requestIdWorkTokenMap = new HashMap<>();
220+
Map<Long, Long> expectedRequestIdWorkTokenMap = new HashMap<>();
221+
for (int i = 0; i < 5; ++i) {
222+
Windmill.WorkItemCommitRequest parsedRequest =
223+
Windmill.WorkItemCommitRequest.parseFrom(
224+
request.getCommitChunk(i).getSerializedWorkItemCommit());
225+
requestIdWorkTokenMap.put(
226+
request.getCommitChunk(i).getRequestId(), parsedRequest.getWorkToken());
227+
}
228+
for (int i = 1; i <= 5; ++i) {
229+
expectedRequestIdWorkTokenMap.put((long) i, (long) (i - 1));
230+
}
231+
assertThat(requestIdWorkTokenMap).containsExactlyEntriesIn(expectedRequestIdWorkTokenMap);
221232
}
222233
// Send back that 1 and 5 finished.
223234
streamInfo.responseObserver.onNext(
@@ -232,12 +243,21 @@ public void testCommitWorkItem_retryOnNewStream() throws Exception {
232243
waitForConnectionAndConsumeHeader();
233244
Windmill.StreamingCommitWorkRequest reconnectRequest = reconnectStreamInfo.requests.take();
234245
assertEquals(3, reconnectRequest.getCommitChunkCount());
235-
for (int i = 0; i < 3; ++i) {
236-
assertEquals(i + 2, reconnectRequest.getCommitChunk(i).getRequestId());
237-
Windmill.WorkItemCommitRequest parsedRequest =
238-
Windmill.WorkItemCommitRequest.parseFrom(
239-
reconnectRequest.getCommitChunk(i).getSerializedWorkItemCommit());
240-
assertEquals(i + 1, parsedRequest.getWorkToken());
246+
{
247+
// Check if request ids and work tokens match.
248+
Map<Long, Long> requestIdWorkTokenMap = new HashMap<>();
249+
Map<Long, Long> expectedRequestIdWorkTokenMap = new HashMap<>();
250+
for (int i = 0; i < 3; ++i) {
251+
Windmill.WorkItemCommitRequest parsedRequest =
252+
Windmill.WorkItemCommitRequest.parseFrom(
253+
reconnectRequest.getCommitChunk(i).getSerializedWorkItemCommit());
254+
requestIdWorkTokenMap.put(
255+
reconnectRequest.getCommitChunk(i).getRequestId(), parsedRequest.getWorkToken());
256+
}
257+
for (int i = 2; i <= 4; ++i) {
258+
expectedRequestIdWorkTokenMap.put((long) i, (long) (i - 1));
259+
}
260+
assertThat(requestIdWorkTokenMap).containsExactlyEntriesIn(expectedRequestIdWorkTokenMap);
241261
}
242262
// Send back that 2 and 3 finished.
243263
reconnectStreamInfo.responseObserver.onNext(
@@ -281,12 +301,21 @@ public void testCommitWorkItem_retryOnNewStreamHalfClose() throws Exception {
281301
}
282302
Windmill.StreamingCommitWorkRequest request = streamInfo.requests.take();
283303
assertEquals(5, request.getCommitChunkCount());
284-
for (int i = 0; i < 5; ++i) {
285-
assertEquals(i + 1, request.getCommitChunk(i).getRequestId());
286-
Windmill.WorkItemCommitRequest parsedRequest =
287-
Windmill.WorkItemCommitRequest.parseFrom(
288-
request.getCommitChunk(i).getSerializedWorkItemCommit());
289-
assertEquals(parsedRequest.getWorkToken(), i);
304+
{
305+
// Check if request ids and work tokens match.
306+
Map<Long, Long> requestIdWorkTokenMap = new HashMap<>();
307+
Map<Long, Long> expectedRequestIdWorkTokenMap = new HashMap<>();
308+
for (int i = 0; i < 5; ++i) {
309+
Windmill.WorkItemCommitRequest parsedRequest =
310+
Windmill.WorkItemCommitRequest.parseFrom(
311+
request.getCommitChunk(i).getSerializedWorkItemCommit());
312+
requestIdWorkTokenMap.put(
313+
request.getCommitChunk(i).getRequestId(), parsedRequest.getWorkToken());
314+
}
315+
for (int i = 1; i <= 5; ++i) {
316+
expectedRequestIdWorkTokenMap.put((long) i, (long) (i - 1));
317+
}
318+
assertThat(requestIdWorkTokenMap).containsExactlyEntriesIn(expectedRequestIdWorkTokenMap);
290319
}
291320
// Half-close the logical stream. This shouldn't prevent reconnection of the physical stream
292321
// from succeeding.
@@ -310,12 +339,21 @@ public void testCommitWorkItem_retryOnNewStreamHalfClose() throws Exception {
310339

311340
Windmill.StreamingCommitWorkRequest reconnectRequest = reconnectStreamInfo.requests.take();
312341
assertEquals(3, reconnectRequest.getCommitChunkCount());
313-
for (int i = 0; i < 3; ++i) {
314-
assertEquals(i + 2, reconnectRequest.getCommitChunk(i).getRequestId());
315-
Windmill.WorkItemCommitRequest parsedRequest =
316-
Windmill.WorkItemCommitRequest.parseFrom(
317-
reconnectRequest.getCommitChunk(i).getSerializedWorkItemCommit());
318-
assertEquals(i + 1, parsedRequest.getWorkToken());
342+
{
343+
// Check if request ids and work tokens match.
344+
Map<Long, Long> requestIdWorkTokenMap = new HashMap<>();
345+
Map<Long, Long> expectedRequestIdWorkTokenMap = new HashMap<>();
346+
for (int i = 0; i < 3; ++i) {
347+
Windmill.WorkItemCommitRequest parsedRequest =
348+
Windmill.WorkItemCommitRequest.parseFrom(
349+
reconnectRequest.getCommitChunk(i).getSerializedWorkItemCommit());
350+
requestIdWorkTokenMap.put(
351+
reconnectRequest.getCommitChunk(i).getRequestId(), parsedRequest.getWorkToken());
352+
}
353+
for (int i = 2; i <= 4; ++i) {
354+
expectedRequestIdWorkTokenMap.put((long) i, (long) (i - 1));
355+
}
356+
assertThat(requestIdWorkTokenMap).containsExactlyEntriesIn(expectedRequestIdWorkTokenMap);
319357
}
320358
assertNull(streamInfo.onDone.get());
321359

0 commit comments

Comments
 (0)