I'm looking for examples using a plain JDK11+ http client reading server sent events, without extra dependencies. I can't find anything about sse in the documentation either.
Any hints?
I'm looking for examples using a plain JDK11+ http client reading server sent events, without extra dependencies. I can't find anything about sse in the documentation either.
Any hints?
Java 11 based implementation of SSE (Server-sent Events) client here:
SSE Client
It provides a pretty simple usage of processing of SSE messages.
Example usage:
EventHandler eventHandler = eventText -> { process(eventText); };
SSEClient sseClient =
SSEClient sseClient = SSEClient.builder().url(url).eventHandler(eventHandler)
.build();
sseClient.start();
Note: I am the author of this SSE client.
EDIT 1: Info here and here on the format of the incoming data.
EDIT 2: Updated the code sample to handle the data: part of the protocol. There are also event:, id:, and retry: parts (see links above), but I do not plan to add handling for those.
I can't find an official BodySubscriber to do SSE, but it's not that hard to write one. Here's a rough impl (but note the TODOs):
public class SseSubscriber implements BodySubscriber<Void>
{
protected static final Pattern dataLinePattern = Pattern.compile( "^data: ?(.*)$" );
protected static String extractMessageData( String[] messageLines )
{
var s = new StringBuilder( );
for ( var line : messageLines )
{
var m = dataLinePattern.matcher( line );
if ( m.matches( ) )
{
s.append( m.group( 1 ) );
}
}
return s.toString( );
}
protected final Consumer<? super String> messageDataConsumer;
protected final CompletableFuture<Void> future;
protected volatile Subscription subscription;
protected volatile String deferredText;
public SseSubscriber( Consumer<? super String> messageDataConsumer )
{
this.messageDataConsumer = messageDataConsumer;
this.future = new CompletableFuture<>( );
this.subscription = null;
this.deferredText = null;
}
@Override
public void onSubscribe( Subscription subscription )
{
this.subscription = subscription;
try
{
this.deferredText = "";
this.subscription.request( 1 );
}
catch ( Exception e )
{
this.future.completeExceptionally( e );
this.subscription.cancel( );
}
}
@Override
public void onNext( List<ByteBuffer> buffers )
{
try
{
// Volatile read
var deferredText = this.deferredText;
for ( var buffer : buffers )
{
// TODO: Safe to assume multi-byte chars don't get split across buffers?
var s = deferredText + UTF_8.decode( buffer );
// -1 means don't discard trailing empty tokens ... so the final token will
// be whatever is left after the last \n\n (possibly the empty string, but
// not necessarily), which is the part we need to defer until the next loop
// iteration
var tokens = s.split( "\n\n", -1 );
// Final token gets deferred, not processed here
for ( var i = 0; i < tokens.length - 1; i++ )
{
var message = tokens[ i ];
var lines = message.split( "\n" );
var data = extractMessageData( lines );
this.messageDataConsumer.accept( data );
// TODO: Handle lines that start with "event:", "id:", "retry:"
}
// Defer the final token
deferredText = tokens[ tokens.length - 1 ];
}
// Volatile write
this.deferredText = deferredText;
this.subscription.request( 1 );
}
catch ( Exception e )
{
this.future.completeExceptionally( e );
this.subscription.cancel( );
}
}
@Override
public void onError( Throwable e )
{
this.future.completeExceptionally( e );
}
@Override
public void onComplete( )
{
try
{
this.future.complete( null );
}
catch ( Exception e )
{
this.future.completeExceptionally( e );
}
}
@Override
public CompletionStage<Void> getBody( )
{
return this.future;
}
}
Then to use it:
var req = HttpRequest.newBuilder( )
.GET( )
.uri( new URI( "http://service/path/to/events" )
.setHeader( "Accept", "text/event-stream" )
.build( );
this.client.sendAsync( req, respInfo ->
{
if ( respInfo.statusCode( ) == 200 )
{
return new SseSubscriber( messageData ->
{
// TODO: Handle messageData
} );
}
else
{
throw new RuntimeException( "Request failed" );
}
} );