# Streams in Application Frameworks — Redis Pub/Sub & Streams

Source: https://www.skillbyai.com/en/redis-streams/p-clients

> 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.](assets/figures/redis-streams/section-8-map.svg) — 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.

```java
@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.

**Quiz:** Why should blocking stream reads use a dedicated connection?

- [ ] Blocking reads require TLS
- [ ] Redis allows only one blocking read per server
- [x] 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.
