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

非安全集群下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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 18:24:32