You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.20 07:15:35