You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使用Vertx和Redis持续监听消息?现有代码仅能接收一次消息

How to Implement Continuous Listening with Vert.x Redis blpop?

I've written the following Listener using Vert.x:

public class Listener extends AbstractVerticle { 
    public static void main(String[] args) { 
        Launcher.executeCommand("run", Listener.class.getName()); 
    } 
    @Override 
    public void start() { 
        RedisOptions config = new RedisOptions() 
                .setHost("127.0.0.1"); 
        RedisClient redis = RedisClient.create(vertx, config); 
        redis.blpop("myKey", 3500, System.out::println); 
    } 
}

The current code successfully prints the first received message, but it can't receive subsequent messages. How can I implement continuous listening?


Great question! The core issue here is that blpop is a single-operation blocking command—it only fetches one message and then exits. To maintain persistent listening, you need to re-trigger the blpop call after each successful message retrieval, working within Vert.x's asynchronous callback pattern.

Here's a modified version of your code that achieves continuous listening:

public class Listener extends AbstractVerticle { 
    private RedisClient redis;

    public static void main(String[] args) { 
        Launcher.executeCommand("run", Listener.class.getName()); 
    } 

    @Override 
    public void start() { 
        RedisOptions config = new RedisOptions()
                .setHost("127.0.0.1"); 
        redis = RedisClient.create(vertx, config);
        
        // Kick off the first listening cycle
        startListening();
    }

    private void startListening() {
        // Execute blpop, and in the callback, start the next listen cycle
        redis.blpop("myKey", 3500, ar -> {
            if (ar.succeeded()) {
                // Print the received message
                System.out.println(ar.result());
                // Immediately listen for the next message
                startListening();
            } else {
                // Handle errors (e.g., temporary connection drops)
                System.err.println("Listening failed: " + ar.cause().getMessage());
                // Optional: Retry after a short delay to avoid overwhelming the system
                vertx.setTimer(1000, id -> startListening());
            }
        });
    }
}

Key Details:

  • We extract the listening logic into a dedicated startListening() method, making it easy to re-invoke after each message.
  • When a message is successfully received (ar.succeeded()), we immediately call startListening() again to wait for the next message.
  • For error scenarios, we log the issue and add a short delay before retrying—this prevents your app from entering a tight loop if there's a temporary Redis connection problem.
  • This approach aligns with Vert.x's non-blocking philosophy: each blpop runs asynchronously, and the callback triggers the next operation without blocking the event loop.

If you're using a newer Vert.x Redis client (version 4.x+), you could also explore using RedisConnection with subscription patterns, but the recursive callback approach above is a simple, reliable fix for your current use case.


内容的提问来源于stack exchange,提问作者Josemy

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 03:39:35