Axon Kafka扩展4.5.4:如何设置StreamableMessageSource初始令牌位置?
问题描述
使用Axon Kafka扩展4.5.4版本,通过MultiStreamableMessageSource接收多数据源事件时,尝试通过以下配置设置初始跟踪令牌为head(最新位置),抛出UnsupportedOperationException异常:
val config = TrackingEventProcessorConfiguration.forSingleThreadedProcessing() .andInitialTrackingToken { it.createHeadToken() }
疑问点:
- 是否需要自定义
KafkaTrackingToken实现?有没有示例? - 不知道具体分区偏移量时该如何处理?
- 多
MessageSource场景下,如何避免繁琐的重复实现?
备注:已知SubscribableMessageSource但不适用当前场景,需依赖MultiStreamableMessageSource合并多数据源。
解决方案
1. 异常原因
StreamableMessageSource的默认createHeadToken()方法未适配Kafka场景,直接调用会触发未实现异常,需手动构建KafkaTrackingToken并指定各主题分区的最新偏移量。
2. 自定义KafkaTrackingToken实现示例
通过Kafka Consumer查询目标主题所有分区的最新偏移量,以此构建KafkaTrackingToken,Kotlin示例如下:
import org.axonframework.extensions.kafka.eventhandling.tokenstore.KafkaTrackingToken import org.apache.kafka.clients.consumer.KafkaConsumer import org.apache.kafka.common.TopicPartition import java.util.Properties fun createKafkaHeadToken(consumerProps: Properties, topics: List<String>): KafkaTrackingToken { // 创建临时Consumer查询最新偏移量,use块自动关闭资源 KafkaConsumer<String, ByteArray>(consumerProps).use { consumer -> // 获取目标主题的所有分区 val partitions = topics.flatMap { topic -> consumer.partitionsFor(topic).map { TopicPartition(topic, it.partition()) } } // 定位到每个分区的最新位置 consumer.seekToEnd(partitions) // 构建分区与偏移量的映射关系 val partitionOffsets = partitions.associate { tp -> tp to consumer.position(tp) } // 返回Kafka跟踪令牌 return KafkaTrackingToken(partitionOffsets) } }
在处理器配置中使用该方法:
// 初始化Kafka消费者配置 val kafkaConsumerProps = Properties().apply { put("bootstrap.servers", "your-kafka-broker:9092") put("group.id", "your-processor-group-id") // 添加其他必要的消费者配置(如key/value反序列化器等) } // 目标监听主题列表 val targetTopics = listOf("event-topic-1", "event-topic-2") // 构建处理器配置 val config = TrackingEventProcessorConfiguration.forSingleThreadedProcessing() .andInitialTrackingToken { _ -> createKafkaHeadToken(kafkaConsumerProps, targetTopics) }
3. 未知分区偏移量的处理
上述示例无需手动指定分区,通过Kafka Consumer自动查询目标主题的所有分区,并获取每个分区的最新偏移量,完全适配动态分区场景。
4. 多MessageSource的简化处理
如果使用MultiStreamableMessageSource合并多个KafkaMessageSource,只需收集所有消息源对应的主题,一次性传入createKafkaHeadToken方法即可,无需单独处理每个消息源:
// 初始化多个KafkaMessageSource val source1 = KafkaMessageSource.builder<String, ByteArray>() .consumerConfiguration(kafkaConsumerProps) .topic("event-topic-1") .build() val source2 = KafkaMessageSource.builder<String, ByteArray>() .consumerConfiguration(kafkaConsumerProps) .topic("event-topic-2") .build() // 收集所有主题 val allTopics = listOf(source1.topic(), source2.topic()) // 构建统一的初始令牌配置 val config = TrackingEventProcessorConfiguration.forSingleThreadedProcessing() .andInitialTrackingToken { _ -> createKafkaHeadToken(kafkaConsumerProps, allTopics) } // 配置MultiStreamableMessageSource val multiSource = MultiStreamableMessageSource.builder() .addMessageSource(source1) .addMessageSource(source2) .build()
关键注意事项
- 临时Kafka Consumer仅用于查询偏移量,使用完毕自动关闭,不会长期占用连接
- 确保消费者配置中的
group.id与跟踪处理器的分组一致,避免偏移量存储冲突 KafkaTrackingToken记录的是各分区的最新偏移量,处理器启动后会从该位置开始消费新产生的事件
内容的提问来源于stack exchange,提问作者Sergey Bulavkin
相关产品推荐
相关产品推荐

