Node.js使用GCP PubSub排序功能时同orderingKey消息乱序问题
问题描述
在Node.js环境下使用GCP Pub/Sub的消息排序特性,发送多条携带相同ordering key的消息。按照特性预期,前一条消息完成ack确认后才会投递下一条同key消息,但实际观测到:第一条消息被ack后,剩余同key消息被批量推送至消费端,且投递顺序不符合排序要求,该现象可通过结果截图中的时间戳佐证。
运行结果截图
消息发布端代码
const main = async () => { const pubSubClient = new PubSub() const keys = ['a', 'a', 'a'] for (let i = 0; i < keys.length; ++i) { console.log(chalk.green(`Sending message ${i + 1}, ordering key - ${keys[i]}`)) await pubSubClient .topic('test-topic', { messageOrdering: true }) .publishMessage({ orderingKey: keys[i], data: Buffer.from(`Msg ${i + 1}`) }) } }
消息拉取消费端代码
const main = async () => { const pubSubClient = new PubSub() const handleMsg = async (message: Message) => { const dataStr = Buffer.from((message.data as unknown) as string, 'base64').toString('utf-8') console.log(chalk.green(`Got message ${dataStr}, orderingKey - ${message.orderingKey}, ${new Date().toISOString()}`)) // some processing happens here await sleep(10000) console.log(chalk.yellow(`Message ack ${dataStr}, orderingKey - ${message.orderingKey}, ${new Date().toISOString()}`)) message.ack() } const subscription = pubSubClient.subscription('test-subscription') subscription.on('message', handleMsg) console.log('started listening...') }
问题根因
问题由两方面原因共同导致:
- 对排序投递机制的认知偏差:GCP Pub/Sub的有序消息保证,是同ordering key的消息最终交付给业务逻辑的顺序,和发布端收到发布成功确认的顺序完全一致,并非服务端侧等前一条消息ack后才向下推送下一条。为了平衡吞吐和顺序保证,服务端会通过长连接预推送多条同key消息到消费端SDK本地缓存,由SDK负责顺序调度:开启排序支持的SDK会严格缓存后续消息,等前一条同key消息被业务处理完成、返回ack后,才会把下一条消息抛给业务回调。
- 核心配置缺失:消费端代码创建Subscription实例时,没有传入
enableMessageOrdering: true参数。Node.js版Pub/Sub SDK默认关闭有序消息调度逻辑,这种状态下SDK不会对服务端推送来的消息做任何缓存和顺序校验,收到消息就立刻并发触发message事件回调,自然会出现第一条消息ack后剩余消息批量触发、顺序错乱的现象。 - 附加非核心问题:发布端循环内重复调用
.topic()创建新的主题实例,会重复初始化底层发布连接、增加不必要的开销,不会直接导致顺序问题,但不符合SDK最佳实践。
修复方案
- 消费端创建订阅实例时显式开启有序消息支持:
const subscription = pubSubClient.subscription('test-subscription', { enableMessageOrdering: true })
开启后SDK会自动处理同key消息的串行调度,保证消息按发布顺序进入handleMsg处理逻辑,前一条消息ack前不会交付后续同key消息,无需额外自行实现并发控制。
2. 发布端提前初始化单例主题实例,避免循环内重复创建:
const main = async () => { const pubSubClient = new PubSub() // 提前初始化单例topic复用 const orderedTopic = pubSubClient.topic('test-topic', { messageOrdering: true }) const keys = ['a', 'a', 'a'] for (let i = 0; i < keys.length; ++i) { console.log(chalk.green(`Sending message ${i + 1}, ordering key - ${keys[i]}`)) await orderedTopic.publishMessage({ orderingKey: keys[i], data: Buffer.from(`Msg ${i + 1}`) }) } }
内容的提问来源于stack exchange,提问作者Rusty Gold
相关产品推荐
相关产品推荐

