【发布时间】:2022-01-13 21:20:15
【问题描述】:
我正在实现一个反应式 SQS 侦听器,但我遇到了嵌套订阅问题。
这就是我的听众的样子
@PostConstruct
public void listener() {
Mono<ReceiveMessageResponse> receiveMessageResponseMono =
Mono.fromFuture(
() ->
sqsAsyncClient.receiveMessage(
ReceiveMessageRequest.builder()
.queueUrl(queueUrl)
.maxNumberOfMessages(MAX_NUMBER_OF_MESSAGES)
.waitTimeSeconds(WAIT_TIME_IN_SECONDS)
.build()));
receiveMessageResponseMono
.repeat()
.retry()
.map(ReceiveMessageResponse::messages)
.map(Flux::fromIterable)
.flatMap(messageFlux -> messageFlux)
.subscribe(
message -> {
boolean isProcessed;
isProcessed = process(message.body(), message.attributesAsStrings());
if (isProcessed) {
sqsAsyncClient.deleteMessage(
DeleteMessageRequest.builder()
.queueUrl(queueUrl)
.receiptHandle(message.receiptHandle())
.build());
}
});
}
及处理方法:
Boolean process(String message, Map<String, String> attributes) {
reactiveCrudRepository
.findById(1)
.subscribe(b -> log.info("Fetched item {}", b.getId()));
return true;
}
现在的问题是我需要来自反应式存储库的数据,而父订阅不会等待子订阅完成,而且大多数时候我都没有收到此日志Fetched item。我无法从函数返回 Mono,因为 process 方法应该是通用的,并且可以处理任何事件,并且在 process 方法中,我需要使用响应式 rest-client 和 DB 客户端。
我是响应式编程的新手,所以我不确定如何解决这个问题,或者我是否遵循反模式。
【问题讨论】:
标签: java rx-java spring-webflux project-reactor reactive-streams