如何使用Vertx和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 callstartListening()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
blpopruns 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

