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

MongoDB聚合管道中如何处理双向流的键匹配与合并问题

MongoDB 双向流文档合并解决方案

要解决流方向反转导致无法合并的问题,核心是生成与流方向无关的分组键,无需依赖公私IP判断。以下提供两种可行方案:

方案一:利用$sortArray生成归一化IP/Port对(MongoDB 5.0+)

通过对IP和端口数组排序,无论原流方向是A→B还是B→A,排序后的结果一致,以此作为分组依据,实现逻辑简洁高效。

修改后的聚合管道

[
  // 保留原匹配条件
  {
    $match: {
      $or: [
        { _id: ObjectId('64227c692063fe9b27582cb1') },
        { _id: ObjectId('64227c692063fe9b27582ded') },
        { _id: ObjectId('64227cc62063fe9b2c3356f5') }
      ]
    }
  },
  // 保留原排序逻辑
  { $sort: { timestamp: -1 } },
  // 新增归一化字段:排序IP和端口对
  {
    $addFields: {
      sorted_ips: { $sortArray: { input: ["$sourceIPv4Address", "$destinationIPv4Address"], sortBy: 1 } },
      sorted_ports: { $sortArray: { input: ["$sourceTransportPort", "$destinationTransportPort"], sortBy: 1 } }
    }
  },
  // 基于归一化字段分组
  {
    $group: {
      _id: {
        src_ip: { $arrayElemAt: ["$sorted_ips", 0] },
        dst_ip: { $arrayElemAt: ["$sorted_ips", 1] },
        src_port: { $arrayElemAt: ["$sorted_ports", 0] },
        dst_port: { $arrayElemAt: ["$sorted_ports", 1] },
        protocol: "$protocol"
      },
      arr: { $addToSet: "$tcpFlags" },
      // 可选:保留原始文档,便于后续核查
      original_records: { $addToSet: "$$ROOT" }
    }
  }
]

方案二:条件判断交换字段(兼容低版本MongoDB)

通过字符串比较IP和端口的大小,动态交换src/dst字段,生成统一的归一化标识,兼容MongoDB 5.0以下版本。

修改后的聚合管道

[
  {
    $match: {
      $or: [
        { _id: ObjectId('64227c692063fe9b27582cb1') },
        { _id: ObjectId('64227c692063fe9b27582ded') },
        { _id: ObjectId('64227cc62063fe9b2c3356f5') }
      ]
    }
  },
  { $sort: { timestamp: -1 } },
  // 新增归一化字段:根据IP/Port大小动态交换
  {
    $addFields: {
      normalized_src_ip: {
        $cond: {
          if: { $gt: ["$sourceIPv4Address", "$destinationIPv4Address"] },
          then: "$destinationIPv4Address",
          else: "$sourceIPv4Address"
        }
      },
      normalized_dst_ip: {
        $cond: {
          if: { $gt: ["$sourceIPv4Address", "$destinationIPv4Address"] },
          then: "$sourceIPv4Address",
          else: "$destinationIPv4Address"
        }
      },
      normalized_src_port: {
        $cond: {
          if: { $or: [
            { $gt: ["$sourceIPv4Address", "$destinationIPv4Address"] },
            { $and: [
              { $eq: ["$sourceIPv4Address", "$destinationIPv4Address"] },
              { $gt: ["$sourceTransportPort", "$destinationTransportPort"] }
            ]}
          ]},
          then: "$destinationTransportPort",
          else: "$sourceTransportPort"
        }
      },
      normalized_dst_port: {
        $cond: {
          if: { $or: [
            { $gt: ["$sourceIPv4Address", "$destinationIPv4Address"] },
            { $and: [
              { $eq: ["$sourceIPv4Address", "$destinationIPv4Address"] },
              { $gt: ["$sourceTransportPort", "$destinationTransportPort"] }
            ]}
          ]},
          then: "$sourceTransportPort",
          else: "$destinationTransportPort"
        }
      }
    }
  },
  {
    $group: {
      _id: {
        src_ip: "$normalized_src_ip",
        dst_ip: "$normalized_dst_ip",
        src_port: "$normalized_src_port",
        dst_port: "$normalized_dst_port",
        protocol: "$protocol"
      },
      arr: { $addToSet: "$tcpFlags" },
      original_records: { $addToSet: "$$ROOT" }
    }
  }
]

可选:生成无方向哈希值

如果需要生成唯一的哈希标识,可以在归一化后添加$toHashedField字段,示例如下:

{
  $addFields: {
    // 基于归一化字段生成哈希
    flow_hash: {
      $toHashedField: {
        input: {
          $concat: [
            "$normalized_src_ip", ":", { $toString: "$normalized_src_port" },
            "-", "$normalized_dst_ip", ":", { $toString: "$normalized_dst_port" },
            "/", { $toString: "$protocol" }
          ]
        }
      }
    }
  }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 06:35:46