Class RedisPubSubService

  • All Implemented Interfaces:

    
    public final class RedisPubSubService
    
                        

    An at-most-once messaging service that uses "virtual channels/topics" on top of a single Redis channel.

    This allows for grater scalability where many clients can share a single Redis connection; and automatic re-connect is handled once.

    Messages used must fit in the heap and cannot be streamed using InputStream, because a single received message is distributed/copied among all subscribers.

    • Nested Class Summary

      Nested Classes 
      Modifier and Type Class Description
    • Field Summary

      Fields 
      Modifier and Type Field Description
    • Enum Constant Summary

      Enum Constants 
      Enum Constant Description
    • Method Summary

      Modifier and Type Method Description
      final Unit start(Duration timeout) Starts the service, optionally waiting for a successfully subscribe
      final Unit stop() Stops the service and waits for the background thread to finish
      final RedisPubSubServiceSubscription subscribe(String channel, Function1<ByteArray, Unit> onMsg)
      • Methods inherited from class java.lang.Object

        clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
    • Constructor Detail

      • RedisPubSubService

        RedisPubSubService(RedisConnectionPool connectionPool, String redisChannelName, Long delayBetweenReconnectMs, Integer heartbeatIntervalMs, Function1<Runnable, Runnable> onNewMessageHandler)
        Parameters:
        connectionPool - the connection pool to be used for sending PUBLISH messages
        redisChannelName - the channel to be used on Redis
        delayBetweenReconnectMs - upon connection failure, this delay causes a wait before a reconnect is initiated
        heartbeatIntervalMs - if nothing is received for this long, a PING is sent; if its reply is not received within another interval, the connection is considered dead and a reconnect is initiated
        onNewMessageHandler - each time a new async virtual thread is started, this handler is used to wrap the Runnable.