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

.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.14 13:45:55