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

