package com.ankurm.protocols; import com.ankurm.protocols.grpc.Quote; import com.ankurm.protocols.grpc.QuoteRequest; import com.ankurm.protocols.grpc.QuoteServiceGrpc; import com.ankurm.protocols.grpc.StreamRequest; import io.grpc.Context; import io.grpc.stub.ServerCallStreamObserver; import io.grpc.stub.StreamObserver; import java.util.concurrent.atomic.AtomicLong; import org.springframework.stereotype.Service; /** * The gRPC side. No registration code: the generated {@code ImplBase} is a * {@code BindableService} and Boot registers every such bean with the server. * *
The cancellation check in {@link #streamQuotes} is not decoration. gRPC does not interrupt
* your thread when a client goes away — it sets a flag on the {@link Context} — so a
* loop that never looks keeps producing into a stream nobody is reading. That is covered at
* length in the Spring gRPC article;
* here it matters because the back-pressure benchmark deliberately abandons a stream.
*
* @see docs/04-backpressure.md
*/
@Service
public class GrpcQuoteService extends QuoteServiceGrpc.QuoteServiceImplBase {
/** Counts what the SERVER produced, comparable with the RSocket and WebSocket counters. */
static final AtomicLong PRODUCED = new AtomicLong();
private final QuoteSource source;
GrpcQuoteService(QuoteSource source) {
this.source = source;
}
@Override
public void getQuote(QuoteRequest request, StreamObserver {@code StreamObserver.onNext} on a gRPC server never blocks. If the
* client is not reading, the message is queued in the server's outbound buffer and the loop
* keeps going — HTTP/2 flow control governs the wire, not your code. The only thing
* that connects the two is {@link ServerCallStreamObserver#isReady()}, and using it means
* restructuring the handler around {@code setOnReadyHandler} rather than writing a loop.
*
* The spin here is deliberately the simplest possible demonstration rather than the
* shape you would ship; a real implementation registers an on-ready handler and returns.
*/
@Override
public void streamUnboundedReady(QuoteRequest request, StreamObserver observer) {
observer.onNext(toProto(source.at(request.getSymbol(), 0)));
observer.onCompleted();
}
/**
* An effectively unbounded stream, for the back-pressure comparison. A consumer that stops
* calling {@code next()} closes the HTTP/2 flow-control window, and this loop blocks inside
* {@code onNext} — so the server stops producing, but it stops by blocking a thread
* rather than by being asked to stop.
*/
@Override
public void streamUnbounded(QuoteRequest request, StreamObserver
observer) {
PRODUCED.set(0);
while (!Context.current().isCancelled()) {
observer.onNext(toProto(source.at(request.getSymbol(), PRODUCED.get())));
PRODUCED.incrementAndGet();
}
}
/**
* The fix, and the reason the unfixed version is worth measuring.
*
*
observer) {
PRODUCED.set(0);
ServerCallStreamObserver
ready = (ServerCallStreamObserver
) observer;
while (!ready.isCancelled()) {
if (!ready.isReady()) {
Thread.onSpinWait();
continue;
}
ready.onNext(toProto(source.at(request.getSymbol(), PRODUCED.get())));
PRODUCED.incrementAndGet();
}
}
@Override
public void produced(QuoteRequest request, StreamObserver
observer) {
for (int i = 0; i < request.getCount(); i++) {
if (Context.current().isCancelled()) {
// Do NOT call onCompleted/onError here: the stream is already closed and
// touching it throws IllegalStateException. Just return.
return;
}
observer.onNext(toProto(source.at(request.getSymbol(), i)));
if (request.getDelayMs() > 0) {
try {
Thread.sleep(request.getDelayMs());
}
catch (InterruptedException ex) {
Thread.currentThread().interrupt();
return;
}
}
}
observer.onCompleted();
}
static Quote toProto(com.ankurm.protocols.Quote q) {
return Quote.newBuilder()
.setSymbol(q.symbol()).setSeq(q.seq())
.setBid(q.bid()).setAsk(q.ask())
.setEpochMicros(q.epochMicros())
.build();
}
}