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

Java多线程Kafka生产者实现方案及线程安全问题咨询

Java多线程Kafka生产者实现方案及线程安全问题咨询

我有一个Kafka生产者,通过下方的sendMessage方法发送消息。该方法接收主题名、分组名、键和消息载荷,然后处理并发布JSON数组中的多条消息:

public void sendMessage(String topicName, String groupName, String key, String msg) throws Exception {
    if (msg != null) {
        JSONObject jsonReq = new JSONObject(msg);
        
        if (key == null || topicName == null || "".equals(key) || "".equals(topicName)) {
            throw new CustomException("Invalid key/topicname.");
        }
        if (jsonReq.has(groupName)) {
            JSONArray messages = jsonReq.getJSONArray(groupName);
            if (messages.length() == 0) {
                log.info("No messages found in groupName: {}. Skipping processing.", groupName);
                return;
            }

            List<String> messageIds = new ArrayList<>();
            AtomicBoolean errorFlag = new AtomicBoolean(false);
            long startTime = System.currentTimeMillis();

            for (int i = 0; i < messages.length(); i++) {
                JSONObject obj = messages.getJSONObject(i);
                if (obj.has(key) && StringUtils.isNotBlank(obj.getString(key))) {
                    ProducerRecord<String, String> producerRecord =
                        new ProducerRecord<>(topicName, obj.getString(key), obj.toString());

                    producer.send(producerRecord, new CRKafkaCallBackHandler(producerRecord, errorFlag));
                    messageIds.add(topicName + "\t" + obj.getString(key));
                } else {
                    throw new CustomException("Mandatory property '" + key + "' is missing in message: " + obj.toString());
                }
            }

            long endTime = System.currentTimeMillis();
            log.info("Total Messages: {} For Topic: {} Time Taken: {} ms", messages.length(), topicName, (endTime - startTime));

            if (errorFlag.get()) {
                throw new CustomException("Failed to publish one or more messages to Kafka");
            }
        } else {
            throw new CustomException("Invalid groupName found in request.");
        }
    } else {
        throw new CustomException("Received a NULL or invalid message.");
    }
    log.info("Message processing completed successfully For Topic: {}", topicName);
}

目前这个方法是单线程运行的,我们遇到了超时异常:

org.apache.kafka.common.errors.TimeoutException: Expiring 18 record(s) 120001 ms has passed since batch creation

我猜测这是高负载下的瓶颈导致的。我想通过实现多线程Kafka生产者来提升吞吐量,现在有几个疑问:

  • 实现Java多线程Kafka生产者的最佳方式是什么?应该创建一个固定大小的线程池,提交sendMessage任务吗?还是有更适配Kafka的方案?
  • 如何保证线程安全?由于Kafka的send()方法是异步的,我需要担心生产者实例的并发问题吗?

备注:内容来源于stack exchange,提问作者World of Titans

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 11:27:58