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

如何在单个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 12:43:31