【问题标题】:Spring Data Redis Streams, Cannot figure out what is happening to my unacknowleded messages?Spring Data Redis Streams,无法弄清楚我未确认的消息发生了什么?
【发布时间】:2020-08-11 11:39:15
【问题描述】:

我正在使用以下代码使用 Spring Data Redis 消费者组来消费 Redis 流,但即使我已经注释掉了确认命令,我的消息在服务器重启后也不会重新读取。

我希望如果我没有确认该消息,则应该在服务器被终止并重新启动时重新读取该消息。我在这里错过了什么?

@Bean
@Autowired
public StreamMessageListenerContainer eventStreamPersistenceListenerContainerTwo(RedisConnectionFactory streamRedisConnectionFactory, RedisTemplate streamRedisTemplate) {

        StreamMessageListenerContainer.StreamMessageListenerContainerOptions<String, MapRecord<String, String, String>> containerOptions = StreamMessageListenerContainer.StreamMessageListenerContainerOptions
                        .builder().pollTimeout(Duration.ofMillis(100)).build();

        StreamMessageListenerContainer<String, MapRecord<String, String, String>> container = StreamMessageListenerContainer.create(streamRedisConnectionFactory,
                        containerOptions);

        container.receive(Consumer.from("my-group", "my-consumer"),
                        StreamOffset.create("event-stream", ReadOffset.latest()),
                        message -> {
                                System.out.println("MessageId: " + message.getId());
                                System.out.println("Stream: " + message.getStream());
                                System.out.println("Body: " + message.getValue());
                                //streamRedisTemplate.opsForStream().acknowledge("my-group", message);
                        });

        container.start();

        return container;
}

【问题讨论】:

    标签: spring spring-boot redis spring-data-redis redis-streams


    【解决方案1】:

    在阅读了有关流如何工作的 Redis 文档后,我想出了以下方法来自动处理任何未确认但之前已为消费者传递的消息:

    // Check for any previously unacknowledged messages that were delivered to this consumer.
    log.info("STREAM - Checking for previously unacknowledged messages for " + this.getClass().getSimpleName() + " event stream listener.");
    String offset = "0";
    while ((offset = processUnacknowledgedMessage(offset)) != null) {
            log.info("STREAM - Finished processing one unacknowledged message for " + this.getClass().getSimpleName() + " event stream listener: " + offset);
    }
    log.info("STREAM - Finished checking for previously unacknowledged messages for " + this.getClass().getSimpleName() + " event stream listener.");
    
    

    以及处理消息的方法:

    /**
     * Processes, and acknowledges the next previously delivered message, beginning
     * at the given message id offset.
     *
     * @param offset The last read message id offset.
     * @return The message that was just processed, or null if there are no more messages.
     */
    public String processUnacknowledgedMessage(String offset) {
            List<MapRecord> messages = streamRedisTemplate.opsForStream().read(Consumer.from(groupName(), consumerName()),
                            StreamReadOptions.empty().noack().count(1),
                            StreamOffset.create(streamKey(), ReadOffset.from(offset)));
            String lastMessageId = null;
            for (MapRecord message : messages) {
                    if (log.isDebugEnabled()) log.debug(String.format("STREAM - Processing event(%s) from stream(%s) during startup: %s", message.getId(), message.getStream(), message.getValue()));
                    processRecord(message);
                    if (log.isDebugEnabled()) log.debug(String.format("STREAM - Finished processing event(%s) from stream(%s) during startup.", message.getId(), message.getStream()));
                    streamRedisTemplate.opsForStream().acknowledge(groupName(), message);
                    lastMessageId = message.getId().getValue();
            }
            return lastMessageId;
    }
    
    

    【讨论】:

    • 说有一条毒消息总是在消费者端产生错误。上述方法将无休止地处理毒消息。理想情况下,我们应该使用基于 deliveryCount 的 XPending api,如果它超过 5,那么我们应该将它移动到死信流,或者我们可以 XClaim 它。但我没有找到用于 Xpending 和 Xclaim 的 spring-data-redis 流 api。
    • spring-data-redis的2.3.1.RELASE版本已经有挂起功能。见docs.spring.io/spring-data/redis/docs/current/api/org/…
    猜你喜欢
    • 1970-01-01
    • 2020-07-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-11-24
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多