【发布时间】: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