KStream左连接异常:仅双边有消息时生成记录,左流单独无输出
Kafka Streams左连接未按预期触发输出问题
我用Kafka Streams编写了左连接代码,预期healthStream收到消息时就生成一条记录,但实际仅当两个流都有消息时才会生成记录。调整JoinWindow时长也没改变结果。
val cncStream = streamsBuilder.stream(config.cncManagerTopic, Consumed.with<String, AgentMessage>(Serdes.String(), AgentMessageSerde())) val healthStream = streamsBuilder.stream(config.healthTopic, Consumed.with<String, AgentMessage>(Serdes.String(), AgentMessageSerde())) healthStream.leftJoin(cncStream, { healthMsg, logsMessage -> Agent( healthMsg.generateKey(), healthMsg.context.customerId, healthMsg.context.agentId, healthMsg.healthMsg.status.toString(), healthMsg.healthMsg.message, LocalDateTime.now(), healthMsg.healthMsg.metrics, LogStatus((logsMessage?.logsUploadedMsg?.ok == true), logsMessage?.logsUploadedMsg?.error.orEmpty(), logsMessage?.logsUploadedMsg?.path.orEmpty()) ) }, JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(5)), StreamJoined.`as`(KafkaAgentsConfiguration.GetStoreName(config.agentsTopic))) .to(config.agentsTopic, Produced.with(Serdes.String(), AgentSerde()))
这不是Bug,是你对Kafka Streams流-流窗口左连接的行为理解有误,具体原因和解决方案如下:
流-流窗口左连接的特性:
流-流之间的连接必须通过窗口限定匹配时间范围,左流消息进入窗口后,会等待窗口内右流的同key消息:- 只要右流有匹配key的消息进入窗口,就会立刻输出连接后的记录;
- 只有当窗口关闭时,左流消息仍未匹配到右流的同key消息,才会输出一条左流记录(右值为null)。
你看到的"仅当两个流都有消息时才输出",是因为窗口还未关闭,左流单独的消息还没触发输出逻辑。
符合你预期的解决方案:
如果你希望左流收到消息就立刻输出(不管右流有没有对应消息),应该用KStream与KTable的左连接,而非流-流窗口连接:
KTable是基于key的状态存储,会保留每个key的最新记录,流-表左连接时,左流每条消息进来会立刻与表中当前同key的记录(无匹配则为null)连接并输出。调整后的代码示例:
val cncTable = streamsBuilder.stream(config.cncManagerTopic, Consumed.with<String, AgentMessage>(Serdes.String(), AgentMessageSerde())) .groupByKey() .reduce { _, newValue -> newValue } // 保留每个key的最新消息作为KTable记录 val healthStream = streamsBuilder.stream(config.healthTopic, Consumed.with<String, AgentMessage>(Serdes.String(), AgentMessageSerde())) healthStream.leftJoin(cncTable, { healthMsg, logsMessage -> Agent( healthMsg.generateKey(), healthMsg.context.customerId, healthMsg.context.agentId, healthMsg.healthMsg.status.toString(), healthMsg.healthMsg.message, LocalDateTime.now(), healthMsg.healthMsg.metrics, LogStatus((logsMessage?.logsUploadedMsg?.ok == true), logsMessage?.logsUploadedMsg?.error.orEmpty(), logsMessage?.logsUploadedMsg?.path.orEmpty()) ) }) .to(config.agentsTopic, Produced.with(Serdes.String(), AgentSerde()))
内容的提问来源于stack exchange,提问作者Amit Eliav
相关产品推荐
相关产品推荐

