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

Kafka消息发送报错:No serializer found for class org.apache.kafka.connect.data.Struct

解决方案:Kafka JSON序列化Struct类失败问题

核心问题是Jackson无法自动序列化org.apache.kafka.connect.data.Struct类,因为它不是标准JavaBean结构,没有Jackson可识别的序列化逻辑。以下是可行的解决办法:

方案1:自定义Jackson序列化器处理Struct

编写自定义序列化器,将Struct的字段和值转换为JSON对象:

import com.fasterxml.jackson.core.JsonGenerator;
import com.fasterxml.jackson.databind.JsonSerializer;
import com.fasterxml.jackson.databind.SerializerProvider;
import org.apache.kafka.connect.data.Struct;
import java.io.IOException;

public class StructSerializer extends JsonSerializer<Struct> {
    @Override
    public void serialize(Struct struct, JsonGenerator jsonGenerator, SerializerProvider serializerProvider) throws IOException {
        jsonGenerator.writeStartObject();
        // 遍历Struct的所有字段,写入JSON键值对
        struct.schema().fields().forEach(field -> {
            try {
                jsonGenerator.writeObjectField(field.name(), struct.get(field));
            } catch (IOException e) {
                throw new RuntimeException("序列化Struct字段失败", e);
            }
        });
        jsonGenerator.writeEndObject();
    }
}

修改EventsDTO,给dataAfterWrite字段指定自定义序列化器:

@Data
public class EventsDTO {
    @JsonSerialize(using = StructSerializer.class)
    private Struct dataAfterWrite;
    private String operation;
}

方案2:提前将Struct转换为Map(简单直接)

在发送消息前,把Struct转换为普通Map,避免直接序列化非标准类:

import java.util.Map;
import java.util.stream.Collectors;

public void sendMessage(EventsDTO data) {
    // 将Struct转换为键值对Map
    Struct struct = data.getDataAfterWrite();
    Map<String, Object> structMap = struct.schema().fields().stream()
            .collect(Collectors.toMap(
                field -> field.name(),
                field -> struct.get(field)
            ));
    
    // 用Map替换原DTO中的Struct字段(也可以直接修改原DTO的字段类型为Map)
    MapEventsDTO mapDto = new MapEventsDTO();
    mapDto.setDataAfterWrite(structMap);
    mapDto.setOperation(data.getOperation());
    
    Message<MapEventsDTO> message = MessageBuilder
            .withPayload(mapDto)
            .setHeader(KafkaHeaders.TOPIC, "events-ingestion-topic")
            .build();
    kafkaTemplate.send(message);
    log.info("message send to kafka topic {}", message);
}

// 对应Map类型的DTO
@Data
public class MapEventsDTO {
    private Map<String, Object> dataAfterWrite;
    private String operation;
}

方案3:禁用Jackson空Bean序列化检查(不推荐)

此方法会跳过无法序列化的对象,但会丢失Struct的所有数据,仅适合临时测试:

在application.properties中添加配置:

spring.jackson.serialization.fail-on-empty-beans=false

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 19:58:29