大型网络流数据集生成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_id | flow_start | flow_end | bytes_out | bytes_in | connection_type | connection_category |
|---|---|---|---|---|---|---|
| 123456789 | 1697133312 | 1697133317 | 124 | 342 | a | b |
| 123456789 | 1697133317 | 1697133342 | 345 | 0 | a | b |
| 123456789 | 1697133400 | 1697133454 | 0 | 100 | x | y |
| 123456789 | 1697133454 | 1697133499 | 200 | 200 | x | y |
| 123456789 | 1697133499 | 1697133600 | 0 | 15 | x | y |
需求目标
- 生成
U-ID:同一连接(连续的flow_start等于上一条flow_end)的记录共享相同U-ID;新连接生成唯一U-ID。U-ID由flow_id与该连接的起始flow_start哈希生成。 - 基于
U-ID合并记录,累计bytes_out和bytes_in的值。
目标结果示例:
| U-ID | flow_id | flow_start | flow_end | bytes_out | bytes_in | connection_type | connection_category |
|---|---|---|---|---|---|---|---|
| 88888 | 123456789 | 1697133312 | 1697133317 | 124 | 342 | a | b |
| 88888 | 123456789 | 1697133317 | 1697133342 | 345 | 0 | a | b |
| 99999 | 123456789 | 1697133400 | 1697133454 | 0 | 100 | x | y |
| 99999 | 123456789 | 1697133454 | 1697133499 | 200 | 200 | x | y |
| 99999 | 123456789 | 1697133499 | 1697133600 | 0 | 15 | x | y |
现有问题
提供的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矢量化操作的优化方案,性能和内存占用均有大幅提升:
优化思路
- 先按
flow_id和flow_start排序,确保同一流的记录按时间顺序排列 - 用分组+shift操作标记新连接的起始行
- 用累积求和生成每个连接的组ID,方便后续生成U-ID
- 基于组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
相关产品推荐
相关产品推荐

