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

如何在DuckDB中按特定条件用Partition By合并相似SQL表行

DuckDB合并相似TCP流数据

需求概述

需合并满足以下条件的TCP流记录:

  • 核心四字段完全匹配:sourceIPAddress、destinationIPAddress、sourceTransportPort、destinationTransportPort
  • 后续流的flowStartMilliseconds ≤ 前一条流的flowEndMilliseconds + 3600000(1小时阈值)

合并规则:

  • 对bytes和packets字段求和
  • 取分组内第一条记录的flowStartMilliseconds作为合并后起始时间
  • 取分组内最后一条记录的flowEndMilliseconds作为合并后结束时间
  • 保留protocol、serviceName等一致字段的值

原查询问题

你当前的查询仅计算了每条流与前一条的关联条件,但未完成实际的合并操作。需要通过会话分组的方式,将符合条件的连续流归为同一组,再对组内数据进行聚合。

最终实现查询

WITH flow_groups AS (
    SELECT
        sourceIPAddress,
        destinationIPAddress,
        sourceTransportPort,
        destinationTransportPort,
        protocol,
        serviceName,
        flowStartMilliseconds,
        flowEndMilliseconds,
        bytes,
        packets,
        -- 生成分组标识:当前流与前一条不满足合并条件时,分组ID递增
        SUM(CASE
            WHEN lag(flowEndMilliseconds) OVER (
                PARTITION BY sourceIPAddress, destinationIPAddress, sourceTransportPort, destinationTransportPort
                ORDER BY flowStartMilliseconds
            ) + 3600 * 1000 >= flowStartMilliseconds THEN 0
            ELSE 1
        END) OVER (
            PARTITION BY sourceIPAddress, destinationIPAddress, sourceTransportPort, destinationTransportPort
            ORDER BY flowStartMilliseconds ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
        ) AS group_id
    FROM read_parquet('//path/*/*/*/*/*/*.parquet', hive_partitioning = 1)
    WHERE protocol = 'TCP'
)
SELECT
    sourceIPAddress,
    destinationIPAddress,
    sourceTransportPort,
    destinationTransportPort,
    protocol,
    serviceName,
    MIN(flowStartMilliseconds) AS flowStartMilliseconds,
    MAX(flowEndMilliseconds) AS flowEndMilliseconds,
    SUM(bytes) AS total_bytes,
    SUM(packets) AS total_packets
FROM flow_groups
GROUP BY
    sourceIPAddress,
    destinationIPAddress,
    sourceTransportPort,
    destinationTransportPort,
    protocol,
    serviceName,
    group_id
ORDER BY
    sourceIPAddress,
    destinationIPAddress,
    sourceTransportPort,
    destinationTransportPort,
    flowStartMilliseconds;

逻辑说明

  1. 分组阶段(flow_groups):
    • 使用lag()窗口函数获取同组(核心四字段相同)前一条流的结束时间,判断当前流是否符合合并条件
    • 通过累加CASE结果生成group_id:每遇到不满足合并条件的流,分组ID加1,确保连续符合条件的流被归为同一组
  2. 聚合阶段:
    • 按核心四字段、protocol、serviceName和group_id分组
    • 用MIN()取分组最早的起始时间,MAX()取分组最晚的结束时间,SUM()计算总字节数和总包数

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 09:10:18