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

Quarkus中KStream与KTable关联后无消息输出问题

Kafka Streams 无输出问题

背景与问题

基于Quarkus Kafka Streams相关工程构建新服务,核心逻辑为根据关联主题的消息条件输出消息(而非示例中的聚合操作)。因JSON消息Serde使用遇到阻碍,已简化消息结构:

  • t1(KStream):纯文本值
  • t2(KTable):分号分隔的纯文本值
  • t3(KStream):CSV字符串

使用kafkacat并通过-k选项设置键(如k1)发送消息后,t3无任何输出。

核心代码(TopologyProducer.java)

@Produces
public Topology buildTopology() {
  StreamsBuilder builder = new StreamsBuilder();

  ObjectMapperSerde<stream1> stream1 = new ObjectMapperSerde<>(stream1.class);
  ObjectMapperSerde<topic1> topic1 = new ObjectMapperSerde<>(topic1.class);
  ObjectMapperSerde<output1> output1 = new ObjectMapperSerde<>(output1.class);

  GlobalKTable<String, topic1> topic1 = builder.globalTable(
      t2,
      Consumed.with(Serdes.String(), topic1));

  builder.stream(t1,
      Consumed.with(Serdes.String(), stream1))
      .join(t2,
          (paramName, paramValue) -> paramName,
          (paramValue, paramLimits) -> {
            // Add some logic to return conditionally
            return new output1("paramName", 0.0, 0.0, true);
          })
      .to(t3,
          Produced.with(Serdes.String(), output1));
  return builder.build();
}
}

内容的提问来源于stack exchange,提问作者cogitoergosum

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 02:15:40