Package com.dehnes.kotlinredis.pubsub
Class RedisPubSubService
-
- All Implemented Interfaces:
public final class RedisPubSubServiceAn 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.
-
-
Constructor Summary
Constructors Constructor Description RedisPubSubService(RedisConnectionPool connectionPool, String redisChannelName, Long delayBetweenReconnectMs, Integer heartbeatIntervalMs, Function1<Runnable, Runnable> onNewMessageHandler)
-
Method Summary
Modifier and Type Method Description final Unitstart(Duration timeout)Starts the service, optionally waiting for a successfully subscribe final Unitstop()Stops the service and waits for the background thread to finish final RedisPubSubServiceSubscriptionsubscribe(String channel, Function1<ByteArray, Unit> onMsg)-
-
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 messagesredisChannelName- the channel to be used on RedisdelayBetweenReconnectMs- upon connection failure, this delay causes a wait before a reconnect is initiatedheartbeatIntervalMs- 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 initiatedonNewMessageHandler- each time a new async virtual thread is started, this handler is used to wrap the Runnable.
-
-
-
-