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

如何从React Js调用Kafka Listener?后端为Springboot

React通过Spring Boot后端向Kafka发送消息的实现方案

首先明确:前端不能直接连接Kafka集群,Kafka本身不提供标准REST接口供前端调用,必须通过你的Spring Boot后端作为中间层,由后端封装Kafka的生产/消费逻辑,前端只需要调用后端提供的API即可。

一、修正你的React前端代码

你当前的代码存在格式错误、请求体不规范、缺少错误处理等问题,以下是修正后的可运行版本:

async apiKafkaProducer(path: string, topicId: string, message: any) {
  const ip = await publicIpv4();
  const headers = {
    'Content-Type': 'application/json',
    Authorization: `Bearer ${localStorage.getItem('token')}`,
    'x-real-ip': `${ip}`,
  };

  // 请求体需包含目标topic和具体消息内容,而非仅传topicId
  const requestBody = {
    topic: topicId,
    content: message
  };

  const requestOptions = {
    method: 'POST',
    headers: headers,
    body: JSON.stringify(requestBody),
  };

  try {
    const response = await fetch(`${process.env.baseUrl}${path}`, requestOptions);
    if (!response.ok) {
      throw new Error(`请求失败: ${response.statusText}`);
    }
    const result = await response.json();
    console.log('消息发送成功:', result);
    return result;
  } catch (error) {
    console.error('发送消息出错:', error);
    throw error; // 向上抛出错误让调用方处理
  }
}

说明:

  • 新增message参数用于传递要发送的具体内容
  • 移除不必要的responseType: 'blob'配置,后端返回JSON格式响应
  • 改用async/await替代链式then,代码可读性更高
  • 完善错误捕获逻辑,处理请求失败场景

二、Spring Boot后端实现

1. 引入Kafka依赖(pom.xml)

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka</artifactId>
</dependency>

2. 配置Kafka生产者(application.yml)

spring:
  kafka:
    bootstrap-servers: your-kafka-broker:9092 # 替换为你的Kafka集群地址
    producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer
      properties:
        spring.json.add.type.headers: false

3. 编写API接口Controller

import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.http.ResponseEntity;
import org.springframework.web.bind.annotation.PostMapping;
import org.springframework.web.bind.annotation.RequestBody;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

@RestController
@RequestMapping("/api/kafka")
public class KafkaProducerController {

    private final KafkaTemplate<String, Object> kafkaTemplate;

    // 构造注入KafkaTemplate
    public KafkaProducerController(KafkaTemplate<String, Object> kafkaTemplate) {
        this.kafkaTemplate = kafkaTemplate;
    }

    @PostMapping("/send")
    public ResponseEntity<String> sendMessage(@RequestBody KafkaMessageRequest request) {
        try {
            // 发送消息到指定topic
            kafkaTemplate.send(request.getTopic(), request.getContent());
            return ResponseEntity.ok("消息已成功发送到Kafka");
        } catch (Exception e) {
            return ResponseEntity.internalServerError().body("发送消息失败: " + e.getMessage());
        }
    }

    // 接收前端请求的DTO类
    public static class KafkaMessageRequest {
        private String topic;
        private Object content;

        // Getter和Setter
        public String getTopic() { return topic; }
        public void setTopic(String topic) { this.topic = topic; }
        public Object getContent() { return content; }
        public void setContent(Object content) { this.content = content; }
    }
}

三、关键注意事项

  • 权限验证:前端传递的Authorization头,后端需通过Spring Security或拦截器验证token有效性,防止未授权请求
  • Topic检查:可在后端添加逻辑,校验目标topic是否存在,不存在时返回明确错误
  • 消息可靠性:若需保证消息不丢失,可配置Kafka生产者acks=all,并开启事务或同步发送
  • 异常处理:后端需捕获Kafka发送过程中的异常,返回友好的错误信息给前端

内容的提问来源于stack exchange,提问作者Yasiru Dahanayaka

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 06:25:33