【问题标题】:Apache Kafka Streams : Out-of-Order messagesApache Kafka Streams:乱序消息
【发布时间】:2021-07-13 10:23:45
【问题描述】:

我有一个写入主题 A (TA) 的 Apache Kafka 2.6 Producer。 我还有一个 Kafka 流应用程序,它从 TA 消费并写入主题 B (TB)。 在流应用程序中,我有一个自定义时间戳提取器,它从消息负载中提取时间戳。

对于我的一个故障处理测试用例,我在应用程序运行时关闭了 Kafka 集群。

当生产者应用程序尝试向 TA 写入消息时,它不能,因为集群已关闭,因此(我假设)缓冲了消息。 假设它以递增的时间顺序接收 4 条消息 m1,m2,m3,m4。 (即 m1 是第一个,m4 是最后一个)。

当我让 Kafka 集群重新上线时,生产者将缓冲的消息发送到主题,但它们不是按顺序排列的。例如,我收到 m2,然后是 m3,然后是 m1,然后是 m4。

这是为什么呢?是否因为生产者中的缓冲是多线程的,每个生产者同时对主题进行生产?

我认为自定义时间戳提取器有助于在使用消息时对消息进行排序。但他们没有。或者我对时间戳提取器的理解是错误的。

我从 SO here 获得了一个解决方案,将所有事件从 tA 流式传输到另一个中间主题(比如 tA'),该主题将使用时间戳提取器到另一个主题。但我不确定这是否会导致事件根据提取的时间戳重新排序。

我的Producer代码如下(我使用Spring Cloud创建Producer): Producer.java

@Service
public class Producer {

    private String topicName = "input-topic";
        
    private ApplicationProperties appProps;
    
    @Autowired
    private KafkaTemplate<String, MyEvent> kafkaTemplate;
    
    public Producer() {
        super();        
    }
    
    @Autowired
    public void setAppProps(ApplicationProperties appProps) {
        this.appProps = appProps;
        this.topicName = appProps.getInput().getTopicName();
    }

    public void sendMessage(String key, MyEvent ce) {
        ListenableFuture<SendResult<String,MyEvent>> future = this.kafkaTemplate.send(this.topicName, key, ce); 
        
    }
}

【问题讨论】:

    标签: apache-kafka timestamp extractor


    【解决方案1】:

    这是为什么呢?是否因为生产者中的缓冲是多线程的,每个生产者同时对主题进行生产?

    默认情况下,生产者最多允许向代理发送 5 个并行进行中的请求,因此如果某些请求失败并被重试,请求顺序可能会发生变化。

    为避免此重新排序问题,您可以设置max.in.flight.requests.per.connection = 1(可能会影响性能)或设置enable.idempotence = true

    顺便说一句:你没有说你的主题是有一个分区还是多个分区,你的消息是否有一个键?如果您的主题有多个分区,并且您的消息被发送到不同的分区,则无论如何都无法保证读取的顺序,因为仅在一个分区内保证偏移顺序。

    我认为自定义时间戳提取器有助于在使用消息时对消息进行排序。但他们没有。或者我对时间戳提取器的理解是错误的。

    时间戳提取器仅提取时间戳。 Kafka Streams 不会对任何消息重新排序,而是始终按偏移顺序处理消息。

    如果不是,那么时间戳提取器的具体用途是什么?只是为了将时间戳与事件相关联?

    正确。

    我从 SO here 获得了一个解决方案,将所有事件从 tA 流式传输到另一个中间主题(比如 tA'),该主题将使用 TimeStamp 提取器到另一个主题。但我不确定这是否会导致事件根据提取的时间戳重新排序。

    不,它不会进行任何重新排序。另一个 SO 问题即将更改时间戳,但如果您按 a、b、c 顺序读取消息,则结果将按 a、b、c 顺序写入(只是时间戳不同,但应保留偏移顺序)。

    本次演讲解释了更多细节:https://www.confluent.io/kafka-summit-san-francisco-2019/whats-the-time-and-why/

    【讨论】:

    • 对于您关于我们是否使用密钥的问题,是的,我们使用。我尝试使用 max.in.flight.requests.per.connection=1 并在重试期间保留顺序。那么,enable.idempotence=true 在保持消息顺序的同时也会导致性能下降吗?
    • enable.idempotence=true 也可能有一个性能命中,但它可能小于max.in.flight.request.per.connection=1——此外,如果您启用幂等写入,您还可以在重试的情况下防止重复(否则,重试可能会导致重复附加到主题,因为写入可能已成功,但只有 ack 丢失触发重试)。
    • 另外,我的主题有 2 个分区。如果我的消息有键,那么为什么 Kafka 不保持顺序(即使它们正在重试)?我认为拥有消息密钥的好处之一是确保它们最终都在同一个分区中,并在该分区中排序。是否有任何参考文档可以向我展示重试机制的工作原理?
    • 如果你的消息有相同的键,你很好。但是在你的问题中你没有说任何关于钥匙的事情,所以不清楚他们是否可能有钥匙,或者可能有不同或相同的钥匙。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-08-17
    • 1970-01-01
    • 2020-04-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多