非安全集群下Spring集成Kafka动态Producer构造报错问题
问题根因
你的Kafka Producer构造失败核心是序列化器配置错误,同时存在Bean重复定义、Producer用法不符合规范的问题,和Kafka服务本身连通性无关。
具体错误点
- Key泛型与序列化器不匹配:你声明的Producer泛型为
Producer<String, KafkaMessage>,Key类型是String,但配置中KEY_SERIALIZER_CLASS_CONFIG指定的是LongSerializer,类型不匹配。 - Value序列化器配置非法:Kafka要求
VALUE_SERIALIZER_CLASS_CONFIG必须传入实现org.apache.kafka.common.serialization.Serializer接口的类,你直接传入了普通POJO类KafkaMessage.class,该类没有实现序列化接口,Kafka反射实例化序列化器时直接抛出异常,这是构造失败的直接原因。从日志打印的value.serializer = class org.kafka.producer.KafkaMessage也能印证这个错误。 - 冗余Bean定义+不规范的Producer生命周期管理:你同时手动定义了原型作用域的原生KafkaProducer、Spring封装的ProducerFactory和KafkaTemplate,两套配置混用容易产生冲突;且每次发送消息都新建Producer、发送完成立即关闭,会产生大量TCP连接创建销毁开销,而Kafka Producer本身是线程安全的,完全可以全局单例复用。
修复方案
1. 实现Value的序列化器
你可以自己实现Kafka的Serializer接口,完成KafkaMessage对象到字节数组的转换:
import org.apache.kafka.common.serialization.Serializer; import com.fasterxml.jackson.databind.ObjectMapper; public class KafkaMessageSerializer implements Serializer<KafkaMessage> { private final ObjectMapper objectMapper = new ObjectMapper(); @Override public byte[] serialize(String topic, KafkaMessage data) { if (data == null) { return null; } try { return objectMapper.writeValueAsBytes(data); } catch (Exception e) { throw new RuntimeException("序列化KafkaMessage失败", e); } } }
如果不想自定义实现,也可以直接用Spring Kafka自带的JsonSerializer,不需要额外写代码。
2. 修正配置类,删除冗余Bean
删掉手动定义的原生KafkaProducer Bean,统一使用Spring Kafka提供的组件,修正序列化器配置:
@Configuration public class KafkaProducerConfig { @Autowired private Kafka kafkaConfig; @Bean public ProducerFactory<String, KafkaMessage> producerFactory() { Map<String, Object> configs = new HashMap<>(); configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaConfig.getBootstrapServers()); // 修正Key序列化器:Key为String类型,使用StringSerializer configs.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); // 修正Value序列化器:使用自定义的KafkaMessageSerializer configs.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, KafkaMessageSerializer.class.getName()); // 如果用自带JsonSerializer,替换上面一行为以下配置即可 // configs.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class.getName()); // configs.put(JsonSerializer.TYPE_MAPPING, "kafkaMessage:org.kafka.producer.KafkaMessage"); configs.put(ProducerConfig.CLIENT_ID_CONFIG, "ABC"); configs.put(ProducerConfig.RETRIES_CONFIG, 3); configs.put(ProducerConfig.LINGER_MS_CONFIG, 5); configs.put(ProducerConfig.ACKS_CONFIG, "all"); return new DefaultKafkaProducerFactory<>(configs); } @Bean public KafkaTemplate<String, KafkaMessage> kafkaTemplate() { return new KafkaTemplate<>(producerFactory()); } }
3. 简化消息发送逻辑
移除原生Producer的注入,直接使用KafkaTemplate发送消息,由Spring管理Producer的生命周期,避免重复创建销毁连接:
@Service public class MessageService { private static final Logger LOGGER = LoggerFactory.getLogger(MessageService.class); private final KafkaTemplate<String, KafkaMessage> kafkaTemplate; public MessageService(KafkaTemplate<String, KafkaMessage> kafkaTemplate) { this.kafkaTemplate = kafkaTemplate; } public String publishSync(final @NotEmpty String topicName, final @NotNull String message) { return publishRawMessage(topicName, message, false); } private @NotNull String publishRawMessage(final @NotEmpty String topicName, final @NotNull String message, final boolean async) { KafkaMessage kafkaMessage = new KafkaMessage(); kafkaMessage.setId(UUID.randomUUID().hashCode()); kafkaMessage.setMessage(message); ProducerRecord<String, KafkaMessage> record = new ProducerRecord<>(topicName, kafkaMessage); try { if (async) { kafkaTemplate.send(record, (metadata, exception) -> { if (exception != null) { LOGGER.warn("异步发送消息失败", exception); return; } LOGGER.info("Record ASYNCHRONOUSLY sent to partition {} with offset {}", metadata.partition(), metadata.offset()); }).get(); } else { kafkaTemplate.send(record).get(); } } catch (ExecutionException e) { LOGGER.warn("发送消息执行失败", e); } catch (InterruptedException e) { LOGGER.warn("发送消息被中断", e); Thread.currentThread().interrupt(); } return "Successfully Inserted"+message; } }
内容的提问来源于stack exchange,提问作者utkarsh saraf
相关产品推荐
相关产品推荐

