【问题标题】:How to make every message process successfully?如何使每个消息处理成功?
【发布时间】:2021-11-06 00:41:53
【问题描述】:

下面是一个包含 3 个 Go-routines 的服务,用于处理来自 Kafka 的消息:


Channel-1 和 Channel-2 是 Go 中的无缓冲数据通道。 Channel 就像一种排队机制。

Goroutine-1 从 kafka 主题读取消息,在消息验证后将其消息负载扔到 Channel-1。

Goroutine-2 从 Channel-1 读取并处理 payload 并将处理后的 payload 扔到 Channel-2 上。

Goroutine-3 从 Channel-2 读取数据,并将处理后的 payload 封装成 http 数据包,然后向另一个服务执行 http 请求(使用 http 客户端)。

上述流程中的漏洞:在我们的例子中,由于服务之间的网络连接不良或远程服务尚未准备好接受来自 Go-routine3(http 客户端超时)的 http 请求,因此处理失败,因此,上述服务丢失该消息(已从 Kafka 主题中读取)。


Goroutine-1 当前订阅来自 Kafka 的消息没有向 Kafka 发送确认(通知 Goroutine-3 已成功处理特定消息)

正确性比性能更重要。


如何保证每条消息都处理成功?

【问题讨论】:

    标签: go apache-kafka message-queue


    【解决方案1】:

    例如,通过新的 Channel-3 将来自 Goroutine-3 的反馈添加到 Goroutine-1。 Goroutine-1 会一直阻塞,直到得到 Channel-3 的确认。

    // in gorouting 1
    channel1 <- data
    select {
        case <-channel3:
        case <-ctx.Done(): // or smth else to prevent deadlock 
    }
    ...
    // in gorouting 3
    data := <-channel2
    for {
        if err := sendData(data); err == nil {
            break
        }
    }
    channel3<-struct{}{}
    

    【讨论】:

    • 是的,我更喜欢反馈方法
    • 但只是想知道,在服务关闭期间,关闭原则(通道)如何应用?按照关闭原则,Goroutine3 应该关闭 channel3
    • 这取决于实现。在 goroutine 3 中,您可以从两个通道中进行选择:通道 2 处理数据和关闭通道(例如 ctx.Done())关闭通道 3 并返回
    • 我在反馈方法中发现了一个问题:我无法将其扩展到多个管道(其中每个管道是上面显示的 Goroutine-1、Goroutine-2 和 Go-routine-3),因为如果一个管道发送否定确认,然后再次重新处理该消息中的所有消息
    • Goroutine 在成功执行其工作之前不会发送确认。如果需要一种“否定”确认,则应根据上级逻辑来处理该否定确认:例如,这些消息是否可以延迟处理,或者它们必须按顺序处理,或者它们是否可以在一段时间后被丢弃不成功的尝试次数,或者其他
    【解决方案2】:

    为确保正确性,您需要在处理成功完成后提交(=确认)消息。
    对于处理未成功完成的情况——一般情况下,您需要自己实现重试机制。
    这应该特定于您的用例,但通常您将消息扔回专用的 Kafka 重试主题(您创建),添加睡眠并再次处理消息。如果在 x 次后处理失败 - 您将消息扔到 DLQ(=死信队列)。
    你可以在这里阅读更多:
    https://eng.uber.com/reliable-reprocessing/
    https://www.confluent.io/blog/error-handling-patterns-in-kafka/

    【讨论】:

    • 而不是重试,如果处理没有成功完成,那么我们以false而不是true向kafka提交(=ack),这样消息就会再次从kafka消费。
    • @overexchange 这样你就会陷入死锁,我认为你需要在 x 次后停止,这可以通过重试机制来实现。
    猜你喜欢
    • 2020-06-11
    • 1970-01-01
    • 2021-09-23
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-07-26
    • 2021-10-20
    • 2020-08-31
    相关资源
    最近更新 更多