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

如何基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 02:18:19