如何使用Apache Kafka实现支持服务扩缩容的请求/响应模式(Node.js环境)
基于Apache Kafka实现支持扩缩容的请求/响应模式方案
核心替代方案(对应RabbitMQ临时队列能力)
Kafka没有原生的临时队列断连自动删除机制,我们可以通过共享Topic+实例标识过滤的方式模拟同等能力,首先规避一个常见的错误方案:
注意:不推荐为每个实例创建临时Topic的方案,Kafka的Topic创建属于重量级操作,频繁扩缩容场景下会大幅提升集群管控压力,远不如共享Topic的方案高效。
推荐的实现思路如下:
- 为每个业务服务创建独立的请求Topic,所有请求方统一往对应Topic发送请求,消息头携带关联请求ID(correlation_id)和请求方实例唯一标识两个字段
- 所有服务实例共享一个全局响应Topic,不需要为每个实例创建单独的消费资源,实例收到响应后只处理消息头中目标实例ID和自己一致的消息
- 实例异常下线后,Kafka会自动超时清理对应消费组的元数据,实现和RabbitMQ临时队列一致的自动清理效果
Node.js 技术栈具体实现(基于kafkajs客户端)
Node.js生态推荐使用kafkajs作为Kafka客户端,是目前活跃度最高、适配性最好的实现,具体代码逻辑如下:
- 实例唯一标识生成
K8s环境下直接使用Pod名称作为实例唯一标识即可,Pod名称天然全局唯一,可以通过os.hostname()直接读取,也可以启动时注入环境变量。
const os = require('os') const INSTANCE_ID = process.env.INSTANCE_ID || os.hostname()
- 请求发送方逻辑
const { Kafka } = require('kafkajs') const kafka = new Kafka({ brokers: ['你的Kafka broker地址列表'] }) const producer = kafka.producer() const pendingRequests = new Map() // 存储待响应的请求回调 // 封装请求发送方法 const sendRequest = async (serviceName, payload, timeout = 10000) => { await producer.connect() const correlationId = Math.random().toString(36).slice(2, 18) // 发送请求到对应服务的请求Topic await producer.send({ topic: `request-${serviceName}`, messages: [ { value: JSON.stringify(payload), headers: { correlation_id: correlationId, reply_instance_id: INSTANCE_ID } } ] }) // 封装超时等待响应的Promise return new Promise((resolve, reject) => { const timer = setTimeout(() => { pendingRequests.delete(correlationId) reject(new Error('请求超时')) }, timeout) pendingRequests.set(correlationId, (res) => { clearTimeout(timer) resolve(res) }) }) }
- 服务提供方逻辑
const serviceName = '你的服务名' const consumer = kafka.consumer({ groupId: `service-${serviceName}-group` }) const runService = async () => { await consumer.connect() await consumer.subscribe({ topic: `request-${serviceName}` }) await consumer.run({ eachMessage: async ({ message }) => { const correlationId = message.headers.correlation_id.toString() const replyInstanceId = message.headers.reply_instance_id.toString() // 处理业务逻辑 const requestPayload = JSON.parse(message.value.toString()) const responsePayload = await handleBusinessLogic(requestPayload) // 发送响应到全局响应Topic await producer.send({ topic: 'response-global', messages: [ { value: JSON.stringify(responsePayload), headers: { correlation_id: correlationId, target_instance_id: replyInstanceId } } ] }) } }) }
- 响应接收逻辑(所有实例都需要启动)
const runResponseConsumer = async () => { // 每个实例用自己的唯一ID作为消费组ID,设置短会话超时时间 const responseConsumer = kafka.consumer({ groupId: `response-${INSTANCE_ID}`, sessionTimeout: 15000 // 15秒会话超时,实例下线后自动清理消费组 }) await responseConsumer.connect() // 只消费最新的响应,历史过期响应不需要处理 await responseConsumer.subscribe({ topic: 'response-global', fromBeginning: false }) await responseConsumer.run({ eachMessage: async ({ message }) => { const targetInstanceId = message.headers.target_instance_id.toString() // 过滤掉不是发给当前实例的响应 if (targetInstanceId !== INSTANCE_ID) return const correlationId = message.headers.correlation_id.toString() const callback = pendingRequests.get(correlationId) if (callback) { callback(JSON.parse(message.value.toString())) pendingRequests.delete(correlationId) } } }) }
扩缩容适配说明
- 扩容:新Pod启动后自动生成唯一实例ID,启动请求处理消费者和响应消费者后即可正常工作,Kafka会自动将请求Topic的分区负载均衡到新加入的消费者实例,不需要额外调整配置
- 缩容:Pod正常下线时,
kafkajs会自动触发消费者优雅关闭,提交已消费的偏移量;如果是异常crash下线,15秒会话超时后,Kafka会自动清理该实例对应的响应消费组元数据,不会残留无效资源,和RabbitMQ临时队列的自动清理效果一致 - 优化建议:给全局响应Topic设置较短的消息保留时间(比如30秒,超过请求超时时间的响应没有业务价值),Kafka会自动清理过期消息,降低存储和消费压力
内容的提问来源于stack exchange,提问作者Arthur
相关产品推荐
相关产品推荐

