Skip to content

Commit bcc8961

Browse files
authored
[Dataflow Streaming Java] Fix possible IllegalStateException when grpc streams have deadline exceeded. (#36170)
1 parent 893e9cb commit bcc8961

4 files changed

Lines changed: 95 additions & 13 deletions

File tree

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/ResettableThrowingStreamObserver.java

Lines changed: 22 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -115,14 +115,24 @@ public void onNext(T t) throws StreamClosedException, WindmillStreamShutdownExce
115115
logger.debug("Stream was shutdown during send.", cancellationException);
116116
return;
117117
}
118+
if (delegateStreamObserver == delegate) {
119+
if (isCurrentStreamClosed) {
120+
logger.debug("Stream is already closed when encountering error with send.");
121+
return;
122+
}
123+
isCurrentStreamClosed = true;
124+
}
118125
}
119126

127+
// Either this was the active observer the current observer that requires closing, or this was
128+
// a previous
129+
// observer which we attempt to close and ignore possible exceptions.
120130
try {
121131
delegate.onError(cancellationException);
122132
} catch (IllegalStateException onErrorException) {
123133
// The delegate above was already terminated via onError or onComplete.
124-
// Fallthrough since this is possibly due to queued onNext() calls that are being made from
125-
// previously blocked threads.
134+
// Fallthrough since this is possibly due to queued onNext() calls that are being made
135+
// from previously blocked threads.
126136
} catch (RuntimeException onErrorException) {
127137
logger.warn(
128138
"Encountered unexpected error {} when cancelling due to error.",
@@ -134,14 +144,20 @@ public void onNext(T t) throws StreamClosedException, WindmillStreamShutdownExce
134144

135145
public synchronized void onError(Throwable throwable)
136146
throws StreamClosedException, WindmillStreamShutdownException {
137-
delegate().onError(throwable);
138-
isCurrentStreamClosed = true;
147+
try {
148+
delegate().onError(throwable);
149+
} finally {
150+
isCurrentStreamClosed = true;
151+
}
139152
}
140153

141154
public synchronized void onCompleted()
142155
throws StreamClosedException, WindmillStreamShutdownException {
143-
delegate().onCompleted();
144-
isCurrentStreamClosed = true;
156+
try {
157+
delegate().onCompleted();
158+
} finally {
159+
isCurrentStreamClosed = true;
160+
}
145161
}
146162

147163
synchronized boolean isClosed() {

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/observers/DirectStreamObserver.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -182,8 +182,8 @@ public void onError(Throwable t) {
182182
Preconditions.checkState(!isUserClosed);
183183
isUserClosed = true;
184184
if (!isOutboundObserverClosed) {
185-
outboundObserver.onError(t);
186185
isOutboundObserverClosed = true;
186+
outboundObserver.onError(t);
187187
}
188188
}
189189
}

runners/google-cloud-dataflow-java/worker/src/main/java/org/apache/beam/runners/dataflow/worker/windmill/client/grpc/observers/StreamObserverCancelledException.java

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -21,15 +21,15 @@
2121

2222
@Internal
2323
public final class StreamObserverCancelledException extends RuntimeException {
24-
StreamObserverCancelledException(Throwable cause) {
24+
public StreamObserverCancelledException(Throwable cause) {
2525
super(cause);
2626
}
2727

28-
StreamObserverCancelledException(String message, Throwable cause) {
28+
public StreamObserverCancelledException(String message, Throwable cause) {
2929
super(message, cause);
3030
}
3131

32-
StreamObserverCancelledException(String message) {
32+
public StreamObserverCancelledException(String message) {
3333
super(message);
3434
}
3535
}

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

Lines changed: 69 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -18,12 +18,15 @@
1818
package org.apache.beam.runners.dataflow.worker.windmill.client;
1919

2020
import static org.junit.Assert.assertThrows;
21+
import static org.mockito.ArgumentMatchers.any;
2122
import static org.mockito.ArgumentMatchers.eq;
2223
import static org.mockito.ArgumentMatchers.isA;
24+
import static org.mockito.Mockito.doThrow;
2325
import static org.mockito.Mockito.spy;
2426
import static org.mockito.Mockito.verify;
2527
import static org.mockito.Mockito.verifyNoInteractions;
2628

29+
import org.apache.beam.runners.dataflow.worker.windmill.client.grpc.observers.StreamObserverCancelledException;
2730
import org.apache.beam.runners.dataflow.worker.windmill.client.grpc.observers.TerminatingStreamObserver;
2831
import org.junit.Test;
2932
import org.junit.runner.RunWith;
@@ -51,6 +54,53 @@ public void terminate(Throwable terminationException) {}
5154
});
5255
}
5356

57+
@Test
58+
public void testOnNext_simple() throws Exception {
59+
ResettableThrowingStreamObserver<Integer> observer = newStreamObserver();
60+
TerminatingStreamObserver<Integer> spiedDelegate = newDelegate();
61+
observer.reset(spiedDelegate);
62+
observer.onNext(1);
63+
verify(spiedDelegate).onNext(eq(1));
64+
observer.onNext(2);
65+
verify(spiedDelegate).onNext(eq(2));
66+
observer.onCompleted();
67+
verify(spiedDelegate).onCompleted();
68+
}
69+
70+
@Test
71+
public void testOnError_success() throws Exception {
72+
ResettableThrowingStreamObserver<Integer> observer = newStreamObserver();
73+
TerminatingStreamObserver<Integer> spiedDelegate = newDelegate();
74+
observer.reset(spiedDelegate);
75+
Throwable t = new RuntimeException("Test exception");
76+
observer.onError(t);
77+
verify(spiedDelegate).onError(eq(t));
78+
79+
assertThrows(
80+
ResettableThrowingStreamObserver.StreamClosedException.class, () -> observer.onNext(1));
81+
assertThrows(
82+
ResettableThrowingStreamObserver.StreamClosedException.class, observer::onCompleted);
83+
assertThrows(
84+
ResettableThrowingStreamObserver.StreamClosedException.class,
85+
() -> observer.onError(new RuntimeException("ignored")));
86+
}
87+
88+
@Test
89+
public void testOnCompleted_success() throws Exception {
90+
ResettableThrowingStreamObserver<Integer> observer = newStreamObserver();
91+
TerminatingStreamObserver<Integer> spiedDelegate = newDelegate();
92+
observer.reset(spiedDelegate);
93+
observer.onCompleted();
94+
verify(spiedDelegate).onCompleted();
95+
assertThrows(
96+
ResettableThrowingStreamObserver.StreamClosedException.class, () -> observer.onNext(1));
97+
assertThrows(
98+
ResettableThrowingStreamObserver.StreamClosedException.class, observer::onCompleted);
99+
assertThrows(
100+
ResettableThrowingStreamObserver.StreamClosedException.class,
101+
() -> observer.onError(new RuntimeException("ignored")));
102+
}
103+
54104
@Test
55105
public void testPoison_beforeDelegateSet() {
56106
ResettableThrowingStreamObserver<Integer> observer = newStreamObserver();
@@ -97,9 +147,7 @@ public void testOnCompleted_afterPoisonedThrows() {
97147
}
98148

99149
@Test
100-
public void testReset_usesNewDelegate()
101-
throws WindmillStreamShutdownException,
102-
ResettableThrowingStreamObserver.StreamClosedException {
150+
public void testReset_usesNewDelegate() throws Exception {
103151
ResettableThrowingStreamObserver<Integer> observer = newStreamObserver();
104152
TerminatingStreamObserver<Integer> firstObserver = newDelegate();
105153
observer.reset(firstObserver);
@@ -113,6 +161,24 @@ public void testReset_usesNewDelegate()
113161
verify(secondObserver).onNext(eq(2));
114162
}
115163

164+
@Test
165+
public void testOnNext_streamCancelledException_closesStream() throws Exception {
166+
ResettableThrowingStreamObserver<Integer> observer = newStreamObserver();
167+
TerminatingStreamObserver<Integer> spiedDelegate = newDelegate();
168+
StreamObserverCancelledException streamObserverCancelledException =
169+
new StreamObserverCancelledException("Test error");
170+
doThrow(streamObserverCancelledException).when(spiedDelegate).onNext(any());
171+
observer.reset(spiedDelegate);
172+
observer.onNext(1);
173+
174+
verify(spiedDelegate).onError(eq(streamObserverCancelledException));
175+
assertThrows(
176+
ResettableThrowingStreamObserver.StreamClosedException.class,
177+
() -> observer.onError(new Exception()));
178+
assertThrows(
179+
ResettableThrowingStreamObserver.StreamClosedException.class, observer::onCompleted);
180+
}
181+
116182
private <T> ResettableThrowingStreamObserver<T> newStreamObserver() {
117183
return new ResettableThrowingStreamObserver<>(LoggerFactory.getLogger(getClass()));
118184
}

0 commit comments

Comments
 (0)