本文档介绍 RocketMQ 生产者的使用建议,推荐在使用消息队列 RocketMQ版进行消息生产与消费之前,阅读以下使用建议,提高接入效率和业务稳定性。
- 建议组合使用 Topic 和 tags,以减少 Topic 的使用。
- 仅当生产者在发送消息时设置了 Tag,消费者在订阅消息时才可以利用 Tag 进行消息过滤,例如 message.setTags("TagA")。消费者在 Broker 侧根据 Tag 的 hashcode 进行初步过滤,在消费端根据字符串过滤。
每个消息在业务层面的唯一标识码要设置到 keys 字段,便于定位消息丢失等问题。消息队列 RocketMQ版服务端会为每个消息创建索引,您可以在控制台中通过 topic + key 来查询这条消息的内容,以及消息被谁消费。使用消息 Key 时,请注意:
- 建议消息生产者为每条消息设置具有业务区分度的 Key。
- 索引类型为哈希索引,所以务必保证 key 尽可能唯一,以确保以此避免潜在的哈希冲突。
- 相同的 Key 的消息数量尽量少,最大不超过 64 条,否则消息查询结果不完整。
设置消息 Key 的方式请参考:
// 订单Id String orderId = "20034568923546";
message.setKeys(orderId);
RocketMQ 发送消息返回的 SendResult 里面会有两个消息 ID,一个是 msgId,一个是 offsetMsgId。
- msgId:客户端生成的唯一消息 ID,即便消息重发,消息 ID 也不会发生变化,一般可以作为唯一键用来消息去重。 msgId 生成规则主要包括客户端 IP、进程 ID、加载 MessageClientIDSetter 的类加载器的 hashcode、当前时间与系统启动时间的差值、自增序号等
- offsetMsgId:Broker 生成的消息 ID,主要是记录的 Broker 的地址和消息的物理偏移量。不能保证唯一,消息重发就会导致相同的消息有不一样的 msgId。
建议在消息发送成功或者失败时打印消息日志,日志中应包含 SendResult 和 Key 字段。可根据实际情况来选择是否打印消息体,如果消息内容比较重要,在消息发送失败时推荐打印消息体。
说明
对于发送结果为 SEND_OK 的消息,可以不打印消息日志,以免造成日志过多,浪费存储资源。
目前消息队列 RocketMQ版提供了三种消息发送模式,说明如下:
说明
- 异步发送和单向发送由于不需要等待返回结果就可以继续发送,消息的吞吐量会比较高,但是容易造成broker的发送线程池处理不过来,造成队列满了任务被拒绝。
send 消息方法只要不抛出异常,就代表发送成功。发送成功会有多个状态,在 sendResult 里定义。每个状态的说明如下:
Producer 的 send 方法本身支持内部重试,重试逻辑如下:
- 同步发送:重试 2 次,总共发送 3 次。通过 retryTimesWhenSendFailed 参数可以配置重试逻辑。重试时会选择另一个 Broker master 进行重试。
- 异步发送:重试 2 次,总共发送 3 次。通过 retryTimesWhenSendAsyncFailed 可配置重试次数。
以上策略在一定程度上可保证消息成功发送。如果业务对消息可靠性要求比较高,建议应用增加相应的重试逻辑。例如调用 send 同步方法发送失败时,则尝试将消息存储到数据库,然后由后台线程定时重试,确保消息一定到达 Broker。
说明
如果服务端 Broker 返回的 response code 为 SYSTEM_BUSY,客户端会直接抛出异常。例如以下情况:
- 消息存入 Broker pagecache 耗时超过 1000ms,导致消息快速失败被清理。
- 消息在发送队列里面等待时间超过 waitTimeMillsInSendQueue(默认 1000ms)。
- Broker pagecache busy,消息到达 Broker 直接被拒绝请求。
- Broker 发送线程池队列任务已满,新任务被拒绝提交到线程池。
Producer 默认的消息发送超时时间是 3000ms,可通过 sendMsgTimeout 配置。在生产环境,不建议将超时时间配置的很短,因为 RocketMQ 只能保证 P99 的延迟在几毫秒以内,部分毛刺的时间可能会比较大。如果时间配置较短,容易导致消息发送失败。