NestJS E2E测试因Kafka连接异常失败求助
问题:Docker容器中NestJS E2E测试Kafka连接失败
主Nest应用与Kafka交互正常,能够发送消息并订阅响应,但在Docker容器中运行E2E测试时,因Kafka连接问题导致测试失败。去掉网关中调用this.client.send的代码后,测试即可通过。
错误日志
[Nest] 29 - 09/09/2022, 10:04:17 AM WARN [ClientKafka] WARN [undefined] KafkaJS v2.0.0 switched default partitioner. To retain the same partitioning behavior as in previous versions, create the producer with the option "createPartitioner: Partitioners.LegacyPartitioner". See the migration guide at https://kafka.js.org/docs/migration-guide-v2.0.0#producer-new-default-partitioner for details. Silence this warning by setting the environment variable "KAFKAJS_NO_PARTITIONER_WARNING=1" {"timestamp":"2022-09-09T10:04:17.902Z","logger":"kafkajs"} ms-gateway_1 | [Nest] 29 - 09/09/2022, 10:04:17 AM LOG [NestApplication] Nest application successfully started +68ms ms-gateway_1 | [Nest] 29 - 09/09/2022, 10:07:28 AM ERROR [ClientKafka] ERROR [Connection] Response Heartbeat(key: 12, version: 3) {"timestamp":"2022-09-09T10:07:28.205Z","logger":"kafkajs","broker":"kafka:9092","clientId":"gateway-client-client","error":"The group is rebalancing, so a rejoin is needed","correlationId":48,"size":10} ms-gateway_1 | [Nest] 29 - 09/09/2022, 10:07:28 AM ERROR [ClientKafka] ERROR [Connection] Response Heartbeat(key: 12, version: 3) {"timestamp":"2022-09-09T10:07:28.214Z","logger":"kafkajs","broker":"kafka:9092","clientId":"gateway-client-client","error":"The group is rebalancing, so a rejoin is needed","correlationId":49,"size":10} ms-gateway_1 | [Nest] 29 - 09/09/2022, 10:07:28 AM ERROR [ClientKafka] ERROR [Connection] Response Heartbeat(key: 12, version: 3) {"timestamp":"2022-09-09T10:07:28.216Z","logger":"kafkajs","broker":"kafka:9092","clientId":"gateway-client-client","error":"The group is rebalancing, so a rejoin is needed","correlationId":50,"size":10} ms-gateway_1 | [Nest] 29 - 09/09/2022, 10:07:28 AM WARN [ClientKafka] WARN [Runner] The group is rebalancing, re-joining {"timestamp":"2022-09-09T10:07:28.218Z","logger":"kafkajs","groupId":"gateway-consumer-client","memberId":"gateway-client-client-eaa78df1-177b-484d-9c22-739e7f4a0e49","error":"The group is rebalancing, so a rejoin is needed"} ms-gateway_1 | [Nest] 29 - 09/09/2022, 10:07:28 AM LOG [ClientKafka] INFO [ConsumerGroup] Consumer has joined the group {"timestamp":"2022-09-09T10:07:28.243Z","logger":"kafkajs","groupId":"gateway-consumer-client","memberId":"gateway-client-client-eaa78df1-177b-484d-9c22-739e7f4a0e49","leaderId":"gateway-client-client-eaa78df1-177b-484d-9c22-739e7f4a0e49","isLeader":true,"memberAssignment":{"AccountBalanceQuery.reply":[0],"AccountCreateCommand.reply":[0],"AccountBalanceTransferCommand.reply":[0]},"groupProtocol":"NestReplyPartitionAssigner","duration":23}
网关正常运行的代码
@Query(returns => AccountBalance) async accountBalance2() { // tests passes without this line let result = await firstValueFrom(this.client.send(BROKER_MESSAGES.ACCOUNT_BALANCE_QUERY, { account_id: '05fbebcd-710c-48fc-a994-7fd0152bdc89' })); console.log('resss:', result) return (result >= 0) ? { balance: result } : { error:"Not Found!" }; }
E2E测试代码
jest.setTimeout(30000) describe('App Module (e2e)', () => { let app: INestApplication; beforeAll(async () => { const moduleFixture: TestingModule = await Test.createTestingModule({ imports: [ AppModule, ], }) .compile(); app = moduleFixture.createNestApplication(); await app.init(); await app.startAllMicroservices(); }); afterAll(async () => { await app.close(); }); describe('test graphql', () => { it('/graphql (GET)', () => { const CREATE_ACCOUNT_MUTATION = 'query { accountBalance2 { balance } }' return request(app.getHttpServer()) .post('/graphql') .send({ query:CREATE_ACCOUNT_MUTATION, }) .expect(200) }); }) });
排查与解决方案
1. Mock Kafka Client(推荐)
E2E测试无需依赖真实Kafka服务,直接MockClientKafka的send方法,避免真实网络调用带来的不稳定:
beforeAll(async () => { const moduleFixture: TestingModule = await Test.createTestingModule({ imports: [AppModule], }) .overrideProvider(ClientKafka) .useValue({ send: jest.fn().mockResolvedValue(100), // 模拟返回预期的余额数值 }) .compile(); app = moduleFixture.createNestApplication(); await app.init(); await app.startAllMicroservices(); });
2. 等待Kafka消费者就绪
测试启动后,应用虽已启动,但Kafka消费者可能还在组重平衡阶段,可添加延迟确保消费者完成加入组:
beforeAll(async () => { // ... 现有模块编译、应用初始化逻辑 await app.startAllMicroservices(); // 添加5秒延迟,等待Kafka消费者完成重平衡 await new Promise(resolve => setTimeout(resolve, 5000)); });
3. 检查Docker网络配置
确保测试容器与Kafka容器处于同一Docker网络,保证kafka:9092地址可被正确解析。在Docker Compose中配置自定义网络,让所有服务加入该网络。
4. 调整Kafka消费者配置
在ClientKafka的配置中增加重平衡超时相关参数,延长等待时间,避免心跳失败:
// 在Kafka模块配置中添加 { client: { brokers: ['kafka:9092'], }, consumer: { groupId: 'gateway-consumer-client', sessionTimeout: 30000, rebalanceTimeout: 60000, }, }
内容的提问来源于stack exchange,提问作者Vahid Najafi
相关产品推荐
相关产品推荐

