如何连接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
相关产品推荐
相关产品推荐

