【发布时间】:2019-07-06 23:11:29
【问题描述】:
在使用 spring kafka 时,我可以使用以下代码根据时间戳读取主题中的消息 -
ConsumerRecords<String, String> records = consumer.poll(100);
if (flag) {
Map<TopicPartition, Long> query = new HashMap<>();
query.put(new TopicPartition(kafkaTopic, 0), millisecondsFromEpochToReplay);
Map<TopicPartition, OffsetAndTimestamp> result = consumer.offsetsForTimes(query);
if(result != null)
{
records = ConsumerRecords.empty();
}
result.entrySet().stream()
.forEach(entry -> consumer.seek(entry.getKey(), entry.getValue().offset()));
flag = false;
}
如何使用 spring integration DSL 实现相同的功能 - 使用 KafkaMessageDrivenChannelAdapter?
我们如何设置集成流程并根据时间戳从主题中读取消息?
【问题讨论】:
标签: spring-integration kafka-consumer-api spring-integration-dsl