创建Kafka Stream按名称Key统计消息条数遇到输出异常如何解决
问题原因
- 冗余的map操作:统计消息条数无需提取value中的amount字段,该操作不影响统计结果,但属于无效逻辑。
- 序列化匹配错误:你当前使用
Serdes.Long()序列化统计得到的Long类型计数值,该序列化器会将数字转为二进制字节存储,如果你用字符串反序列化器读取输出topic的value,就会得到乱码。 - 输出格式不符合预期:你的需求是输出JSON结构的统计结果,当前直接输出Long类型的值,完全不满足格式要求。
修复方案
根据你需要的输出格式选择对应修复方式:
方案1:输出指定JSON格式
输出{"key":"xxx", "count": xx}格式
修改拓扑代码,将计数结果转为JSON结构后用JSON Serde输出:
public Topology createTopology(){ StreamsBuilder builder = new StreamsBuilder(); // json Serde final Serializer<JsonNode> jsonSerializer = new JsonSerializer(); final Deserializer<JsonNode> jsonDeserializer = new JsonDeserializer(); final Serde<JsonNode> jsonSerde = Serdes.serdeFrom(jsonSerializer, jsonDeserializer); ObjectMapper objectMapper = new ObjectMapper(); KStream<String, JsonNode> textLines = builder.stream("bank-transactions", Consumed.with(Serdes.String(), jsonSerde)); KTable<String, Long> transactionCounts = textLines // 直接按key分组计数,不需要处理value .groupByKey(Serialized.with(Serdes.String(), jsonSerde)) .count(); // 将计数结果转为要求的JSON格式 transactionCounts.toStream() .map((key, count) -> { ObjectNode result = objectMapper.createObjectNode(); result.put("key", key); result.put("count", count); return KeyValue.pair(key, result); }) .to("person-transaction-frequency", Produced.with(Serdes.String(), jsonSerde)); return builder.build(); }
如果要输出{"Mike": 2}格式的JSON,只需要调整map里的构造逻辑即可:
.map((key, count) -> { ObjectNode result = objectMapper.createObjectNode(); result.put(key, count); return KeyValue.pair(key, result); })
方案2:仅解决乱码,输出可读数字字符串
如果不需要JSON结构,只是解决乱码问题,将Long类型的计数转成字符串输出即可:
public Topology createTopology(){ StreamsBuilder builder = new StreamsBuilder(); // json Serde final Serializer<JsonNode> jsonSerializer = new JsonSerializer(); final Deserializer<JsonNode> jsonDeserializer = new JsonDeserializer(); final Serde<JsonNode> jsonSerde = Serdes.serdeFrom(jsonSerializer, jsonDeserializer); KStream<String, JsonNode> textLines = builder.stream("bank-transactions", Consumed.with(Serdes.String(), jsonSerde)); KTable<String, Long> transactionCounts = textLines .groupByKey(Serialized.with(Serdes.String(), jsonSerde)) .count(); // 转成字符串输出,避免乱码 transactionCounts.toStream() .mapValues(String::valueOf) .to("person-transaction-frequency", Produced.with(Serdes.String(), Serdes.String())); return builder.build(); }
内容的提问来源于stack exchange,提问作者ks_on_v
相关产品推荐
相关产品推荐

