Spring WebFlux REST Controller仅处理前两个订阅问题排查
The Root Cause
Chances are, the issue boils down to how you're creating and managing your Flux with Flux.create, or a misconfiguration in how your data stream handles multiple subscribers. Here are the most likely culprits:
- Reusing a Singleton Flux Instance: If your
topWordsStream.getTopWords()returns a pre-created singletonFlux, you're in trouble.Flux.createis a cold stream by default—meaning it should initialize a new data source for every subscriber. Reusing the same instance might cause the data generation logic to only run once (for the first two subscribers) and never restart for new requests. - Shared Resources in the Create Callback: If you're spinning up a single thread, opening a shared connection, or holding onto a limited resource inside the
Flux.createconsumer, that resource could get exhausted or locked after two subscriptions. For example, a single data-pushing thread might be hardcoded to only serve two clients before shutting down. - Missing Multicast for Shared Streams: If you want all clients to receive the same real-time
TopWordsdata (a hot stream) but haven't converted your coldFluxto support multicast, only the first few subscribers might get data before the stream's buffer is exhausted.
Fixes You Can Implement
1. Create a Fresh Flux for Every Subscriber
Instead of returning a singleton Flux, generate a new one each time getTopWords() is called. This ensures every client gets its own independent data stream:
// Inside your TopWordsStream class public Flux<TopWords> getTopWords() { // Create a new Flux instance for each subscriber return Flux.create(emitter -> { // Spin up a dedicated thread for this subscriber's data feed Thread dataGenerator = new Thread(() -> { try { while (!emitter.isCancelled()) { // Replace with your actual TopWords generation logic TopWords latestWords = computeRealTimeTopWords(); emitter.next(latestWords); Thread.sleep(1000); // Simulate delay between updates } } catch (InterruptedException e) { Thread.currentThread().interrupt(); } finally { emitter.complete(); } }); dataGenerator.start(); // Clean up the thread when the subscriber disconnects emitter.onDispose(dataGenerator::interrupt); }); }
2. Convert to a Hot Stream for Shared Data
If you want all clients to receive the same real-time stream (e.g., everyone sees the same TopWords updates), convert your cold Flux to a hot stream using share():
// Inside your TopWordsStream class private final Flux<TopWords> sharedTopWordsStream; public TopWordsStream() { // Create the base cold stream Flux<TopWords> coldStream = Flux.create(emitter -> { // Your data generation logic here (runs once when the first subscriber connects) // ... }); // Convert to a hot stream that multicasts to all subscribers sharedTopWordsStream = coldStream.share(); } public Flux<TopWords> getTopWords() { return sharedTopWordsStream; }
3. Fix Resource Management
- Avoid shared resources in the
Flux.createcallback: Every subscriber should get its own thread, connection, or data source instance. - Clean up resources properly: Use
emitter.onDispose()oremitter.onCancel()to release resources when a subscriber disconnects—this prevents leaks that could block new subscriptions. - Check for hardcoded limits: Ensure you don't have any code that caps the number of subscribers (e.g., a counter that stops accepting new connections after 2).
4. Tweak WebFlux Server Configurations
Double-check your server settings to ensure you're not limiting concurrent connections. For example, in application.properties:
# For Tomcat (if using Tomcat instead of Netty) server.tomcat.max-threads=200 server.tomcat.accept-count=100 # For Netty (default in Spring Boot WebFlux) server.netty.max-connections=1000 server.netty.connection-timeout=30s
How to Verify the Fix
- Spin up 3+ independent client processes and hit the
/top-wordsendpoint—all should receive data. - Check server logs for resource exhaustion errors (e.g., thread pool limits, connection timeouts).
- Debug the
Flux.createcallback to confirm it's triggered for every new subscriber.
内容的提问来源于stack exchange,提问作者Tomas Pinos

