.NET Core中用Elasticsearch 8.17.1实现多分组多聚合查询
.NET Elasticsearch 8.17.1 分组聚合查询实现方案
需求明确
- 分组字段:
sessionId、destinationPort、sourcePort、protocol、destinationPhysicalAddress、frameProtocol、projectId、sourcePhysicalAddress - 聚合需求:
- 求和
packageNumber,别名TotalPackageNumber - 求和
length,别名TotalLength - 取
timestamp的最小值,别名Timestamp
- 求和
索引字段注意事项
索引中sessionId、protocol、destinationPhysicalAddress、frameProtocol、sourcePhysicalAddress为text类型,分组时必须使用其keyword子字段(如sessionId.keyword),否则会因分词导致分组逻辑失效。数值类型字段(destinationPort、sourcePort、projectId)可直接用于分组。
.NET 客户端代码实现
using Elastic.Clients.Elasticsearch; using Elastic.Clients.Elasticsearch.Aggregations; using Elastic.Transport; // 初始化Elasticsearch客户端 var settings = new ElasticsearchClientSettings(new Uri("http://your-es-host:9200")) .Authentication(new BasicAuthentication("your-username", "your-password")); // 按需配置认证信息 var client = new ElasticsearchClient(settings); // 构建分组聚合查询 var searchResponse = client.Search<SessionDocument>(s => s .Index("sessions") .Size(0) // 无需返回原始文档,仅获取聚合结果 .Aggregations(a => a .Terms("group_by_fields", t => t .Composite(c => c .Sources( new TermsCompositeSource("sessionId", sc => sc.Field(f => f.SessionId.Suffix("keyword"))), new TermsCompositeSource("destinationPort", sc => sc.Field(f => f.DestinationPort)), new TermsCompositeSource("sourcePort", sc => sc.Field(f => f.SourcePort)), new TermsCompositeSource("protocol", sc => sc.Field(f => f.Protocol.Suffix("keyword"))), new TermsCompositeSource("destinationPhysicalAddress", sc => sc.Field(f => f.DestinationPhysicalAddress.Suffix("keyword"))), new TermsCompositeSource("frameProtocol", sc => sc.Field(f => f.FrameProtocol.Suffix("keyword"))), new TermsCompositeSource("projectId", sc => sc.Field(f => f.ProjectId)), new TermsCompositeSource("sourcePhysicalAddress", sc => sc.Field(f => f.SourcePhysicalAddress.Suffix("keyword"))) ) ) .Aggregations(aggs => aggs .Sum("TotalPackageNumber", sum => sum.Field(f => f.PackageNumber)) .Sum("TotalLength", sum => sum.Field(f => f.Length)) .Min("Timestamp", min => min.Field(f => f.Timestamp)) ) ) ) ); // 处理查询结果 if (searchResponse.IsValidResponse) { var compositeAgg = searchResponse.Aggregations.Terms("group_by_fields")?.Composite; if (compositeAgg != null) { foreach (var bucket in compositeAgg.Buckets) { // 提取分组字段值 var sessionId = bucket.Key["sessionId"]?.ToString(); var destinationPort = bucket.Key["destinationPort"]?.ToInt64(); var sourcePort = bucket.Key["sourcePort"]?.ToInt64(); var protocol = bucket.Key["protocol"]?.ToString(); // 提取聚合结果 var totalPackages = bucket.Aggregations.Sum("TotalPackageNumber")?.Value; var totalLength = bucket.Aggregations.Sum("TotalLength")?.Value; var minTimestamp = bucket.Aggregations.Min("Timestamp")?.Value; // 业务逻辑处理示例 Console.WriteLine($"SessionId: {sessionId}, 总包数: {totalPackages}, 总长度: {totalLength}, 最早时间: {minTimestamp}"); } } } else { Console.WriteLine($"查询失败: {searchResponse.DebugInformation}"); } // 对应文档实体类 public class SessionDocument { public long DestinationPort { get; set; } public long SourcePort { get; set; } public string FrameProtocol { get; set; } public long Length { get; set; } public string SessionId { get; set; } public string SourcePhysicalAddress { get; set; } public long PackageNumber { get; set; } public string DestinationIp { get; set; } public string Protocol { get; set; } public string SourceIp { get; set; } public string DestinationPhysicalAddress { get; set; } public long ProjectId { get; set; } public DateTimeOffset Timestamp { get; set; } }
对应Elasticsearch DSL查询(用于直接验证)
POST /sessions/_search { "size": 0, "aggs": { "group_by_fields": { "composite": { "sources": [ { "sessionId": { "terms": { "field": "sessionId.keyword" } } }, { "destinationPort": { "terms": { "field": "destinationPort" } } }, { "sourcePort": { "terms": { "field": "sourcePort" } } }, { "protocol": { "terms": { "field": "protocol.keyword" } } }, { "destinationPhysicalAddress": { "terms": { "field": "destinationPhysicalAddress.keyword" } } }, { "frameProtocol": { "terms": { "field": "frameProtocol.keyword" } } }, { "projectId": { "terms": { "field": "projectId" } } }, { "sourcePhysicalAddress": { "terms": { "field": "sourcePhysicalAddress.keyword" } } } ] }, "aggs": { "TotalPackageNumber": { "sum": { "field": "packageNumber" } }, "TotalLength": { "sum": { "field": "length" } }, "Timestamp": { "min": { "field": "timestamp" } } } } } }
说明:选用
composite聚合而非多层terms聚合,是因为前者支持分页获取大量分组结果,避免内存溢出问题;若确定结果集较小,也可替换为多层terms聚合,但composite是更通用的生产级方案。
内容的提问来源于stack exchange,提问作者Vicente García Diez
相关产品推荐
相关产品推荐

