【发布时间】: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