पाठ 22 / 25
Streams in Application Frameworks
Consume streams with common client libraries and frameworks.
Library support
All major Redis clients support Pub/Sub and Streams. In Java, Lettuce and Jedis expose the commands, and Spring Data Redis offers StreamMessageListenerContainer and StreamReceiver for consumer-group listeners with automatic polling, plus RedisMessageListenerContainer for Pub/Sub. In Node.js, node-redis and ioredis provide xAdd/xReadGroup and subscription APIs; BullMQ builds a job system on Redis. In Python, redis-py supports both, including asyncio. In .NET, StackExchange.Redis supports Streams commands (it does not support blocking XREAD, so you poll) and Pub/Sub through its subscriber API. Whatever the library: create consumer groups idempotently, use separate connections for blocking reads and Pub/Sub, set sensible timeouts, handle reconnects (re-subscribe for Pub/Sub, resume from pending for Streams), and give consumers stable, unique names.
One Redis, many client libraries
Different languages share the same streams and channels through their client libraries.
A Spring Data Redis stream listener (Java)
The container polls with XREADGROUP and calls the listener for each record.
@Bean
Subscription orderSubscription(RedisConnectionFactory cf, OrderListener listener) {
var options = StreamMessageListenerContainer.StreamMessageListenerContainerOptions
.builder()
.pollTimeout(Duration.ofSeconds(2))
.batchSize(50)
.build();
var container = StreamMessageListenerContainer.create(cf, options);
var subscription = container.receive( // manual ack: we call XACK
Consumer.from("billing", hostName()),
StreamOffset.create("events:orders", ReadOffset.lastConsumed()),
listener);
container.start();
return subscription;
}
@Component
class OrderListener implements StreamListener<String, MapRecord<String, String, String>> {
@Autowired StringRedisTemplate redis;
public void onMessage(MapRecord<String, String, String> record) {
createInvoice(record.getValue());
redis.opsForStream().acknowledge("billing", record);
}
}Dedicated connections for blocking calls
A blocking XREADGROUP ... BLOCK 5000 occupies its connection for up to five seconds. Do not share that connection with request-path commands, or every request will wait behind it.
त्वरित जाँच: Why should blocking stream reads use a dedicated connection?
- Blocking reads require TLS
- Redis allows only one blocking read per server
- They hold the connection while waiting, delaying other commands sent on it
- They disable pipelining globally
Answer
They hold the connection while waiting, delaying other commands sent on it — A blocked connection cannot serve other commands until the call returns.