KafkaJS仅事务模式报producer is disconnected错误咨询
KafkaJS事务生产者抛出"The producer is disconnected"错误排查
问题根因
报错核心原因是创建事务型生产者实例后,未显式建立连接就直接初始化事务。
你当前代码只完成了生产者的配置初始化,没有和Kafka broker建立实际网络连接:
- 普通生产者调用
send()时,内部会自动检测连接状态,未连接时隐式发起连接,因此测试普通发消息无报错 - 事务型生产者的
transaction()初始化方法没有内置自动连接逻辑,调用时直接校验连接状态,未连接就直接抛出The producer is disconnected错误
另外代码存在不影响运行的笔误:clientId配置值app-transactional-clienet拼写错误,正确应为app-transactional-client。你配置的idempotent: true、maxInFlightRequests: 1是事务生产者的强制要求配置,这部分没有问题。
修复方案
在创建事务生产者实例后、调用producer.transaction()前,显式执行await producer.connect()完成连接建立即可。
修复后的核心代码如下:
const client = new Kafka({ clientId: 'app-transactional-client', // 修正拼写笔误 brokers: ['host.docker.internal:9092'], }) const producer = client.producer({ transactionalId: 'transaction.connectedservice.producer', maxInFlightRequests: 1, idempotent: true }) // 新增:显式建立生产者连接 await producer.connect() const transaction = await producer.transaction(); try { const connectedAccountMsg: ConnectedAcccountMsg = { userId: req.session.userId, accountId: profile.id, accountUsername: profile.username, accountType: 'twitter', authorized: true, }; await transaction.send({ topic: 'user.connectedaccount.status', messages: [{ value: JSON.stringify(connectedAccountMsg) }], }) await transaction.commit() // 可选:事务提交后如果后续不再使用生产者,可主动断开连接 // await producer.disconnect() return done(null, profile); } catch (error) { console.log(error) await transaction.abort() // 异常场景下也可按需断开生产者连接 // await producer.disconnect() }
额外注意事项
- 事务型生产者不要每次发消息都新建实例、重复执行connect/disconnect,建议全局复用单实例,减少
transactionalId和broker事务协调器握手的开销 - 确保Kafka broker版本在0.11以上才支持事务消息能力,你当前普通生产者能正常发消息,说明broker连接、监听器配置均无问题,不需要额外调整broker配置
内容的提问来源于stack exchange,提问作者nick
相关产品推荐
相关产品推荐

