简介:这份资源面向计算机、人工智能及网络安全方向的学习者与研究人员,提供一套基于动态图神经网络的异常流量检测完整实现方案,用于解决传统静态拓扑方法在动态网络环境中准确率与效率不足的问题。压缩包共141个文件,约34.94MB,以60个Python源码文件为核心,辅以56个编译缓存、8个模型权重文件、4份CSV数据集、3份JSON配置及论文PDF、说明文档等,覆盖数据预处理、特征提取、模型训练到评估的完整链路。项目通过多层网络结构学习流量时间序列中的节点与边动态特性,并借助预训练模型参数支持快速推理。目前已有69人学习下载,适合作为毕业设计、课程设计或科研参考。读者可获得可运行的源码、逐行注释、项目说明与论文资料,便于理解DGNN工作机制、复现实验并迁移到实际网络监控场景。
1. 动态图神经网络做异常流量检测:为什么静态模型总在真实网络里翻车
线上流量从来不是一张固定的图。早上九点办公网里全是内网横向的 SMB 会话,凌晨三点变成几台机器往境外 IP 打小包,周末又换成 CDN 回源的长连接。你拿一周前的拓扑训一个 GCN,上线第一天 AUC 就掉到 0.6,这不是模型不行,是图本身在动。动态图神经网络(Dynamic Graph Neural Network,DGNN)要解决的就是这件事:把「谁在跟谁通信」建成随时间演化的图,节点是 IP 或主机,边是会话,边权是流量统计量,然后让模型在时间维度上学习正常与异常的差异。这套方法适合做内网东西向流量检测、DDoS 早期识别、C2 心跳发现,也适合已经有一批 NetFlow 或 Zeek 日志、想从规则阈值升级到模型的人。标题里那份 python 源码加论文的组合,本质就是给你一条从原始流量到图快照、再到 DGNN 训练和推理的完整链路,下面我按自己复现过的顺序把它拆开讲。
2. 从 pcap 到动态图:数据管线的四个关键决策
2.1 为什么选 NetFlow 聚合而不是逐包建图
逐包建图听起来最保真,但一个千兆口一天就是几亿条边,DGNN 的邻居聚合根本跑不动。常见做法是按固定时间窗做流聚合,比如 5 秒或 1 分钟一个快照,每个快照内把五元组相同的包合并成一条边。窗口太短,正常的长连接会被切碎,模型学到一堆碎片;窗口太长,DDoS 的突发特征被平均掉。我一般先用 10 秒窗跑一遍统计,看每个快照的边数是否稳定在几千到几万量级,这个量级单卡才能吃得下。
聚合时至少要保留这些字段:源 IP、目的 IP、源端口、目的端口、协议号、包数、字节数、持续时间、TCP 标志位计数。前五个用来确定边的两端,后四个是边特征。端口和协议号不要直接当数值喂进去,要做 one-hot 或 embedding,否则模型会以为 443 比 80 大。
import pandas as pd # 假设原始 flow 记录已经带时间戳,按 10 秒窗口聚合 def build_edge_features(flow_df, window='10s'): flow_df['ts'] = pd.to_datetime(flow_df['ts']) flow_df = flow_df.set_index('ts') # 同一窗口内同五元组聚合 grouped = flow_df.groupby([ pd.Grouper(freq=window), 'src_ip', 'dst_ip', 'src_port', 'dst_port', 'proto' ]).agg( pkt_count=('pkt_count', 'sum'), byte_count=('byte_count', 'sum'), duration=('duration', 'max'), syn_cnt=('syn_flag', 'sum'), rst_cnt=('rst_flag', 'sum') ).reset_index() return grouped这段代码的逻辑是先把时间戳设为索引,再用pd.Grouper按窗口切分,同一窗口内相同五元组的记录合并。pkt_count和byte_count用求和,duration取最大值,标志位计数求和。参数上window是最需要调的,建议从 5s、10s、30s、60s 各跑一遍,看后续模型验证集上的表现再定。
2.2 节点特征怎么构造才不丢信息
边建好了,节点特征不能只放一个 IP 字符串。我一般给每个 IP 节点拼三类特征:一是该窗口内作为源和作为目的的次数,二是收发字节比,三是历史窗口的滑动统计,比如过去 5 个窗口的平均出度。滑动统计是让模型感知「这个 IP 平时很安静,突然开始扫段」的关键。没有历史对比,单窗口的绝对数值很难区分正常的高流量和异常的突发。
# 为每个窗口的每个 IP 计算节点特征 def build_node_features(edge_df): src_stat = edge_df.groupby(['window', 'src_ip']).agg( out_deg=('dst_ip', 'nunique'), out_pkt=('pkt_count', 'sum'), out_byte=('byte_count', 'sum') ).reset_index().rename(columns={'src_ip': 'ip'}) dst_stat = edge_df.groupby(['window', 'dst_ip']).agg( in_deg=('src_ip', 'nunique'), in_pkt=('pkt_count', 'sum'), in_byte=('byte_count', 'sum') ).reset_index().rename(columns={'dst_ip': 'ip'}) node_df = pd.merge(src_stat, dst_stat, on=['window', 'ip'], how='outer').fillna(0) # 滑动窗口均值,窗口大小 5 node_df = node_df.sort_values(['ip', 'window']) for col in ['out_deg', 'out_pkt', 'out_byte', 'in_deg', 'in_pkt', 'in_byte']: node_df[f'{col}_ma5'] = node_df.groupby('ip')[col].transform( lambda x: x.rolling(5, min_periods=1).mean() ) return node_df这里out_deg是该 IP 在窗口内连了多少个不同目的 IP,扫描行为会让这个值异常高。_ma5后缀是 5 窗口滑动均值,用来给模型提供基线。注意fillna(0)之后要检查一下有没有 IP 只在目的侧出现,这类节点如果直接丢,会漏掉被攻击目标。
2.3 快照序列的切分与标签对齐
动态图模型吃的是快照序列,不是单张图。假设你用 10 秒窗,那 1 小时就是 360 个快照。训练时通常取连续 20 到 50 个快照作为一个样本,预测下一个或下一段快照里哪些边是异常的。标签来自你手里的告警日志或人工标注,常见做法是把告警时间点前后各扩一个窗口,避免边界漏标。
def make_sequences(node_df, edge_df, label_df, seq_len=30): windows = sorted(edge_df['window'].unique()) samples = [] for i in range(len(windows) - seq_len): seq_windows = windows[i:i+seq_len] next_window = windows[i+seq_len] seq_edges = edge_df[edge_df['window'].isin(seq_windows)] seq_nodes = node_df[node_df['window'].isin(seq_windows)] # 下一窗口的边标签 next_edges = edge_df[edge_df['window'] == next_window].copy() next_edges = next_edges.merge(label_df, on=['src_ip','dst_ip'], how='left') next_edges['label'] = next_edges['label'].fillna(0) samples.append((seq_nodes, seq_edges, next_edges)) return samplesseq_len是序列长度,太小模型看不到演化趋势,太大显存吃不消。我一般从 20 开始试,显存够就加到 50。标签对齐时how='left'保证正常边不会被丢掉,fillna(0)把未标注边当负样本,但这里有个坑:未标注不等于正常,后面避坑章节会细说。
2.4 图快照的存储格式与加载速度
快照多了以后,CSV 读写会成为瓶颈。我一般把每个快照存成单独的.npz或.pt,节点特征矩阵和边索引分开存。边索引用 COO 格式,两行分别是源节点在节点表中的下标和目的节点下标。加载时用内存映射,不要一次性全读进内存。
import numpy as np import torch def save_snapshot(window, node_feat, edge_index, edge_attr, path): np.savez_compressed( f'{path}/{window}.npz', node_feat=node_feat.astype(np.float32), edge_index=edge_index.astype(np.int64), edge_attr=edge_attr.astype(np.float32) ) def load_snapshot(window, path): data = np.load(f'{path}/{window}.npz') return ( torch.from_numpy(data['node_feat']), torch.from_numpy(data['edge_index']), torch.from_numpy(data['edge_attr']) )np.savez_compressed比 pickle 安全,也比 CSV 快一个数量级。edge_index用 int64 是 PyG 的默认要求,node_feat和edge_attr转 float32 省显存。如果快照数量超过几千,建议再建一个索引文件记录每个快照的节点数和边数,加载时按需分配。
3. DGNN 模型选型与训练:把演化信息真正用起来
3.1 为什么在流量场景里 GRU 加 GCN 比纯 Transformer 稳
纯 Transformer 做动态图不是不行,但流量图的节点数每个窗口都在变,位置编码很难对齐,而且自注意力的 O(N²) 在几万节点时直接爆显存。我试过在同样数据上跑纯 Transformer,验证集 loss 震荡得厉害,换成 GRU 沿时间维更新节点状态、GCN 做空间聚合之后,收敛曲线平滑很多。常见做法是每个快照先用 GCN 聚合邻居,再把聚合后的节点表示送进 GRU,GRU 的隐藏状态跨窗口传递。
import torch.nn as nn import torch.nn.functional as F from torch_geometric.nn import GCNConv class DGNN(nn.Module): def __init__(self, node_dim, edge_dim, hidden=64, num_layers=2): super().__init__() self.gcn_layers = nn.ModuleList([ GCNConv(node_dim if i == 0 else hidden, hidden) for i in range(num_layers) ]) self.edge_encoder = nn.Linear(edge_dim, hidden) self.gru = nn.GRU(hidden, hidden, batch_first=True) self.classifier = nn.Linear(hidden * 2, 1) def forward(self, node_feats, edge_indices, edge_attrs): # node_feats: [T, N, node_dim] T, N, _ = node_feats.shape h = None outputs = [] for t in range(T): x = node_feats[t] ei = edge_indices[t] for gcn in self.gcn_layers: x = F.relu(gcn(x, ei)) x = x.unsqueeze(0) # [1, N, hidden] out, h = self.gru(x, h) outputs.append(out.squeeze(0)) return torch.stack(outputs, dim=0)GCNConv的层数不要超过 3,流量图里两跳邻居已经能覆盖大部分通信关系,再深会过平滑。GRU的隐藏维度跟 GCN 输出保持一致,省去投影。classifier里hidden * 2是因为后面要把源节点和目的节点的表示拼接起来判断这条边是否异常。
3.2 边分类头的设计:源节点和目的节点怎么拼
异常检测最终要落到边上,所以分类头不能只看单个节点。我一般把边的两个端点表示取出来,拼上边特征,再过两层 MLP。拼接方式有 concat、hadamard 积、差值三种,实测 concat 最稳,hadamard 积在稀疏图上容易丢信息。
def edge_predict(self, node_emb, edge_index, edge_attr): src, dst = edge_index[0], edge_index[1] src_emb = node_emb[src] dst_emb = node_emb[dst] edge_emb = self.edge_encoder(edge_attr) combined = torch.cat([src_emb, dst_emb, edge_emb], dim=-1) return self.classifier(combined).squeeze(-1)edge_encoder把原始边特征投影到跟节点表示同维,避免量纲差异。classifier输出 logit,训练时用BCEWithLogitsLoss,推理时过 sigmoid 得到概率。如果正负样本极度不均衡,可以在 loss 里加pos_weight,一般设成负样本数除以正样本数。
3.3 训练循环里必须监控的三个量
训练 DGNN 不能只看 loss。我一般同时盯验证集 AUC、正样本召回率、以及每个 epoch 的梯度范数。AUC 看整体排序能力,召回率看漏报,梯度范数看有没有梯度爆炸。流量数据里正样本通常不到 1%,AUC 高但召回低是常态,这时候要调阈值或加 focal loss。
from sklearn.metrics import roc_auc_score, recall_score def train_epoch(model, loader, optimizer, pos_weight): model.train() total_loss = 0 for seq_nodes, seq_edges, next_edges in loader: optimizer.zero_grad() node_emb = model(seq_nodes, seq_edges, seq_edges) logits = model.edge_predict(node_emb[-1], next_edges.edge_index, next_edges.edge_attr) loss = F.binary_cross_entropy_with_logits( logits, next_edges.label, pos_weight=pos_weight ) loss.backward() torch.nn.utils.clip_grad_norm_(model.parameters(), max_norm=5.0) optimizer.step() total_loss += loss.item() return total_loss / len(loader)clip_grad_norm_的max_norm设 5.0 是经验值,梯度再大就说明学习率可能偏高。pos_weight用torch.tensor([neg/pos])传进去。验证时不要用训练阈值,要在验证集上扫一遍找最佳 F1 对应的阈值,再固定下来用于测试。
3.4 推理阶段的滑动窗口与在线更新
线上推理不是拿一个固定序列跑一次就完事。我一般维护一个长度为seq_len的滑动窗口队列,每来一个新快照就弹出最旧的、压入最新的,然后跑一次前向。模型参数可以定期用新数据微调,但不要每个窗口都更新,否则容易灾难性遗忘。
from collections import deque class OnlineDetector: def __init__(self, model, seq_len=30): self.model = model self.seq_len = seq_len self.buffer = deque(maxlen=seq_len) def update(self, node_feat, edge_index, edge_attr): self.buffer.append((node_feat, edge_index, edge_attr)) if len(self.buffer) < self.seq_len: return None seq_nodes = [b[0] for b in self.buffer] seq_edges = [b[1] for b in self.buffer] with torch.no_grad(): node_emb = self.model(seq_nodes, seq_edges, seq_edges) logits = self.model.edge_predict( node_emb[-1], edge_index, edge_attr ) return torch.sigmoid(logits)deque的maxlen自动控制窗口长度。推理时torch.no_grad()省显存。返回的概率要跟固定阈值比较,阈值来自验证集。如果线上分布漂移明显,可以每周用最近数据重新扫一次阈值。
4. 避坑与排查:复现时最容易翻车的五个地方
4.1 未标注边当负样本导致召回率虚低
现象是验证集 AUC 0.95 但召回只有 0.3。原因是标签只覆盖了告警涉及的边,大量正常边没标注,被fillna(0)当成负样本,模型学到「没标注就是正常」,遇到真正异常时反而不敢报。解决方法是做负样本采样,只从确认正常的边里抽,或者用 PU learning 的思路给未标注边低权重。
4.2 节点表每个窗口重建导致 embedding 对不上
现象是训练时 loss 正常,推理时同一 IP 的表示每次都不一样。原因是每个快照单独建节点表,IP 到下标的映射变了。解决方法是维护全局 IP 字典,所有快照共用一套下标,新 IP 追加到末尾,旧 IP 不删除。
4.3 时间窗边界把一条会话切成两半
现象是某些正常长连接被反复报异常。原因是 10 秒窗刚好切在会话中间,两个窗口各拿到一半包数,滑动均值波动大。解决方法是聚合时按会话开始时间归窗,或者把窗口边界做重叠,比如步长 5 秒、窗长 10 秒。
4.4 正负样本比失衡时 AUC 骗人
现象是 AUC 很高但实际漏报严重。原因是负样本太多,模型只要把少数正样本排前面 AUC 就好看,但阈值一卡就漏。解决方法是看 PR 曲线和召回率,用 focal loss 或调整pos_weight,不要只盯 AUC。
4.5 显存不够时盲目减层
现象是 OOM 后把 GCN 从 2 层减到 1 层,效果掉很多。原因是瓶颈在快照序列长度和节点数,不在层数。解决方法是先减seq_len,再用邻居采样,最后才考虑减层。邻居采样用 PyG 的NeighborLoader,每层采 10 到 15 个邻居就够。
5. 把模型压到线上:量化、阈值扫描与一个可复用的验证脚本
模型训完只是第一步,能不能上线看推理延迟和阈值稳定性。我一般先做动态量化,把 GRU 和 Linear 层转成 int8,延迟能降三到四成,AUC 掉不到 0.01。量化后再扫一遍阈值,因为量化会轻微改变输出分布。
import torch.quantization def quantize_model(model): model.eval() quantized = torch.quantization.quantize_dynamic( model, {nn.GRU, nn.Linear}, dtype=torch.qint8 ) return quantized def scan_threshold(model, val_loader): model.eval() all_probs, all_labels = [], [] with torch.no_grad(): for seq_nodes, seq_edges, next_edges in val_loader: node_emb = model(seq_nodes, seq_edges, seq_edges) logits = model.edge_predict( node_emb[-1], next_edges.edge_index, next_edges.edge_attr ) all_probs.append(torch.sigmoid(logits)) all_labels.append(next_edges.label) probs = torch.cat(all_probs).numpy() labels = torch.cat(all_labels).numpy() best_f1, best_th = 0, 0.5 for th in np.arange(0.1, 0.9, 0.01): pred = (probs > th).astype(int) f1 = f1_score(labels, pred, zero_division=0) if f1 > best_f1: best_f1, best_th = f1, th return best_th, best_f1quantize_dynamic只量化 GRU 和 Linear,GCN 的稀疏操作量化收益不大。scan_threshold在 0.1 到 0.9 之间以 0.01 步长扫,返回最佳 F1 和对应阈值。这个脚本我每次换数据都会跑一遍,阈值不要写死在代码里,放配置文件。
验证方法上,除了离线指标,我习惯再做一个「时间外推」测试:用前 70% 时间的数据训练,后 30% 测试,看性能衰减。如果衰减超过 15%,说明模型对时间漂移敏感,需要加滑动统计特征或缩短重训周期。这个习惯帮我提前发现过好几次线上翻车,比只看随机划分的验证集靠谱得多。希望帮到你。
本文还有配套的精品资源,点击获取