पाठ 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 central Redis box connected to four differently shaped client icons around it.
Figure 8.1 — Clients in several languages sharing streams and channels.

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.