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

如何用Jest模拟EventHubConsumerClient的回调参数

如何模拟Event Hub消费者回调中的events和context参数?

我正在使用微软Azure官方文档中的示例代码创建Event Hub消费者,想知道应当如何模拟传入回调函数的events和context参数,相关代码如下:

const { EventHubConsumerClient, earliestEventPosition } = require("@azure/event-hubs");
const { ContainerClient } = require("@azure/storage-blob");    
const { BlobCheckpointStore } = require("@azure/eventhubs-checkpointstore-blob");

const connectionString = "EVENT HUBS NAMESPACE CONNECTION STRING";    
const eventHubName = "EVENT HUB NAME";
const consumerGroup = "$Default"; // 默认消费者组名称
const storageConnectionString = "AZURE STORAGE CONNECTION STRING";
const containerName = "BLOB CONTAINER NAME";

async function main() {
  // 创建Blob容器客户端和基于该客户端的检查点存储
  const containerClient = new ContainerClient(storageConnectionString, containerName);
  const checkpointStore = new BlobCheckpointStore(containerClient);

  // 指定检查点存储,创建Event Hub消费者客户端
  const consumerClient = new EventHubConsumerClient(consumerGroup, connectionString, eventHubName, checkpointStore);

  // 订阅事件,并指定事件处理和错误处理的回调函数
  const subscription = consumerClient.subscribe({
      processEvents: async (events, context) => {
        if (events.length === 0) {
          console.log(`在等待时间内未收到事件,等待下一个周期`);
          return;
        }

        for (const event of events) {
          console.log(`收到事件: '${event.body}',来自分区: '${context.partitionId}',消费者组: '${context.consumerGroup}'`);
        }
        // 更新检查点
        await context.updateCheckpoint(events[events.length - 1]);
      },

      processError: async (err, context) => {
        console.log(`错误 : ${err}`);
      }
    },
    { startPosition: earliestEventPosition }
  );

  // 30秒后停止处理
  await new Promise((resolve) => {
    setTimeout(async () => {
      await subscription.close();
      await consumerClient.close();
      resolve();
    }, 30000);
  });
}

main().catch((err) => {
  console.log("发生错误: ", err);
});

模拟参数的方法

1. 模拟events参数

events是事件对象数组,每个对象只需包含业务逻辑中用到的核心字段,示例如下:

const mockEvents = [
  {
    body: "测试消息1",
    partitionKey: "partition-0",
    offset: "12345",
    sequenceNumber: 1,
    enqueuedTimeUtc: new Date()
  },
  {
    body: "测试消息2",
    partitionKey: "partition-0",
    offset: "67890",
    sequenceNumber: 2,
    enqueuedTimeUtc: new Date()
  }
];

2. 模拟context参数

context需要包含代码中用到的属性和方法,比如partitionId、consumerGroup,以及updateCheckpoint异步方法。如果只是验证业务逻辑,可给方法写一个简单模拟实现:

const mockContext = {
  partitionId: "partition-0",
  consumerGroup: "$Default",
  // 模拟检查点更新逻辑,可根据需求添加具体行为
  updateCheckpoint: async (event) => {
    console.log(`模拟更新检查点,offset: ${event.offset}, sequenceNumber: ${event.sequenceNumber}`);
  }
};

3. 测试调用回调函数

建议把processEvents单独提取成独立函数,方便直接传入模拟参数测试:

// 单独提取事件处理函数
const processEvents = async (events, context) => {
  if (events.length === 0) {
    console.log(`在等待时间内未收到事件,等待下一个周期`);
    return;
  }

  for (const event of events) {
    console.log(`收到事件: '${event.body}',来自分区: '${context.partitionId}',消费者组: '${context.consumerGroup}'`);
  }
  await context.updateCheckpoint(events[events.length - 1]);
};

// 传入模拟参数测试
processEvents(mockEvents, mockContext).catch(err => console.error(err));

如果需要更复杂的场景模拟,也可以用Jest等测试框架Mock相关类,但上述手动模拟方式足够覆盖大部分业务逻辑验证需求。

内容的提问来源于stack exchange,提问作者RMMD12

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 20:20:03