【问题标题】:Reactive SQS Poller with backpressue具有背压的反应式 SQS 轮询器
【发布时间】:2018-11-25 19:03:29
【问题描述】:

尝试创建 SQS 轮询器:

  • 进行指数轮询(如果队列中没有消息,则减少请求数量)
  • 如果队列中有大量消息,则更频繁地查询 SQS
  • 如果收到一定数量的消息会产生背压,它会停止轮询
  • 不受 AWS API 速率限制的限制

作为一个例子,我正在使用this JavaRx 实现,它很容易转换为 Project Reactor 并通过背压丰富它。

private static final Long DEFAULT_BACKOFF = 500L;
private static final Long MAX_BACKOFF = 8000L;
private static final Logger LOGGER = LoggerFactory.getLogger(SqsPollerService.class);
private static volatile boolean stopRequested;

public Flux<Message> pollMessages(GetQueueUrlResult q)
{
    return Flux.create(sink -> {
        long backoff = DEFAULT_BACKOFF;

        while (!stopRequested)
        {
            if (sink.isCancelled())
            {
                sink.error(new RuntimeException("Stop requested"));
                break;
            }

            Future<ReceiveMessageResult> future = sink.requestedFromDownstream() > 0
                    ? amazonSQS.receiveMessageAsync(createRequest(q))
                    : completedFuture(new ReceiveMessageResult());

            try
            {
                ReceiveMessageResult result = future.get();

                if (result != null && !result.getMessages().isEmpty())
                {
                    backoff = DEFAULT_BACKOFF;

                    LOGGER.info("New messages found in queue size={}", result.getMessages().size());

                    result.getMessages().forEach(m -> {
                        if (sink.requestedFromDownstream() > 0L)
                        {
                            sink.next(m);
                        }
                    });
                }
                else
                {
                    if (backoff < MAX_BACKOFF)
                    {
                        backoff = backoff * 2;
                    }

                    LOGGER.debug("No messages found on queue.  Sleeping for {} ms.", backoff);

                    // This is to prevent rate limiting by the AWS api
                    Thread.sleep(backoff);
                }
            }
            catch (InterruptedException e)
            {
                stopRequested = true;
            }
            catch (ExecutionException e)
            {
                sink.error(e);
            }
        }
    });
}

实施似乎有效,但有几个问题:

  • 看起来可以使用 Reactor Primitives 在循环中查询 Future 结果,使用 Flux.generate 进行了尝试,但无法控制对 SqsClient 的异步调用次数
  • 如果使用Flux.interval 方法,不了解如何正确实施退避策略
  • 不喜欢Thread.sleep打电话有什么想法怎么替换吗?
  • 如何在取消信号的情况下正确停止循环?现在使用sink.error 来涵盖该案例。

【问题讨论】:

  • SQS 不需要退避,看来您可能误解了某些 SQS 行为。将 WaitTimeSeconds 设置为最大值 (20),将 MaxNumberOfMessages 设置为最大值 (10)。如果队列中有任何消息,SQS 将立即返回最多 10 条消息,否则将等待 1 条消息到达并立即返回(如果它们可能 > 1非常非常接近地到达)。如果没有到达 20 秒,则响应返回空。这样,您可以立即收到消息,而一个完全空闲的队列每小时只会轮询 180 次。
  • 我了解 SQS 部分并计划对生产代码使用“长轮询”,但仍然对如何使用响应式方法解决此类问题感兴趣。

标签: java amazon-sqs project-reactor


【解决方案1】:

您对以下解决方案有何看法:

    private static final Integer batchSize = 1;
    private static final Integer intervalRequest = 3000;
    private static final Integer waitTimeout = 10;
    private static final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(1);

    private static final SqsAsyncClient sqsAsync =
       SqsAsyncClient
         .builder()
         .endpointOverride(URI.create(queueUrl))
         .build();

    public static Flux<Message> sqsPublisher =
        Flux.create(sink -> {
                if (sink.isCancelled()) {
                    sink.error(new RuntimeException("Stop requested"));
                }

            scheduler.scheduleWithFixedDelay(() -> {
                long numberOfRequests = Math.min(sink.requestedFromDownstream(), batchSize);
                if (numberOfRequests > 0) {
                    ReceiveMessageRequest request = ReceiveMessageRequest
                            .builder()
                            .queueUrl(queueUrl)
                            .maxNumberOfMessages((int) numberOfRequests)
                            .waitTimeSeconds(waitTimeout).build();

                    CompletableFuture<ReceiveMessageResponse> response = sqsAsync.receiveMessage(request);

                    response.thenApply(responseValue -> {
                        if (responseValue != null && responseValue.messages() != null && !responseValue.messages().isEmpty()) {
                            responseValue.messages().stream().limit(numberOfRequests).forEach(sink::next);
                        }
                        return responseValue;
                    });

                }
            }, intervalRequest, intervalRequest, TimeUnit.MILLISECONDS);
        });

【讨论】:

    猜你喜欢
    • 2022-01-13
    • 2016-09-20
    • 2021-11-28
    • 2018-08-24
    • 1970-01-01
    • 1970-01-01
    • 2021-06-07
    • 2022-01-13
    • 1970-01-01
    相关资源
    最近更新 更多