【问题标题】:Frontend for spring kafka using angular使用角度的弹簧卡夫卡前端
【发布时间】:2021-12-11 21:38:15
【问题描述】:

我有一个消费者使用我的 kafka 集群中的 spring 来使用数据。现在我想使用 Angular 创建 UI。我已经写完了,但是我在后端(java)中得到了这个错误

我知道 kafka 不适用于多线程,当我尝试在浏览器(前端)中调用 UI 时会发生此错误

这是我的 service.ts 代码,将自身与 spring 连接起来

  private usersUrl: string;

  constructor(private http: HttpClient) {
  this.usersUrl = 'http://localhost:8080/receiptsData';
 }

  public findAll(): Observable<ReceiptData[]> {
  return this.http.get<ReceiptData[]>(this.usersUrl, { responseType: 'json' });
 }

在我的后端,我还将它添加到我的消费者以避免访问阻塞

@CrossOrigin(origins = "http://localhost:4200")

如果我不能从其他来源(角度、浏览器)对我的 kafka 进行多线程访问,我可以访问数据的其他可能性是什么?

  • 对不起我的英语,我尽力解释了

编辑:

我在我的 java 中使用线程,因为我不想使用 while(true) 这是我的后端 Controller Java 代码:

public class ExampleKafkaConsumerController extends 
ShutdownableThread {

private final KafkaConsumer<Integer, String> consumer;
private static final String RECEIVED_MESSAGE = "Received message: (";
private static final long DURATION = 5000;

public ExampleKafkaConsumerController() {

  super(DataUtils.GROUP_ID, false);
  ExampleKafkaConsumerConfig exKafkaConsumerConfig = new 
  ExampleKafkaConsumerConfig();
  consumer = exKafkaConsumerConfig.kafkaConfiguration();
}

@GetMapping("/receiptsData")
@Override
public void doWork() {
  
  consumer.subscribe(Collections.singletonList(DataUtils.TOPIC_NAME));
    ConsumerRecords<Integer, String> records = 
      consumer.poll(Duration.ofMillis(DURATION));
    System.out.println(records.count());
    for (ConsumerRecord<Integer, String> record : records) {
      System.out.println(RECEIVED_MESSAGE + record.value());
      KafkaJsonConverter kafkaJsonConverter = new 
        KafkaJsonConverter(record.value());
      System.out.println(kafkaJsonConverter.
        convertStringToJsonObject().toString());
    }
  }
}

这是我的 Main 类:

@SpringBootApplication
public class ExampleConsumerRunner implements CommandLineRunner {

  public static void main(String[] args) {
   SpringApplication.run(ExampleConsumerRunner.class, args);
 }

 @Autowired
 ExampleKafkaConsumerController exKafkaConsumerController;

 @Override
 public void run(String... args) {
   try {
     exKafkaConsumerController.start();
  } catch (Exception ex) {
     System.out.println(ex.getMessage());
  }
}

【问题讨论】:

  • 错误来自您的服务器,与您的 UI 或您如何访问它无关。您的 Java 代码中的某些内容正在使用多个线程,因此您需要显示该代码
  • 你是对的,我在我的控制器中使用线程,因为我想避免'While(true)',你能看看我的java代码
  • 如果你真的想将消费者暴露为 HTTP 端点azkarrastreams.io,我建议从这个库开始,否则,Kafka REST 代理已经为你解决了这个问题,你不需要编写任何代码

标签: java angular spring apache-kafka


【解决方案1】:

根据示例,您正尝试通过 GET 请求从 Kafka 主题接收最新条目。您显示的错误是因为您试图在控制器的 doWork() 方法中使用单个使用者实例。要解决这个特定问题,您应该为每个请求(在 doWork() 中)拥有自己的消费者实例,但我认为您可能想要重新审视解决问题的基本方法,而不是这样做。

原因是将端点直接连接到主题订阅不是很典型 - 因为您可能不希望每个请求单独连接到 Kafka(尤其是具有相同的 groupId),这种方法有很多缺点,并且不会导致高性能或确定性 API。

我会推荐你​​有的设计 a) 服务器端数据存储(单例、内存数据库或持久数据库,取决于您的需要) b) 服务器端消费者订阅主题并将记录写入数据库。 Spring @KafkaListener 很容易做到这一点。 c) REST 控制器读取步骤 a) 中完成的数据存储

通过这种设计,您可以为主题设置一个侦听器,将 API 与侦听器分离,并允许您对 API 实施例如查询过滤器(最新的,或自日期以来的条目或偏移量等)。

【讨论】:

  • 我移动了subscribe to topic 并根据@OneCricketeer 所说的删除了我的CommandLineRunner ,这个问题解决了,但它在Angular 上读取空数据,但在Java 中正确读取
  • @Space Angular 不是问题。使用 cURL 或 Postman 测试您的 API。具体来说,您的 GetMapping 不会返回 Response 对象。打印只到服务器控制台,而不是 HTTP 服务器响应......此外,资源类不应该是线程对象; Springboot 会自动为您处理
  • @OneCricketeer 你能把这个写成答案,这样我就可以接受它作为解决方案,非常感谢
  • @Space 嗯,你的问题解决了吗?
  • 是的,我的问题解决了,因为我忘记了 GetMapping 来返回响应对象
猜你喜欢
  • 2020-02-08
  • 2022-11-28
  • 2020-02-22
  • 2018-11-08
  • 1970-01-01
  • 1970-01-01
  • 2020-05-22
  • 1970-01-01
  • 2016-07-22
相关资源
最近更新 更多