如何基于Kafka原生API实现消费消息时调用远程REST API
基于Kafka原生API的文件逐行生产-消费链路实现方案
链路逻辑拆解
整个消息流转链路拆成3个独立环节,全部基于Kafka官方3.x版本Java API实现,无额外第三方组件依赖:
- 生产环节:逐行读取目标文件内容,每一行封装为独立Kafka消息,发往源业务主题
- 消费处理环节:消费端订阅源业务主题,拉取到消息后调用远程REST接口完成业务逻辑处理
- 结果转发环节:拿到REST接口返回的处理结果后,复用生产者实例将结果发往下游业务主题,供后续链路节点消费
核心实现代码
基础配置说明
以下配置为必填项,参数值替换为实际业务环境的配置即可
- 生产者端配置:
bootstrap.servers(Kafka集群连接地址)、key.serializer/value.serializer(普通文本场景直接用StringSerializer)、acks=all(保证消息多副本落盘后再返回成功,避免丢数) - 消费者端配置:
bootstrap.servers、key.deserializer/value.deserializer(对应使用StringDeserializer)、group.id(消费组ID,同组内节点自动负载均衡)、auto.offset.reset=earliest(首次启动无消费偏移量时从最早消息开始消费)、enable.auto.commit=false(关闭自动提交偏移量,避免消息没处理完就提交导致丢数)
生产者端:文件逐行发送实现
KafkaProducer是线程安全的,全局初始化一次即可,别每发一条消息新建一个实例,会把集群连接打满。逐行读文件用缓冲流,不要一次性把整个文件加载到内存,大文件场景会直接OOM。
import org.apache.kafka.clients.producer.*; import java.io.BufferedReader; import java.io.FileReader; import java.io.IOException; import java.util.Properties; public class FileLineProducer { public static void main(String[] args) { Properties producerProps = new Properties(); producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的Kafka集群地址"); producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); producerProps.put(ProducerConfig.ACKS_CONFIG, "all"); // 全局复用一个生产者实例 KafkaProducer<String, String> producer = new KafkaProducer<>(producerProps); String sourceTopic = "源业务主题名"; String filePath = "待读取的本地文件绝对路径"; try (BufferedReader br = new BufferedReader(new FileReader(filePath))) { String line; Long lineNum = 0L; while ((line = br.readLine()) != null) { lineNum++; if (line.isBlank()) continue; // 空行按业务需求决定是否跳过 // 消息key用行号即可,有业务主键的话替换成业务主键,相同key的消息会落到同一个分区保证顺序 ProducerRecord<String, String> record = new ProducerRecord<>(sourceTopic, lineNum.toString(), line); // 异步发送加失败回调,别直接发了不管,失败了要做重试或者落盘记录 producer.send(record, (metadata, e) -> { if (e != null) { System.err.printf("第%d行发送失败,内容:%s,异常:%s%n", lineNum, line, e.getMessage()); // 生产环境这里加告警、失败重试逻辑 } }); } } catch (IOException e) { throw new RuntimeException("文件读取失败", e); } finally { producer.flush(); // 把缓冲区里没发完的消息全部发出去 producer.close(); } } }
消费者端:消息消费+REST调用+结果转发实现
消费端同样复用生产者实例发下游结果,REST调用一定要设置超时时间,别用默认的无超时配置,不然接口卡住会直接堵死整个消费线程。偏移量一定要等REST调用成功、结果消息发送成功之后再提交,别图省事开自动提交。
import org.apache.kafka.clients.consumer.*; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.common.serialization.StringDeserializer; import java.net.URI; import java.net.http.HttpClient; import java.net.http.HttpRequest; import java.net.http.HttpResponse; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class RestProcessConsumer { public static void main(String[] args) { // 初始化消费者 Properties consumerProps = new Properties(); consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的Kafka集群地址"); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "你的业务消费组ID"); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); KafkaConsumer<String, String> consumer = new KafkaConsumer<>(consumerProps); // 初始化结果转发用的生产者,全局复用 Properties producerProps = new Properties(); producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "你的Kafka集群地址"); producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.StringSerializer"); producerProps.put(ProducerConfig.ACKS_CONFIG, "all"); KafkaProducer<String, String> resultProducer = new KafkaProducer<>(producerProps); String sourceTopic = "源业务主题名"; String resultTopic = "结果下发主题名"; String restApiUrl = "你的远程REST接口地址"; consumer.subscribe(Collections.singletonList(sourceTopic)); try { while (true) { // 拉取消息,拉取间隔按业务延迟要求设置,一般设100-1000ms即可 ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> record : records) { String content = record.value(); try { // 1. 调用REST接口处理业务 String processResult = callRestApi(restApiUrl, content); // 2. 同步发送结果消息,确保发送成功 ProducerRecord<String, String> resultRecord = new ProducerRecord<>(resultTopic, record.key(), processResult); resultProducer.send(resultRecord).get(); // 3. 处理+发送都成功了再提交偏移量 consumer.commitSync(); } catch (Exception e) { System.err.printf("消息处理失败,分区:%d,偏移量:%d,异常:%s%n", record.partition(), record.offset(), e.getMessage()); // 生产环境这里加重试逻辑,重试3次还失败就把消息发到死信主题,别一直卡着这条消息阻塞消费 } } } } finally { consumer.close(); resultProducer.close(); } } private static String callRestApi(String apiUrl, String content) throws Exception { HttpClient client = HttpClient.newHttpClient(); HttpRequest request = HttpRequest.newBuilder() .uri(URI.create(apiUrl)) .timeout(Duration.ofSeconds(10)) // 超时时间按业务接口响应时间设置 .header("Content-Type", "application/json") .POST(HttpRequest.BodyPublishers.ofString(content)) .build(); HttpResponse<String> response = client.send(request, HttpResponse.BodyHandlers.ofString()); if (response.statusCode() != 200) { throw new RuntimeException("接口返回异常,状态码:" + response.statusCode()); } return response.body(); } }
生产环境踩坑提醒
- 所有客户端实例(KafkaProducer、KafkaConsumer、HTTP客户端)全部全局单例复用,不要每次处理消息新建,不然会出现连接泄漏、吞吐量上不去的问题
- 消息发送、REST调用都要加合理的重试机制,重试次数别设太多,一般3次足够,避免异常消息阻塞整个消费链路
- 大文件、高吞吐场景可以调整生产者的
batch.size、linger.ms参数攒批发送,提升发送效率;消费端可以调整max.poll.records控制单次拉取的消息数量,避免单次拉太多处理不过来触发消费组重平衡 - 多次重试仍然失败的消息不要直接丢弃,单独发到死信主题存储,后续人工排查处理,避免数据丢失
- 如果有消息顺序要求,发送时给同业务标识的消息指定相同的key,保证这些消息落到同一个分区,消费端按分区顺序处理即可
内容的提问来源于stack exchange,提问作者Test Admin
相关产品推荐
相关产品推荐

