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

如何在Spring Boot中通过Kafka Streams的KStream读取Kafka主题JSON值?

解决方案:Kafka Streams读取JSON主题并提取字段

先解决反序列化错误(State transition from PENDING_ERROR to ERROR)

你遇到的错误大概率是键的反序列化不匹配导致的:当前代码指定键为Long类型,但实际Kafka主题中的键可能是String、null或者其他类型,和Serdes.Long()不兼容,触发序列化异常。

第一步:排查键的类型

先临时将键和值都用Serdes.String()读取,确认流能正常启动:

KStream<String, String> kStream = streamsBuilder.stream("TopicName", Consumed.with(Serdes.String(), Serdes.String()));

如果这样能正常运行,说明之前的键类型配置错误,需要根据生产者实际发送的键类型调整Serdes(比如生产者用StringSerializer就用Serdes.String(),用LongSerializer就保留Serdes.Long())。

高效提取JSON字段的两种方法

方法1:POJO + Jackson Serde(推荐,类型安全)

这种方法适合需要频繁操作JSON字段的场景,通过POJO映射JSON结构,操作更高效。

1.1 创建对应JSON的POJO类

public class WorkbookStatus {
    private Long Id;
    private String workbook;
    private String state;

    // 必须保留无参构造器(Jackson反序列化需要)
    public WorkbookStatus() {}

    // Getter & Setter
    public Long getId() { return Id; }
    public void setId(Long id) { Id = id; }
    public String getWorkbook() { return workbook; }
    public void setWorkbook(String workbook) { this.workbook = workbook; }
    public String getState() { return state; }
    public void setState(String state) { this.state = state; }
}

1.2 配置Jackson Serde

在Spring Boot环境中,可以直接使用JsonSerde(spring-kafka提供)或者自定义Serde:

// 自定义Jackson Serde
Serde<WorkbookStatus> workbookStatusSerde = Serdes.serdeFrom(
    new JsonSerializer<>(),
    new JsonDeserializer<>(WorkbookStatus.class)
);

1.3 读取并操作POJO流

确认键类型后,直接读取为POJO类型的KStream,之后可以直接提取字段:

// 假设键是Long类型,如果是String就换成Serdes.String()
KStream<Long, WorkbookStatus> kStream = streamsBuilder.stream("TopicName", 
    Consumed.with(Serdes.Long(), workbookStatusSerde));

// 提取state字段示例
kStream.mapValues(WorkbookStatus::getState)
       .foreach((key, state) -> System.out.println("任务状态:" + state));

方法2:直接用Jackson解析String值(适合临时提取字段)

如果不需要完整POJO,可直接用ObjectMapper解析JSON字符串:

ObjectMapper objectMapper = new ObjectMapper();
KStream<String, String> kStream = streamsBuilder.stream("TopicName", Consumed.with(Serdes.String(), Serdes.String()));

kStream.mapValues(value -> {
    try {
        JsonNode jsonNode = objectMapper.readTree(value);
        // 提取需要的字段
        Long id = jsonNode.get("Id").asLong();
        String workbook = jsonNode.get("workbook").asText();
        String state = jsonNode.get("state").asText();
        return String.format("ID: %d, 工作簿: %s, 状态: %s", id, workbook, state);
    } catch (JsonProcessingException e) {
        // 处理解析失败的消息:可记录日志或跳过
        log.error("JSON解析失败: {}", value, e);
        return null;
    }
})
.filter((key, value) -> value != null) // 过滤解析失败的消息
.foreach((key, value) -> System.out.println(value));

额外排查要点

  • 查看Kafka Streams的详细日志,找到SerializationException的具体栈跟踪,确认是键还是值的反序列化问题。
  • 确保生产者发送消息时的序列化器和消费者的Serdes一致:比如生产者键用StringSerializer,消费者就必须用Serdes.String()。
  • Spring Boot中如果配置了全局默认Serde,要注意和Consumed.with()指定的Serdes不冲突。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 19:56:19