【问题标题】:How to balance multiple message queues如何平衡多个消息队列
【发布时间】:2019-06-19 22:52:58
【问题描述】:

我有一项可能需要长时间运行(数小时)的任务。该任务由从消息队列(在我的情况下为 AWS SQS)读取的多个工作人员(在我的情况下为 AWS ECS 实例)执行。我有多个用户将消息添加到队列中。问题是,如果 Bob 向队列中添加 5000 条消息,足以让工作人员忙 3 天,然后 Alice 出现并想要处理 5 个任务,Alice 需要等待 3 天才能开始 Alice 的任何任务。

我想在 Alice 提交任务后立即以相等的速率从 Alice 和 Bob 向工作人员提供消息。

我已经在另一个上下文中解决了这个问题,方法是为每个用户(甚至用户提交的每个批次)创建多个队列(子队列),并在消费者请求下一条消息时在所有子队列之间交替。

至少在我的世界里,这似乎是一个常见问题,我想知道是否有人知道解决它的既定方法。

我没有看到 ActiveMQ 的任何解决方案。我看过 Kafka,它能够在一个主题中循环分区,这可能会奏效。现在,我正在使用 Redis 实现一些东西。

【问题讨论】:

  • 你没有了解queueing theory
  • @GuyCoder 如果您能指出适用于此处的排队理论的某些方面,那将很有帮助。然而,即使我的问题涉及队列,我认为排队理论不会有帮助,因为我知道我想要的行为。我只是想找到一种实现行为的最佳方式。
  • 有 5003 条消息,但只有两个参与者。以循环方式为参与者提供服务,一次一条消息(或一批消息,以加快速度;您决定批量大小)。这就是你实际写的内容,这是一个很好的方法。
  • 感谢@dialecticus 的分析。我已经使用 Redis 提交了一个有效的实现。需要更多测试以了解它如何与许多并发用户、错误条件等一起扩展,但对于第一次尝试它很有好处。我将演员队列汇集到一个单独的就绪队列中,我将其保持在不超过活动演员队列数量的大小。漏斗逻辑使用最近服务最多的参与者队列的运行状态。

标签: algorithm message-queue


【解决方案1】:

我会推荐 Cadence Workflow 而不是队列,因为它支持长时间运行的操作和开箱即用的状态管理。

在您的情况下,我将为每个用户创建一个工作流实例。每个新任务都将通过信号 API 发送到用户工作流。然后工作流实例会将接收到的任务排队并一一执行。

下面是实现的概要:

public interface SerializedExecutionWorkflow {

    @WorkflowMethod
    void execute();

    @SignalMethod
    void addTask(Task t);
}

public interface TaskProcessorActivity {
    @ActivityMethod
    void process(Task poll);
}

public class SerializedExecutionWorkflowImpl implements SerializedExecutionWorkflow {

    private final Queue<Task> taskQueue = new ArrayDeque<>();
    private final TaskProcesorActivity processor = Workflow.newActivityStub(TaskProcesorActivity.class);

    @Override
    public void execute() {
        while(!taskQueue.isEmpty()) {
            processor.process(taskQueue.poll());
        }
    }

    @Override
    public void addTask(Task t) {
        taskQueue.add(t);
    }
}

然后是通过信号方法将该任务排入工作流的代码:

private void addTask(WorkflowClient cadenceClient, Task task) {
    // Set workflowId to userId
    WorkflowOptions options = new WorkflowOptions.Builder().setWorkflowId(task.getUserId()).build();
    // Use workflow interface stub to start/signal workflow instance
    SerializedExecutionWorkflow workflow = cadenceClient.newWorkflowStub(SerializedExecutionWorkflow.class, options);
    BatchRequest request = cadenceClient.newSignalWithStartRequest();
    request.add(workflow::execute);
    request.add(workflow::addTask, task);
    cadenceClient.signalWithStart(request);
}

与使用队列进行任务处理相比,Cadence 提供了许多其他优势。

  • 内置指数重试,无限期间隔
  • 故障处理。例如,如果在配置的时间间隔内两次更新都无法成功,它允许执行通知另一个服务的任务。
  • 支持长时间运行的心跳操作
  • 能够实现复杂的任务依赖。例如,在发生不可恢复的故障时实现调用链或补偿逻辑 (SAGA)
  • 提供对当前更新状态的完整可见性。例如,当使用队列时,您知道队列中是否有一些消息,并且您需要额外的数据库来跟踪整体进度。使用 Cadence 记录每个事件。
  • 能够取消正在进行的更新。

请参阅 the presentation,了解 Cadence 编程模型。

【讨论】:

  • 谢谢! Cadence Workflow 对于我的当前需求来说可能有点过于繁重,我需要得到整个组织的支持,但作为一种选择非常好了解。
  • 如果您想了解更多信息,请随时与我们联系。
猜你喜欢
  • 2017-09-06
  • 2012-11-27
  • 2016-07-15
  • 1970-01-01
  • 2016-08-21
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-02-07
相关资源
最近更新 更多