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

创建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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 01:45:04