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