You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.25 03:33:15