如何在单个Quarkus Kafka Streams应用中部署多个Topology?
在单个Quarkus应用中运行多个Kafka Streams Topology的解决方案
Quarkus默认的Kafka Streams扩展期望注入单个Topology实例,所以直接通过多个@Produces方法暴露Topology会触发AmbiguousResolutionException。你可以通过以下两种方式实现多拓扑部署:
方法1:使用@Named限定符区分拓扑
通过@Named注解为每个Topology实例标记唯一名称,然后在注入时指定名称,手动创建并管理KafkaStreams实例的生命周期。
步骤1:定义带限定符的拓扑生产者
@ApplicationScoped public class UserTopologyProducer { @Produces @Named("user-topology") public Topology createUserTopology() { Topology topology = new Topology(); // 构建用户相关的拓扑逻辑 topology.addSource("USER_SOURCE", "raw-user-topic") .addProcessor("USER_PROCESSOR", UserProcessor::new, "USER_SOURCE") .addSink("USER_SINK", "processed-user-topic", "USER_PROCESSOR"); return topology; } }
@ApplicationScoped public class OrderTopologyProducer { @Produces @Named("order-topology") public Topology createOrderTopology() { Topology topology = new Topology(); // 构建订单相关的拓扑逻辑 topology.addSource("ORDER_SOURCE", "raw-order-topic") .addProcessor("ORDER_PROCESSOR", OrderProcessor::new, "ORDER_SOURCE") .addSink("ORDER_SINK", "processed-order-topic", "ORDER_PROCESSOR"); return topology; } }
步骤2:创建流管理器统一控制生命周期
@ApplicationScoped public class MultiStreamsManager { private KafkaStreams userStreams; private KafkaStreams orderStreams; @Inject public MultiStreamsManager( @Named("user-topology") Topology userTopology, @Named("order-topology") Topology orderTopology, @ConfigProperty(name = "kafka.bootstrap.servers") String bootstrapServers) { // 配置用户流(注意application.id必须唯一) Properties userConfig = new Properties(); userConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, "user-streams-app"); userConfig.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); userConfig.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); userConfig.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); userStreams = new KafkaStreams(userTopology, userConfig); // 配置订单流 Properties orderConfig = new Properties(); orderConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, "order-streams-app"); orderConfig.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); orderConfig.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); orderConfig.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); orderStreams = new KafkaStreams(orderTopology, orderConfig); } @Startup public void startStreams() { userStreams.start(); orderStreams.start(); } @Shutdown public void stopStreams() { userStreams.close(Duration.ofSeconds(10)); orderStreams.close(Duration.ofSeconds(10)); } }
方法2:手动创建拓扑与KafkaStreams实例
完全不依赖Quarkus的@Produces注入机制,直接在管理器中构建拓扑并初始化KafkaStreams实例,同样手动管理生命周期。
@ApplicationScoped public class ManualMultiStreamsManager { private KafkaStreams userStreams; private KafkaStreams orderStreams; @Inject @ConfigProperty(name = "kafka.bootstrap.servers") String bootstrapServers; @PostConstruct public void initStreams() { // 构建并初始化用户流 Topology userTopology = new Topology(); userTopology.addSource("USER_SOURCE", "raw-user-topic") .addProcessor("USER_PROCESSOR", UserProcessor::new, "USER_SOURCE") .addSink("USER_SINK", "processed-user-topic", "USER_PROCESSOR"); Properties userConfig = new Properties(); userConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, "user-streams-app"); userConfig.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); userConfig.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); userConfig.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); userStreams = new KafkaStreams(userTopology, userConfig); // 构建并初始化订单流 Topology orderTopology = new Topology(); orderTopology.addSource("ORDER_SOURCE", "raw-order-topic") .addProcessor("ORDER_PROCESSOR", OrderProcessor::new, "ORDER_SOURCE") .addSink("ORDER_SINK", "processed-order-topic", "ORDER_PROCESSOR"); Properties orderConfig = new Properties(); orderConfig.put(StreamsConfig.APPLICATION_ID_CONFIG, "order-streams-app"); orderConfig.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); orderConfig.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); orderConfig.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); orderStreams = new KafkaStreams(orderTopology, orderConfig); } @Startup public void start() { userStreams.start(); orderStreams.start(); } @Shutdown public void stop() { userStreams.close(Duration.ofSeconds(10)); orderStreams.close(Duration.ofSeconds(10)); } }
关键注意事项
- 每个
KafkaStreams实例的application.id必须唯一,否则会导致集群中流实例的冲突。 - 必须手动处理流的启动和关闭逻辑,确保应用启停时流实例能正确初始化和销毁。
- 可以通过Quarkus的配置属性(如
application.properties)外部化所有Kafka相关配置,避免硬编码。
内容的提问来源于stack exchange,提问作者Tea Kaa
相关产品推荐
相关产品推荐

