【问题标题】:Reactive SQS Listener反应式 SQS 监听器
【发布时间】: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


    【解决方案1】:

    您可以正确猜到,findById(1) 操作没有时间完成。原因是.subscribe 启动进程并立即返回。

    我无法从函数返回 Mono,因为 process 方法 应该是通用的

    好吧,如果你想使用响应式存储库,你必须返回Mono

     Mono<Boolean> process(String message, Map<String, String> attributes) {
        return reactiveCrudRepository
                 .existsById(1); //consider using existsById instead of findById
      }
    

    然后像这样使用它:

     .filterWhen(messageFlux -> process(...))
     .flatMap(message -> sqsAsyncClient.deleteMessage(
                    DeleteMessageRequest.builder()
                        .queueUrl(queueUrl)
                        .receiptHandle(message.receiptHandle())
                        .build()))
      .subscribe();
    

    编辑

    我总是必须返回该数据类型的 Mono,但数据可以是 不同类型的事件不同。

    您可以完全按照使用简单 Java 方法的方式进行操作。假设您有两种类型的事件:PlainEventAdvancedEvent 和父类Event

    class Event {
        private String id;
    }
    
    class PlainEvent extends Event {
        private String name;
    }
    
    class AdvancedEvent extends Event {
        private String name;
    }
    

    返回类型为Mono&lt;Event&gt;:

    Mono<Event> process() {
        return Mono.just(new PlainEvent("id", "name"));
    }
    
    .flatMap(event -> process())
    .flatMap(event -> {
        if (event instanceof PlainEvent) {
            log.info("plain {}", event);
        } else {
            log.info("advanced {}", event);
        }
    
        return reactiveCrudRepository.save(event);
    })
    

    【讨论】:

    • 但是如果我想在事件上使用反应式存储库存储一些数据,我该怎么做?我总是必须返回该数据类型的 Mono,但不同类型的事件的数据可能不同。
    • @Professor 编辑完成
    • 非常感谢@lkatiforis
    猜你喜欢
    • 2019-09-26
    • 2021-10-15
    • 2019-08-05
    • 2021-08-18
    • 2018-10-18
    • 1970-01-01
    • 2011-06-19
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多