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
相关产品推荐
相关产品推荐

