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

如何连接Kafka端口并读取主题与消息完成集成测试

在Aspire集成测试中连接Kafka并读取主题与消息

1. 安装Kafka客户端依赖

给测试项目安装Confluent.Kafka NuGet包,这是操作Kafka的标准客户端库。

2. 获取Kafka连接配置

在集成测试中,通过Aspire应用主机实例获取Kafka资源的连接信息,替代硬编码端口,确保测试环境兼容性:

// 获取指定名称的Kafka资源
var kafkaResource = _appHost.Resources.OfType<KafkaResource>().First(r => r.Name == "kafka");
// 提取包含Bootstrap Servers的连接字符串
var bootstrapServers = kafkaResource.GetConnectionString();

3. 读取所有Kafka主题

使用Kafka AdminClient获取所有主题列表:

var adminConfig = new AdminClientConfig
{
    BootstrapServers = bootstrapServers
};

using var adminClient = new AdminClientBuilder(adminConfig).Build();
// 获取Kafka集群元数据
var metadata = adminClient.GetMetadata(TimeSpan.FromSeconds(10));
// 提取所有主题名称
var allTopics = metadata.Topics.Select(t => t.Topic).ToList();

// 输出主题(可用于日志或测试断言)
foreach (var topic in allTopics)
{
    Console.WriteLine($"已发现主题:{topic}");
}

4. 读取特定主题的消息

创建Kafka Consumer客户端,订阅目标主题并拉取消息:

var consumerConfig = new ConsumerConfig
{
    BootstrapServers = bootstrapServers,
    GroupId = "integration-test-group", // 测试专用消费组ID,避免与其他环境冲突
    AutoOffsetReset = AutoOffsetReset.Earliest // 从主题起始位置读取消息
};

using var consumer = new ConsumerBuilder<Ignore, string>(consumerConfig).Build();
// 替换为你需要读取的目标主题名称
var targetTopic = "your-target-topic";
consumer.Subscribe(targetTopic);

try
{
    // 拉取消息示例(可根据需求调整拉取次数或超时时间)
    for (int i = 0; i < 10; i++)
    {
        var consumeResult = consumer.Consume(TimeSpan.FromSeconds(5));
        if (consumeResult != null)
        {
            Console.WriteLine($"读取到消息:{consumeResult.Message.Value}");
            // 可添加测试断言,比如验证消息内容符合预期
        }
    }
}
finally
{
    consumer.Close();
}

完整IntegrationTest.cs补充代码

将上述逻辑整合到现有测试代码中,完整片段如下:

_appHost = await DistributedApplicationTestingBuilder.CreateAsync<Projects.PI_AspireSolution_AppHost>();
_appHost.Services.ConfigureHttpClientDefaults(clientBuilder =>
{
    clientBuilder.AddStandardResilienceHandler();
});
_app = await _appHost.BuildAsync();
await _app.StartAsync();
await _resourceNotificationService.WaitForResourceAsync("kafka", KnownResourceStates.Running).WaitAsync(TimeSpan.FromSeconds(120));

// --- 新增Kafka操作逻辑 ---
// 获取Kafka连接配置
var kafkaResource = _appHost.Resources.OfType<KafkaResource>().First(r => r.Name == "kafka");
var bootstrapServers = kafkaResource.GetConnectionString();

// 读取所有主题
var adminConfig = new AdminClientConfig { BootstrapServers = bootstrapServers };
using var adminClient = new AdminClientBuilder(adminConfig).Build();
var metadata = adminClient.GetMetadata(TimeSpan.FromSeconds(10));
var allTopics = metadata.Topics.Select(t => t.Topic).ToList();

// 读取特定主题消息
var consumerConfig = new ConsumerConfig
{
    BootstrapServers = bootstrapServers,
    GroupId = "integration-test-group",
    AutoOffsetReset = AutoOffsetReset.Earliest
};

using var consumer = new ConsumerBuilder<Ignore, string>(consumerConfig).Build();
consumer.Subscribe("your-target-topic"); // 替换为实际主题名

try
{
    var consumeResult = consumer.Consume(TimeSpan.FromSeconds(10));
    if (consumeResult != null)
    {
        // 示例断言:验证消息不为空
        Assert.IsNotNull(consumeResult.Message.Value);
    }
}
finally
{
    consumer.Close();
}

注意事项:

  • 务必将your-target-topic替换为实际需要读取的主题名称
  • 消费组ID使用测试专用标识,避免干扰其他环境的消费进度
  • 根据消息生产频率,调整消费超时时间和拉取次数

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 04:12:03