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
相关产品推荐
相关产品推荐

