【问题标题】:Manually acknowledge Kafka Event A consuming after producing event B在产生事件 B 后手动确认 Kafka 事件 A 消费
【发布时间】:2020-02-20 12:42:57
【问题描述】:

我有一种情况,我必须消耗事件 A 并进行一些处理,然后生成事件 B。所以我的问题是处理崩溃并且应用程序在消耗已经消耗 A 时无法生成 B。我的方法是在成功发布 B 后确认,我是正确的还是应该针对这种情况实施其他解决方案?

@KafkaListener(
        id = TOPIC_ID,
        topics = TOPIC_ID,
        groupId = GROUP_ID,
        containerFactory = LISTENER_CONTAINER_FACTORY
)
public void listen(List<Message<A>> messages, Acknowledgment acknowledgment) {

    try {
        final AEvent aEvent = messages.stream()
                .filter(message -> null != message.getPayload())
                .map(Message::getPayload)
                .findFirst()
                .get();

        processDao.doSomeProcessing() // returns a Mono<Example> by calling an externe API
                .subscribe(
                        response -> {
                            ProducerRecord<String, BEvent> BEventRecord = new ProducerRecord<>(TOPIC_ID, null, BEvent);

                            ListenableFuture<SendResult<String, BEvent>> future = kafkaProducerTemplate.send(buildBEvent());
                            future.addCallback(new ListenableFutureCallback<SendResult<String, BEvent>>() {
                                @Override
                                public void onSuccess(SendResult<String, BEvent> BEventSendResult) {
                                    //TODO: do when event published successfully
                                }

                                @Override
                                public void onFailure(Throwable exception) {
                                    exception.printStackTrace();
                                    throw new ExampleException();
                                }
                            });
                        },
                        error -> {
                            error.printStackTrace();
                            throw new ExampleException();
                        }
                );
        acknowledgment.acknowledge(); // ??
    } catch (ExampleException) {
        exception.printStackTrace();
    }
}

【问题讨论】:

    标签: java spring spring-boot apache-kafka spring-kafka


    【解决方案1】:

    在使用 reactor 等异步代码时,您无法管理 kafka “确认”。

    Kafka 不管理每个主题/分区的离散确认,只管理分区的最后提交偏移量。

    如果您异步处理两条记录,您将争夺首先提交哪个偏移量。

    您需要在侦听器容器线程上执行发送以保持正确的顺序。

    【讨论】:

    • 在这种情况下如何使用监听线程?你能否给我一个提示,我应该如何重构这段代码以获得正确的顺序?谢谢
    • 我现在反应不太好,但我认为您需要在Mono&lt;?&gt; 上调用block() 而不是subscribe() 并使用结果发送消息。然后,您可以在未来完成时调用确认,或者抛出异常并配置 SeekToCurrentErrorHandler 以便重新传递记录。不过,您确实不需要使用手动确认,AckMode.BATCH(默认)将导致容器在侦听器正常退出时提交偏移量。
    • 如果要并行处理列表;您将需要更多复杂性 - 保留您的订阅者并使用某种闩锁来阻止侦听器线程,直到所有处理完成。例如new CountDownLatch(messages.size()) 并在每次完成时倒计时锁存器,以便在它们全部完成后释放侦听器线程。
    • 如果我理解得很好,通过使用错误处理程序SeekToCurrentErrorHandler (如here),我的应用程序将再次收听相同的消息,直到进程正常完成,所以不需要实施手动知识,对吗?
    • 正确 - 但请记住,使用批处理侦听器; SeekToCurrentBatchErrorHandler 无法知道批次中的哪条记录失败,因此它永远无法恢复,并且将无限期地重新传递记录。这对于瞬态错误是可以的,但通常,对于批处理侦听器,错误处理/恢复必须在侦听器本身中完成。您仍然可以使用SeekToCurrentBatchErrorHandler,但如果同一条记录一遍又一遍地失败,您将需要在侦听器中进行一些逻辑来放弃。
    猜你喜欢
    • 2019-01-15
    • 1970-01-01
    • 2021-03-14
    • 2018-02-08
    • 2020-11-15
    • 2020-03-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多