消费类型
Apache RocketMQ 支持以下几种消费者类型:PushConsumer、SimpleConsumer 和 PullConsumer。本主题介绍了这三种消费者类型的使用方法、工作原理、重试机制及适用场景。
背景信息
Apache RocketMQ 提供了 PushConsumer、SimpleConsumer 和 PullConsumer 三种消费者类型。这三种消费者类型具有不同的集成和控制方式,您可以利用它们来满足不同业务场景下的消息需求。以下因素可帮助您为业务场景选择合适的消费者类型:
并发消费:消费者如何利用多线程技术实现并发消息消费,以提高消息处理效率?
同步或异步消息处理:在不同的集成场景中,消费者可能需要将接收到的消息异步分发到业务逻辑系统进行处理。如何实现异步消息处理?
可靠的消息处理:消费者在处理消息时如何返回响应结果?当消息处理出错时,如何实现消息重试以确保可靠的消息处理?
有关上述问题的解答,请参阅 PushConsumer 和 SimpleConsumer。
功能概述

上图显示了 Apache RocketMQ 中消费者消费消息的过程涉及以下阶段:接收消息、处理消息以及提交消费状态。
这三种类型的消费者通过提供不同的实现方法和 API 操作,适用于各种消息消费场景。下表描述了这三种消费者类型之间的区别。
仅建议在流处理框架中集成 PullConsumer。PushConsumer 和 SimpleConsumer 可以满足大多数场景需求。
您可以根据业务场景在 PushConsumer 和 SimpleConsumer 之间切换。切换到不同的消费者类型时,Apache RocketMQ 中现有资源的使用和现有的业务处理任务不会受到影响。
严禁在同一个消费者组(consumerGroup)中混用 PullConsumer 和其他类型的消费者。
| 项目 | PushConsumer | SimpleConsumer | PullConsumer |
|---|---|---|---|
| API 操作调用 | 通过使用消息监听器调用回调操作来返回消费结果。消费者只能在消息监听器的作用域内处理消费逻辑。 | 业务应用程序实现消息处理并调用相应的操作来返回消费结果。 | 业务应用程序实现消息拉取和处理,并调用相应的操作来返回消费结果。 |
| 消费并发管理 | 使用 Apache RocketMQ SDK 来管理消息消费的并发线程数。 | 用于消息消费的并发线程数基于各个业务应用程序的消费逻辑。 | 用于消息消费的并发线程数基于各个业务应用程序的消费逻辑。 |
| 负载均衡机制 | 5.0 版本中基于消息的负载均衡,早期版本中基于队列的负载均衡。 | 基于消息的负载均衡。 | 基于队列的负载均衡。 |
| API 灵活性 | API 操作经过封装,灵活性较差。 | 原子操作提供了极大的灵活性。 | 原子操作提供了极大的灵活性。 |
| 适用场景 | 此消费者类型适用于不需要自定义流程的开发场景。 | 此消费者类型适用于需要自定义流程的开发场景。 | 仅建议在流处理框架场景中集成 |
PushConsumer
PushConsumer 是一种提供高度封装的消费者类型。消息消费和消费结果提交仅通过消息监听器进行处理。消息获取、消费状态提交和消费重试均由 Apache RocketMQ 客户端 SDK 完成。
使用方法
PushConsumer 的使用方式是固定的。在初始化 PushConsumer 消费者时注册一个消息监听器,并在消息监听器中实现消息处理逻辑。消息获取、触发监听器调用和消息重试均由 Apache RocketMQ SDK 处理。
示例代码
// Message consumption example: Use a PushConsumer consumer to consume messages.
ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "YourTopic";
FilterExpression filterExpression = new FilterExpression("YourFilterTag", FilterExpressionType.TAG);
PushConsumer pushConsumer = provider.newPushConsumerBuilder()
// Configure consumer group.
.setConsumerGroup("YourConsumerGroup")
// Specify the access point.
.setClientConfiguration(ClientConfiguration.newBuilder().setEndpoints("YourEndpoint").build())
// Specify the pre-bound subscriptions.
.setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
// Set the message listener.
.setMessageListener(new MessageListener() {
@Override
public ConsumeResult consume(MessageView messageView) {
// Consume the messages and return the consumption result.
return ConsumeResult.SUCCESS;
}
})
.build();
PushConsumer 消费者的消息监听器会返回以下结果之一:
消费成功:例如,当使用 Java 版 Apache RocketMQ SDK 且消息被成功消费时,返回
ConsumeResult.SUCCESS。服务端会根据消费结果更新消费进度。消费失败:例如,当使用 Java 版 Apache RocketMQ SDK 且消息消费失败时,返回
ConsumeResult.FAILURE。Apache RocketMQ 是否重试消费该消息取决于消费重试逻辑。意外失败:例如,如果抛出意外异常,则消息消费失败。Apache RocketMQ 是否重试消费该消息取决于消费重试逻辑。
如果消息处理逻辑中的意外错误持续导致 PushConsumer 消费者无法消费消息,SDK 会认为消费超时并强制提交消费失败结果。随后,该消息将根据消费重试逻辑进行处理。有关消费超时的更多信息,请参阅 PushConsumer 重试策略。
当发生消费超时时,SDK 会提交消费失败结果。但是,当前的消费线程可能无法及时响应结果并继续处理该消息。::
工作原理
对于 PushConsumer,实时消息处理基于 SDK 典型的 Reactor 线程模型。SDK 内置了一个长轮询线程,负责拉取消息并将消息存储到队列中。然后,消息从队列传递到各个消息消费线程。消息监听器的行为基于消息消费逻辑。下图显示了 PushConsumer 消费者的消息消费过程。
可靠性重试
对于 PushConsumer,客户端 SDK 与消费逻辑单元之间的通信仅通过消息监听器实现。客户端 SDK 根据消息监听器返回的结果检查消息是否被消费,并根据消费重试逻辑执行重试,以确保消息可靠性。所有消息必须以同步方式进行消费。消费结果在监听器操作调用结束时返回。不允许异步分发。有关消息重试的更多信息,请参阅 PushConsumer 重试策略。
为确保消息传递的可靠性,Apache RocketMQ 禁止 PushConsumer 消费者在消息消费过程中执行以下行为:
在消息消费完成前返回消费结果。例如,对于后续消费失败的消息提前返回消费成功结果。在这种情况下,Apache RocketMQ 无法检查实际的消费结果,也不会重试消费该消息。
从消息监听器将消息分发给其他自定义线程并提前返回消费结果。如果消息消费失败但提前返回了消费成功结果,Apache RocketMQ 无法检查实际的消费结果,也不会重试消费该消息。
确保消息顺序
对于 Apache RocketMQ 中的 FIFO 消息,如果为消费者组配置了顺序消息消费,则 PushConsumer 消费者会按消费顺序消费消息。当 PushConsumer 消费者消费消息时,无需业务应用程序在业务逻辑中定义消费顺序即可确保消费顺序。
在 Apache RocketMQ 中,同步提交是处理顺序消息的前提。如果业务逻辑中定义了异步分发,Apache RocketMQ 将无法确保消息的顺序。 ::
适用场景
PushConsumer 将消息处理限制为同步处理,并限制了每条消息的处理超时时间。PushConsumer 适用于以下场景:
可预测的消息处理耗时:如果消息处理耗时没有限制,对于需要较长处理时间的消息,会持续触发消息重试以确保消息可靠性。这将导致大量的重复消息。
无异步处理且无自定义流程:PushConsumer 将消费逻辑的线程模型限制为 Reactor 线程模型。客户端 SDK 根据最大吞吐量处理消息。这种模型开发简单,但不允许异步或自定义流程。
SimpleConsumer
SimpleConsumer 是一种支持消息处理原子操作的消费者类型。此类消费者调用操作来获取消息、提交消费状态,并根据业务逻辑执行消息重试。
使用方法
SimpleConsumer 涉及多个 API 操作。根据需要调用相应的操作来获取消息并将其分发到业务线程进行处理。然后,调用提交(commit)操作来提交消息处理结果。示例代码
// Consumption example: When a SimpleConsumer consumer consumes normal messages, the consumer obtain messages and commit message consumption results.
ClientServiceProvider provider = ClientServiceProvider.loadService();
String topic = "YourTopic";
FilterExpression filterExpression = new FilterExpression("YourFilterTag", FilterExpressionType.TAG);
SimpleConsumer simpleConsumer = provider.newSimpleConsumerBuilder()
// Configure consumer group.
.setConsumerGroup("YourConsumerGroup")
// Specify the access point.
.setClientConfiguration(ClientConfiguration.newBuilder().setEndpoints("YourEndpoint").build())
// Specify the pre-bound subscriptions.
.setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
// Specify the max await time when receive messages from the server.
.setAwaitDuration(Duration.ofSeconds(1))
.build();
try {
// A SimpleConsumer consumer must obtain and process messages.
List<MessageView> messageViewList = simpleConsumer.receive(10, Duration.ofSeconds(30));
messageViewList.forEach(messageView -> {
System.out.println(messageView);
// After consumption is complete, the consumer must invoke ACK to submit the consumption result.
try {
simpleConsumer.ack(messageView);
} catch (ClientException e) {
logger.error("Failed to ack message, messageId={}", messageView.getMessageId(), e);
}
});
} catch (ClientException e) {
// If the pull fails due to system traffic throttling or other reasons, the consumer must re-initiate the request to obtain the message.
logger.error("Failed to receive message", e);
}
下表描述了为 SimpleConsumer 提供的 API 操作。
| 操作 | 描述 | 可修改参数 |
|---|---|---|
ReceiveMessage | 消费者可以调用此操作从服务端获取消息。注意:由于服务端使用分布式存储,即使请求的消息实际上存在于服务端,服务端也可能返回空结果。您可以再次调用 ReceiveMessage 操作或增加 ReceiveMessage 操作中的并发值。 | 批量拉取大小:一次获取的消息数量。SimpleConsumer 消费者可以获取多条消息进行批量消费。 消息不可见持续时间:消息的最大处理时长。此参数控制消费失败时的消息重试间隔。有关更多信息,请参阅 SimpleConsumer 重试策略。调用 ReceiveMessage 操作时需要此参数。 |
AckMessage | 消费者消费消息后,调用此操作向服务端返回消费成功结果。 | 无 |
ChangeInvisibleDuration | 在消费重试场景中,消费者可以调用此操作来更改消息处理时长,以控制消息重试间隔。 | 消息不可见持续时间:消息的最大处理时间。您可以调用此操作来更改 ReceiveMessage 操作中指定的消息不可见持续时间。在大多数情况下,此操作用于需要增加消息处理时长的场景。 |
可靠性重试
当 SimpleConsumer 消费者消费消息时,客户端 SDK 与 Apache RocketMQ 服务端之间的通信通过 ReceiveMessage 和 AckMessage 操作实现。当客户端 SDK 成功处理消息时,会调用 AckMessage 操作。当消息处理失败时,不返回确认消息,从而在指定的不可见持续时间过后触发消息重试机制。有关更多信息,请参阅 SimpleConsumer 重试策略。
确保消息顺序
在 Apache RocketMQ 中,SimpleConsumer 消费者按照 FIFO 消息 的存储顺序获取它们。如果一组顺序消息中的某条消息未处理完成,则无法获取该组顺序消息中的下一条消息。
适用场景
SimpleConsumer 提供了用于获取消息和提交消费结果的原子 API 操作。与 PushConsumer 相比,SimpleConsumer 提供了更好的灵活性。SimpleConsumer 适用于以下场景:
不可控的消息处理耗时:如果消息处理耗时不可预估,我们建议使用 SimpleConsumer 以防止消息处理时间过长。您可以在消息消费期间指定预估的消息处理时长。如果现有的处理时长不适合您的业务场景,您可以调用相应的 API 操作来更改消息处理时长。
异步处理和批量消费:SimpleConsumer 不涉及 SDK 中复杂的线程封装。业务应用程序可以使用自定义设置。这样,SimpleConsumer 消费者可以实现异步分发、批量消费和其他自定义场景。
自定义消息消费速率:使用 SimpleConsumer 时,业务应用程序调用 ReceiveMessage 操作来获取消息。您可以调整获取消息的频率以控制消息消费速率。
PullConsumer
待续。
使用说明
为 PushConsumer 指定合理的消费耗时限制
我们建议限制 PushConsumer 消费者的消息消费耗时,以防止某条消息处理时间过长。长时间处理消息可能会因消息处理超时而导致重复消息,并使下一条消息持续等待消费。如果消息经常处理过长,建议使用 SimpleConsumer 并根据业务需求指定合适的消息不可见持续时间。