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

大型网络流数据集生成UUID时Jupyter Notebook崩溃求助

大数据集下网络流U-ID生成与合并优化问题

数据集与规则说明

数据集包含网络流详情,字段如下:

  • connection_type(字符串): 连接类型
  • connection_category(字符串): 连接分类
  • bytes_in: 流入字节数
  • bytes_out: 流出字节数
  • flow_start: 流起始时间(epoch时间戳)
  • flow_end: 流结束时间(epoch时间戳)
  • flow_id: 由客户端IP、服务端IP、客户端端口、服务端端口哈希生成,相同IP端口组合对应相同flow_id

新连接判断规则:同一flow_id下,若当前记录的flow_start不等于上一条记录的flow_end,则判定为新连接。

数据集示例:

flow_idflow_startflow_endbytes_outbytes_inconnection_typeconnection_category
12345678916971333121697133317124342ab
123456789169713331716971333423450ab
123456789169713340016971334540100xy
12345678916971334541697133499200200xy
12345678916971334991697133600015xy

需求目标

  1. 生成U-ID:同一连接(连续的flow_start等于上一条flow_end)的记录共享相同U-ID;新连接生成唯一U-ID。U-ID由flow_id与该连接的起始flow_start哈希生成。
  2. 基于U-ID合并记录,累计bytes_out和bytes_in的值。

目标结果示例:

U-IDflow_idflow_startflow_endbytes_outbytes_inconnection_typeconnection_category
8888812345678916971333121697133317124342ab
88888123456789169713331716971333423450ab
99999123456789169713340016971334540100xy
9999912345678916971334541697133499200200xy
9999912345678916971334991697133600015xy

现有问题

提供的Python代码在小数据集上运行正常,但处理大型数据集时Jupyter Notebook会崩溃。

现有代码

import uuid

def buildId(data, indices):
    flow_changed = False
    previous_flow_st = 0
    previous_flow_et = 0
    index = []
    flow_id = 0
    flow_st = 0
    flow_et = 0

    start_time = 0
    end_time = 0
    p = 0
    q = 0

    for i, j in enumerate(data):
        index.append(j[0]) if j[0] not in index else None
        flow_id = j[1]
        flow_st = j[2]
        flow_et = j[3]

        if i == 0:
            start_time = flow_st
            end_time = flow_et
        else:
            if flow_st == previous_flow_et:
                # ------------------------- same flow
                end_time = flow_et
                if i == len(data) - 1:
                    uid = str(uuid.uuid3(uuid.NAMESPACE_DNS, str(str(flow_id) + str(start_time))))
                    indices[uid] = (index, end_time - start_time)
            elif flow_st != previous_flow_et:
                # ------------------------- new flow           
                # compute the uid for the previous flow
                end_time = previous_flow_et
                # remove the current index entry, since this does not belong to the previous flow
                index.remove(j[0])
                uid = str(uuid.uuid3(uuid.NAMESPACE_DNS, str(str(flow_id) + str(start_time))))
                indices[uid] = (index, end_time - start_time)

                # update stuff for the new flow
                start_time = flow_st
                end_time = flow_et
                index = []
                flow_changed = True

        previous_flow_st = flow_st
        previous_flow_et = flow_et
    
    if flow_changed is False:
        uid = str(uuid.uuid3(uuid.NAMESPACE_DNS, str(str(flow_id) + str(start_time))))
        indices[uid] = (index, end_time - start_time)

indices = {} # is dictionary key: UUID value:duration of flow
buildId(df_flow1[df_flow1['flow_id'].isin([id_])][['index', 'flow_id', 'flow_start', 'flow_end']].values.tolist(), indices)

优化方案

原代码的核心问题是用Python循环逐行处理,且频繁进行列表的增删操作,在大数据集下内存和时间开销极大。以下是基于pandas矢量化操作的优化方案,性能和内存占用均有大幅提升:

优化思路

  1. 先按flow_id和flow_start排序,确保同一流的记录按时间顺序排列
  2. 用分组+shift操作标记新连接的起始行
  3. 用累积求和生成每个连接的组ID,方便后续生成U-ID
  4. 基于组ID生成U-ID,最后分组聚合统计

完整优化代码

import pandas as pd

# 1. 确保数据集按flow_id和flow_start排序,保证时间顺序正确
df = df_flow1.sort_values(by=['flow_id', 'flow_start']).reset_index(drop=True)

# 2. 按flow_id分组,标记新连接:当前flow_start不等于上一条的flow_end则为新连接
df['is_new_connection'] = df.groupby('flow_id').apply(
    lambda x: x['flow_start'] != x['flow_end'].shift(1)
).reset_index(drop=True)
# 每组第一行默认是新连接,填充空值
df['is_new_connection'] = df['is_new_connection'].fillna(True)

# 3. 生成每个flow_id下的连接组ID,同一连接的组ID相同
df['connection_group'] = df.groupby('flow_id')['is_new_connection'].cumsum()

# 4. 生成U-ID:用flow_id和连接起始时间的哈希值,替代uuid更高效
def generate_uid(group):
    flow_id = group['flow_id'].iloc[0]
    start_time = group['flow_start'].min()
    # 若需要字符串格式的U-ID,可改用hashlib生成固定长度哈希,或直接拼接字符串
    return hash(f"{flow_id}_{start_time}")

df['U-ID'] = df.groupby(['flow_id', 'connection_group']).apply(generate_uid).reset_index(drop=True)

# 5. 基于U-ID合并统计,累计字节数并保留其他字段
merged_df = df.groupby('U-ID').agg(
    flow_id=('flow_id', 'first'),
    connection_start=('flow_start', 'min'),
    connection_end=('flow_end', 'max'),
    total_bytes_out=('bytes_out', 'sum'),
    total_bytes_in=('bytes_in', 'sum'),
    connection_type=('connection_type', 'first'),  # 假设同一连接类型一致,不一致可改用众数
    connection_category=('connection_category', 'first')
).reset_index()

优化点说明

  • 矢量化操作:全程利用pandas的分组、shift、cumsum等矢量化函数,避免Python循环,处理速度提升几个数量级
  • 内存优化:不需要维护额外的列表和字典,所有操作基于原DataFrame,内存占用大幅降低
  • 高效哈希:用Python内置hash函数替代uuid,生成速度更快;若需要全局唯一的字符串U-ID,可改用hashlib.md5(f"{flow_id}_{start_time}".encode()).hexdigest()
  • 稳定分组:pandas的分组引擎经过优化,处理百万级以上数据仍能保持稳定

内容的提问来源于stack exchange,提问作者Anonymous Raze

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 08:09:53