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

基于Spring Data Redis实现可靠Pub/Sub:消息重试机制求助

Fixing Reliable Message Retries with Spring Data Redis Pub/Sub

Hey there! Let's break down what's not working in your current approach and build a reliable retry mechanism for your Pub/Sub service.

First: Why Your Current Transaction Approach Isn't Working

Redis transactions (multi/exec) are designed for atomicity of key-value operations—they have no effect on Pub/Sub messages. Here's why:

  • Redis Pub/Sub is a fire-and-forget system: once a message is pushed to subscribers, Redis doesn't store it anymore.
  • Transactions only control whether a batch of key-based commands (like set, hget) are executed atomically. They can't "undo" a Pub/Sub message delivery or push it back to the topic.

So we need a different strategy to handle failed message processing.

Solution 1: Manual Retry with Pub/Sub + Redis Storage

We'll add retry tracking and republish failed messages (with a retry limit to avoid infinite loops) and a dead-letter queue for messages that can't be fixed.

Step 1: Update Your Message Listener Adapter

Modify your MyMessageListenerAdapter to inject the publisher and Redis template, then handle failures by republishing messages:

public class MyMessageListenerAdapter extends MessageListenerAdapter{
    private final RedisMessagePublisher messagePublisher;
    private final RedisTemplate<String, Object> redisTemplate;
    private static final int MAX_RETRIES = 3;
    private static final String RETRY_COUNT_KEY_PREFIX = "retry:count:";
    private static final String DEAD_LETTER_QUEUE = "dead:letter:queue";

    // Inject required beans via constructor
    public MyMessageListenerAdapter(RedisMessageSubscriber redisMessageSubscriber, 
                                   RedisMessagePublisher messagePublisher,
                                   RedisTemplate<String, Object> redisTemplate) {
        super(redisMessageSubscriber);
        this.messagePublisher = messagePublisher;
        this.redisTemplate = redisTemplate;
    }

    @Override
    public void onMessage(Message message, byte[] pattern) {
        String messageBody = new String(message.getBody());
        String retryCountKey = RETRY_COUNT_KEY_PREFIX + messageBody.hashCode(); // Use a unique message ID for production apps
        Integer retryCount = redisTemplate.opsForValue().get(retryCountKey);
        if (retryCount == null) retryCount = 0;

        try {
            // Execute your business logic
            String result = doSomeLogic(message);
            // On success: clean up retry tracking
            redisTemplate.delete(retryCountKey);
            System.out.println("Message processed successfully: " + messageBody);
        } catch (Exception e) {
            // On failure: handle retries
            if (retryCount < MAX_RETRIES) {
                retryCount++;
                // Store retry count with an expiration to avoid memory leaks
                redisTemplate.opsForValue().set(retryCountKey, retryCount, 10, TimeUnit.MINUTES);
                // Republish the message to the original topic
                messagePublisher.publish(messageBody);
                System.out.printf("Message failed, republishing (retry %d/%d): %s%n", retryCount, MAX_RETRIES, messageBody);
            } else {
                // Retries exhausted: send to dead-letter queue for manual inspection
                redisTemplate.opsForList().rightPush(DEAD_LETTER_QUEUE, messageBody);
                redisTemplate.delete(retryCountKey);
                System.out.println("Max retries reached, sent to dead-letter queue: " + messageBody);
            }
        }
    }

    // Your business logic method (throws exceptions on failure)
    private String doSomeLogic(Message message) throws Exception {
        String content = new String(message.getBody());
        // Simulate a failure scenario
        if ("error".equals(content)) {
            throw new RuntimeException("Business logic failed for message: " + content);
        }
        return "success";
    }
}

Step 2: Update Your Config to Inject Dependencies

Modify the messageListener bean to pass in the required dependencies:

@Bean
public MessageListenerAdapter messageListener(RedisMessageSubscriber redisMessageSubscriber,
                                             @Qualifier("RedisMessagePublisherBean") RedisMessagePublisher messagePublisher,
                                             RedisTemplate<String, Object> redisTemplate) {
    MyMessageListenerAdapter messageListeneradapter = new MyMessageListenerAdapter(redisMessageSubscriber, messagePublisher, redisTemplate);
    messageListeneradapter.afterPropertiesSet();
    return messageListeneradapter;
}

Solution 2: Use Redis Streams (Better for Reliable Messaging)

If you're open to moving beyond basic Pub/Sub, Redis Streams are purpose-built for reliable, persistent messaging. They include built-in features like:

  • Message persistence (messages are stored until explicitly deleted)
  • Consumer groups (multiple subscribers can share load)
  • Pending Entries List (PEL) for unacknowledged messages
  • Automatic retry and message redelivery

Example Stream Implementation

Stream Listener Component

@Component
public class StreamMessageHandler implements StreamListener<String, MapRecord<String, String, String>> {
    private final StringRedisTemplate stringRedisTemplate;
    private static final String STREAM_NAME = "mq_stream";
    private static final String CONSUMER_GROUP = "mq_consumer_group";

    public StreamMessageHandler(StringRedisTemplate stringRedisTemplate) {
        this.stringRedisTemplate = stringRedisTemplate;
        // Initialize consumer group (run once; ignore if already exists)
        try {
            stringRedisTemplate.opsForStream().createGroup(STREAM_NAME, CONSUMER_GROUP);
        } catch (RedisSystemException e) {
            // Group already exists, skip initialization
        }
    }

    @Override
    public void onMessage(MapRecord<String, String, String> message) {
        try {
            String content = message.getValue().get("content");
            // Execute business logic
            if ("error".equals(content)) {
                throw new RuntimeException("Failed to process message: " + content);
            }
            // Acknowledge message to remove it from the pending list
            stringRedisTemplate.opsForStream().acknowledge(CONSUMER_GROUP, message);
            System.out.println("Stream message processed: " + content);
        } catch (Exception e) {
            // Failed messages stay in PEL and can be retried later with XCLAIM
            System.out.println("Stream message failed, remains in PEL for retry: " + message.getValue());
        }
    }
}

Stream Listener Container Config

@Bean
public RedisMessageListenerContainer streamContainer(LettuceConnectionFactory connectionFactory,
                                                    StreamMessageListenerAdapter<String, MapRecord<String, String, String>> streamListenerAdapter) {
    RedisMessageListenerContainer container = new RedisMessageListenerContainer();
    container.setConnectionFactory(connectionFactory);
    // Subscribe to the stream, starting from the last consumed message
    container.addMessageListener(streamListenerAdapter, StreamOffset.create(STREAM_NAME, ReadOffset.lastConsumed()));
    return container;
}

@Bean
public StreamMessageListenerAdapter<String, MapRecord<String, String, String>> streamListenerAdapter(StreamMessageHandler handler) {
    StreamMessageListenerAdapter<String, MapRecord<String, String, String>> adapter = new StreamMessageListenerAdapter<>(handler);
    adapter.setStreamConverter(new MappingJackson2MessageConverter());
    return adapter;
}

Key Takeaways

  • Avoid using transactions for Pub/Sub: They don't work for message redelivery.
  • Manual retry for Pub/Sub: Track retry counts in Redis, republish failed messages, and use a dead-letter queue for exhausted retries.
  • Redis Streams: A better fit for reliable messaging with built-in persistence, retry, and load-balancing features.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 06:54:30