news 2026/10/11 7:03:02

基于Flume+Spark Streaming的实时日志入侵检测系统

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
基于Flume+Spark Streaming的实时日志入侵检测系统

简介:本资源是一个基于Flume、Spark与Flask构建的分布式实时日志分析与入侵检测系统,面向大数据初学者、毕业设计学生及安全分析实践者,聚焦于Web服务器日志(如access_log)的采集、流式处理、异常行为识别与可视化展示,可支撑课程设计、毕设开发与小型安全监控场景。压缩包共108个文件,含14个编译后class文件、16个导出配置export、7个PNG图表、5个说明类txt与1个核心data样本,辅以Scala/Java源码、Flask Web接口、Spark Streaming作业及配置文件(conf、properties),整体18.91MB,结构完整、模块清晰,便于分层理解数据流转链路。目前已有325人学习下载,资源经本地全链路编译验证,附详细环境配置文档与助教审定保障,提供从日志接入、特征提取、规则匹配到结果呈现的端到端可运行方案,特别适合夯实实时计算与安全日志分析双重能力。

1. 这不是又一个“日志可视化看板”:它用 Flume 实时接 Apache/Nginx 日志,Spark Streaming 做滑动窗口统计 + 规则匹配(含 SQL 式 IOC 提取),Flask 暴露 REST API + 简洁 Web 界面,真正跑在三节点伪分布式环境里、能复现入侵行为链的毕业设计级系统

你可能已经下载过十几个“日志分析系统”压缩包,解压后发现是 Flask 写的静态页面 + 本地 CSV 加载 + 三个 matplotlib 图表——那叫日志展示,不叫实时分析。而这个.zip包里的东西,从 Flume agent 配置开始就踩在真实生产逻辑上:它监听/var/log/nginx/access.log的 tail -F 行为,通过 Avro Sink 推给 Spark Streaming 的 ReceiverInputDStream;Spark 侧不是简单 count(),而是用windowDuration=30s, slideDuration=10s构建滑动窗口,对每个窗口内 IP 做请求频次、404 比率、User-Agent 异常熵值、SQL 注入关键词(如union select,' or '1'='1)正则命中数四维聚合;结果存入内存 MapState 并触发阈值告警(比如 10s 内单 IP 404 > 50 且含注入词),再经 Flask REST 接口推送到前端 ECharts 动态拓扑图。它不依赖 HDFS 高可用,但明确要求 Spark Standalone 三节点(master + 2 worker),所有配置文件(flume-conf.properties、spark-defaults.conf、application.py)都带注释标注端口/路径/序列化方式。适合需要答辩演示“数据从产生到告警闭环”的本科毕设,也足够作为校招项目讲清楚流式架构分层(采集层→计算层→服务层)的边界与协作。


2. 从零搭起三节点 Spark Standalone 集群:为什么必须用spark-submit --master spark://master:7077而不是 local[*],以及如何绕过 YARN 依赖直接跑通 Streaming

2.1 为什么毕业设计必须坚持 Standalone 模式?避开 Hadoop 生态的“伪分布式陷阱”

很多同学在毕设里写“基于 Spark 的实时分析”,实际代码里全是SparkContext(sc = SparkContext("local[4]", "LogAnalysis"))——这根本不算分布式,连单机多线程都算不上,只是把 RDD 分片在本机 CPU 核上跑。而本系统要求的Standalone 模式,是 Spark 官方推荐的轻量级集群管理器,它不依赖 Hadoop 生态(HDFS/YARN),却能真实体现 Driver 与 Executor 的网络通信、Shuffle 数据拉取、Executor 内存隔离等核心机制。尤其对 Streaming 场景,StreamingContext必须连接到集群 Master(spark://master:7077),否则ReceiverInputDStream无法在 Worker 上启动 NetcatReceiver 或 FlumePollingReceiver。本系统spark-streaming-flume_2.12-3.3.2.jar已预编译进 lib 目录,但若你用 Spark 3.4+,需手动替换为对应 Scala 版本的 jar(见后文避坑章)。

提示:本系统默认使用 Spark 3.3.2 + Scala 2.12,所有 pom.xml 和 build.sbt 中 scalaVersion 必须严格匹配,否则ClassNotFoundException: org.apache.spark.streaming.flume.FlumeUtils会直接中断启动。

2.2 三节点部署实操:master、worker1、worker2 的角色分工与关键配置项

我们以三台虚拟机(或同一物理机的三个终端)模拟最小可行集群。假设 IP 分配如下:

节点主机名IP 地址角色
mastermaster192.168.56.101Spark Master + Flume Agent + Flask Web Server
worker1node1192.168.56.102Spark Worker + Flume Agent(可选)
worker2node2192.168.56.103Spark Worker

第一步:统一安装 Spark 并配置环境变量
在三台机器上执行(以 Ubuntu 22.04 为例):

# 下载预编译版(注意 Scala 版本!) wget https://downloads.apache.org/spark/spark-3.3.2/spark-3.3.2-bin-hadoop3.tgz tar -xzf spark-3.3.2-bin-hadoop3.tgz sudo mv spark-3.3.2-bin-hadoop3 /opt/spark echo 'export SPARK_HOME=/opt/spark' >> ~/.bashrc echo 'export PATH=$SPARK_HOME/bin:$PATH' >> ~/.bashrc source ~/.bashrc

第二步:配置 master 节点的$SPARK_HOME/conf/spark-env.sh

# 复制模板并编辑 cp $SPARK_HOME/conf/spark-env.sh.template $SPARK_HOME/conf/spark-env.sh nano $SPARK_HOME/conf/spark-env.sh

填入以下内容(关键参数已加注释):

#!/usr/bin/env bash # 必须指定 JAVA_HOME,否则 start-master.sh 报错 export JAVA_HOME=/usr/lib/jvm/java-11-openjdk-amd64 # Master 绑定地址,不能写 0.0.0.0(安全限制),必须写本机 IP export SPARK_MASTER_HOST=192.168.56.101 # Master Web UI 端口,默认 8080,若被占用可改 export SPARK_MASTER_WEBUI_PORT=8080 # Master RPC 端口,Spark Streaming 必须通过此端口注册 Receiver export SPARK_MASTER_PORT=7077 # JVM 内存参数,避免小内存机器 OOM export SPARK_DAEMON_MEMORY=1g

第三步:配置 worker 节点的$SPARK_HOME/conf/spark-env.sh

# 在 node1 和 node2 上执行 nano $SPARK_HOME/conf/spark-env.sh

填入:

#!/usr/bin/env bash export JAVA_HOME=/usr/lib/jvm/java-11-openjdk-amd64 # 指向 master 的 RPC 地址,格式为 spark://<host>:<port> export SPARK_MASTER=spark://192.168.56.101:7077 # Worker 内存分配,建议至少 2g,Streaming 场景需更多堆外内存 export SPARK_WORKER_MEMORY=2g # Worker 核心数,根据 CPU 核心数设置,避免超卖 export SPARK_WORKER_CORES=2 # Worker Web UI 端口,每台机器必须唯一 export SPARK_WORKER_WEBUI_PORT=8081 # node1 用 8081,node2 改为 8082

第四步:启动集群并验证

在 master 上执行:

$SPARK_HOME/sbin/start-master.sh # 查看日志确认是否绑定 7077 端口 tail -f $SPARK_HOME/logs/spark-*-org.apache.spark.deploy.master.Master-*.out

在 node1 和 node2 上分别执行:

$SPARK_HOME/sbin/start-worker.sh spark://192.168.56.101:7077 # 查看日志确认是否注册成功 tail -f $SPARK_HOME/logs/spark-*-org.apache.spark.deploy.worker.Worker-*.out

此时访问http://192.168.56.101:8080,应看到 Web UI 显示 1 个 Alive Workers,Total Cores=4,Total Memory=4g。这是后续所有 Streaming 任务能运行的前提——如果 UI 里看不到 Worker,spark-submit会 fallback 到 local 模式,整个“分布式”就失效了。

2.3 提交 Streaming 任务:spark-submit的必填参数与常见失败原因

进入系统源码目录src/main/python/streaming/,执行提交命令:

spark-submit \ --master spark://192.168.56.101:7077 \ --deploy-mode client \ --driver-memory 2g \ --executor-memory 2g \ --executor-cores 2 \ --jars $SPARK_HOME/jars/spark-streaming-flume_2.12-3.3.2.jar,$SPARK_HOME/jars/spark-streaming-flume-sink_2.12-3.3.2.jar \ --py-files $SPARK_HOME/python/lib/pyspark.zip,$SPARK_HOME/python/lib/py4j-0.10.9.5-src.zip \ log_streaming_analyzer.py

参数详解与血泪经验:

  • --master:必须是spark://host:port,不能是yarn或local;
  • --deploy-mode client:必须用 client 模式,因为 Streaming 任务需要 Driver 持续运行接收 Flume 数据,cluster 模式下 Driver 退出即任务终止;
  • --jars:显式指定 Flume 集成 jar,Spark 3.3.2 对应的 artifactId 是spark-streaming-flume_2.12,版本号必须与 Spark 一致,否则NoClassDefFoundError;
  • --py-files:PySpark 依赖包,路径必须准确,py4j版本需与 Spark 打包版本一致(3.3.2 对应py4j-0.10.9.5);
  • log_streaming_analyzer.py:主程序,内部StreamingContext初始化时指定了batchDuration=10秒,即每 10 秒处理一个 DStream Batch。

若提交后报错Failed to connect to master 192.168.56.101:7077,请立即检查:
① master 是否真的在运行(jps应看到 Master 进程);
② 防火墙是否放行 7077 端口(sudo ufw allow 7077);
③/etc/hosts中是否将192.168.56.101 master正确映射(Spark 内部用主机名通信)。


3. Flume 实时采集:从spooldir到exec再到avro,为什么本系统强制使用 Avro Source + Sink 架构

3.1 三种采集方式对比:为什么spooldir和exec不适合入侵检测场景

Flume 提供多种 Source,但并非都适用于安全分析:

Source 类型原理是否实时是否可靠是否支持断点续传是否适合本系统
spooldir监控目录,文件移动后读取❌ 秒级延迟(轮询间隔)✅ 文件移动原子性保证✅ 通过 .COMPLETED 标记否:无法捕获正在写入的 access.log
exectail -F命令输出✅ 准实时❌ 进程崩溃即丢失数据❌ 无状态记录否:入侵行为常发生在日志滚动瞬间,tail 可能漏掉最后一行
avroAvro RPC 协议接收数据✅ 毫秒级✅ TCP ACK 保障✅ Flume Channel(FileChannel)持久化✅ 唯一满足“不丢、不重、低延迟”的方案

本系统采用Avro Source(在 Spark 端) + Avro Sink(在 Flume 端)的反向架构:Flume Agent 不主动拉日志,而是作为客户端,将解析后的日志事件(JSON 格式)通过 Avro 协议 Push 到 Spark Streaming 的 AvroSource。这种设计规避了传统exec tail的进程稳定性问题,也绕开了spooldir的文件移动延迟。

3.2 Flume Agent 配置详解:flume-conf.properties中的四个生死参数

进入conf/flume-conf.properties,核心配置如下:

# Agent 名称(必须与启动命令一致) a1.sources = r1 a1.sinks = k1 a1.channels = c1 # Source:使用 exec + tail -F,但加了关键增强 a1.sources.r1.type = exec a1.sources.r1.command = tail -F /var/log/nginx/access.log a1.sources.r1.shell = /bin/bash -c # 【关键】每行日志必须以 \n 结尾,否则 Flume 会粘包 a1.sources.r1.restart = true a1.sources.r1.restart Throttle = 10000 # Channel:必须用 FileChannel 保证可靠性 a1.channels.c1.type = file a1.channels.c1.checkpointDir = /var/flume/checkpoint a1.channels.c1.dataDirs = /var/flume/data # 【关键】事务容量必须 >= source batch size,否则丢数据 a1.channels.c1.transactionCapacity = 1000 a1.channels.c1.capacity = 1000000 # Sink:Avro Sink,Push 到 Spark 的 AvroSource a1.sinks.k1.type = avro a1.sinks.k1.hostname = 192.168.56.101 a1.sinks.k1.port = 41414 # 【关键】批次大小必须与 Spark Streaming 的 batchDuration 匹配 a1.sinks.k1.batch-size = 100 # 【关键】超时设置,避免网络抖动导致阻塞 a1.sinks.k1.connect-timeout = 20000 a1.sinks.k1.request-timeout = 20000 # Bind source → channel → sink a1.sources.r1.channels = c1 a1.sinks.k1.channel = c1

参数生死线解释:

  • transactionCapacity=1000:Channel 每次事务最多写入 1000 条事件。若batch-size=100,则每 10 次 Sink 操作才触发一次 Channel 事务,降低 IO 压力;但若设太小(如 100),高并发下 Channel 会成为瓶颈;
  • batch-size=100:Sink 每次向 Avro Server 发送 100 条日志。该值需与 Spark Streaming 的batchDuration=10s对齐——假设 Nginx 每秒 50 条日志,则 10s 产生 500 条,batch-size=100意味着 5 次网络请求,比batch-size=1(500 次)更高效;
  • connect-timeout=20000:连接超时 20 秒。若 Spark AvroSource 未启动,Flume 会重试 20 秒后报错,而非无限等待;
  • restart=true:确保tail -F进程崩溃后自动重启,这是execSource 唯一的容错手段。

3.3 Spark 端 AvroSource 启动:FlumeUtils.createStream()的隐藏约束

在log_streaming_analyzer.py中,关键初始化代码为:

from pyspark.streaming.flume import FlumeUtils # 创建 StreamingContext,batchDuration=10秒 ssc = StreamingContext(sc, 10) # 【关键】AvroSource 绑定在 41414 端口,必须与 flume-conf.properties 中 sink.port 一致 flumeStream = FlumeUtils.createStream( ssc, hostname="192.168.56.101", # Spark Driver 所在机器 port=41414, # 必须与 Flume sink.port 完全相同 protoClass="org.apache.flume.source.avro.AvroFlumeEvent", # 固定值 transform=lambda x: x # 默认返回原始 event,后续需解析 body )

注意:hostname参数不是 Flume Agent 的地址,而是Spark Driver 绑定的地址。因为 AvroSource 是 Spark 启动的 Netty Server,Flume Sink 是 Client,所以hostname必须是 Spark Driver 所在机器(即 master)能被 Flume 访问到的 IP。若填localhost,Flume 会尝试连127.0.0.1:41414,而该端口只在 master 本机监听,node1/node2 无法访问,导致Connection refused。

提示:启动前务必在 master 上执行sudo lsof -i :41414确认端口空闲;若被占用,修改flume-conf.properties和 Python 代码中的 port 为 41415,并同步更新。


4. 入侵检测规则引擎:从硬编码正则到可热加载 JSON 规则库,以及 Spark 中读取 JSON 的正确姿势

4.1 为什么不用 ML 模型?基于规则的轻量级检测更适合毕业设计场景

当前系统未集成孤立森林或 LSTM 等模型,原因很务实:
①数据冷启动问题:Nginx access.log 缺乏标签(正常/攻击),无法监督训练;
②特征工程成本高:IP 地理位置、UA 设备指纹、Referer 可信度等需外部 API,增加部署复杂度;
③可解释性刚需:答辩时老师问“为什么判定这个 IP 是攻击者?”,展示正则r"(?i)union\s+select"比解释 embedding cosine similarity 更直观。

因此,系统采用四维规则组合:

  • 高频扫描:10s 内单 IP 请求 > 100 次;
  • 异常状态码:404 比率 > 70%;
  • UA 异常熵:len(set([ua[:3] for ua in uas])) < 2(大量 UA 以curl/python-requests开头);
  • IOC 关键词:SQLi/XSS/Path Traversal 三类共 23 个正则(见rules/ioc_rules.json)。

4.2 Spark 中读取 JSON 规则:spark.read.json()的陷阱与broadcast正确用法

规则文件rules/ioc_rules.json格式如下:

[ {"name": "SQLi_UNION_SELECT", "pattern": "(?i)union\\s+select", "severity": "high"}, {"name": "XSS_SCRIPT_TAG", "pattern": "<script[^>]*>", "severity": "medium"}, {"name": "PATH_TRAV_ASLASHDOT", "pattern": "\\.{2}/", "severity": "high"} ]

错误做法(直接在 map 中读取文件):

# ❌ 千万不要这样写!每个 Partition 都会打开文件,IO 爆炸 def check_ioc(log_line): with open("rules/ioc_rules.json") as f: rules = json.load(f) for rule in rules: if re.search(rule["pattern"], log_line): return rule["name"] return None

正确做法(Broadcast + 预编译正则):

# ✅ 在 Driver 端一次性读取并广播 with open("rules/ioc_rules.json") as f: raw_rules = json.load(f) # 预编译正则,避免每个 record 重复 compile compiled_rules = [ (rule["name"], re.compile(rule["pattern"]), rule["severity"]) for rule in raw_rules ] # 广播到所有 Executor bc_rules = sc.broadcast(compiled_rules) # 在 RDD map 中使用 def check_ioc_broadcast(log_line): rules = bc_rules.value # 获取广播变量 for name, pattern, severity in rules: if pattern.search(log_line): return (name, severity, log_line[:100]) # 返回规则名、等级、日志片段 return None # 应用到 DStream alerts = parsed_logs.map(check_ioc_broadcast).filter(lambda x: x is not None)

为什么必须 broadcast?

  • sc.broadcast()将变量序列化后发送到每个 Executor 的内存中,只传输一次;
  • 若用map内部open(),每个 Task(每个 Partition 的每个 record)都会触发磁盘 IO,10 万条日志 = 10 万次文件打开,必然 OOM;
  • 预编译re.compile()避免正则引擎重复解析,提升 3~5 倍匹配速度。

4.3 滑动窗口聚合:reduceByKeyAndWindow的窗口函数与触发时机

入侵检测的核心是“行为链”,单条日志无意义,需看时间窗口内的模式。系统使用reduceByKeyAndWindow实现:

# 每条日志解析为 (ip, 1) 键值对 ip_counts = parsed_logs.map(lambda log: (log["remote_addr"], 1)) # 滑动窗口:窗口长度 30 秒,每 10 秒滑动一次 # reduceFunc: (acc, new) => acc + new (累加) # invReduceFunc: (acc, old) => acc - old (减去离开窗口的老值,需启用 state) windowed_counts = ip_counts.reduceByKeyAndWindow( reduceFunc=lambda x, y: x + y, invReduceFunc=lambda x, y: x - y, # 启用增量计算,必须提供 windowDuration=30, # 窗口长度 30 秒 slideDuration=10, # 每 10 秒计算一次 numPartitions=4 # 避免单 partition 热点 ) # 过滤出高频 IP(>50 次/30s) suspicious_ips = windowed_counts.filter(lambda kv: kv[1] > 50)

关键点:

  • invReduceFunc是性能命脉:若不提供,Spark 每次窗口计算都重新 scan 全量数据,O(n²) 复杂度;提供后变为 O(n),仅更新增删部分;
  • numPartitions=4:默认 1 个 partition,所有 IP hash 到同一 partition,造成单核瓶颈;设为 4 后负载均衡;
  • 窗口时间单位是秒,batchDuration=10意味着每 10 秒生成一个 micro-batch,windowDuration=30即跨 3 个 batch。

5. Flask Web 服务与前端联动:REST API 设计、ECharts 动态渲染、以及部署时的flask run --host=0.0.0.0玄学

5.1 Flask API 接口设计:为什么/api/alerts/latest必须用@app.route而非 WebSocket

系统提供两个核心接口:

接口方法用途数据格式更新频率
/api/alerts/latestGET获取最近 10 条告警JSON Array每 5 秒 AJAX 轮询
/api/top_ipsGET获取当前 Top5 高频 IPJSON Object每 10 秒轮询

为什么不选 WebSocket?

  • WebSocket 需维护长连接,在 Flask 中需额外引入Flask-SocketIO+eventlet,增加部署复杂度;
  • 毕设演示场景下,5 秒轮询延迟完全可接受,且前端 EChartssetOption()支持平滑过渡动画;
  • @app.route更符合 RESTful 规范,便于 Postman 测试和答辩现场 curl 验证。

app.py中关键代码:

from flask import Flask, jsonify import redis # 使用 Redis 作中间缓存,解耦 Spark 与 Flask app = Flask(__name__) # Redis 连接池,存储最新告警(List)和 Top IP(Sorted Set) r = redis.Redis(host='localhost', port=6379, db=0, decode_responses=True) @app.route('/api/alerts/latest') def get_latest_alerts(): # LRANGE 获取最新 10 条,Redis List 天然支持 LIFO alerts = r.lrange('alerts', 0, 9) # 解析 JSON 字符串 return jsonify([json.loads(alert) for alert in alerts]) @app.route('/api/top_ips') def get_top_ips(): # ZREVRANGE 获取分数最高的 5 个 IP(分数 = 请求次数) ips = r.zrevrange('top_ips', 0, 4, withscores=True) return jsonify({ip: int(score) for ip, score in ips})

为什么用 Redis?

  • Spark Streaming 任务持续运行,不能直接暴露DStream.foreachRDD()到 Flask;
  • Redis 作为共享存储,Spark 侧用r.lpush('alerts', json.dumps(alert))和r.zincrby('top_ips', 1, ip)更新;
  • Flask 侧只读,无锁竞争,性能稳定。

5.2 前端 ECharts 渲染:动态拓扑图的节点/边数据构造逻辑

templates/index.html中,ECharts 初始化代码:

// 初始化拓扑图 const chart = echarts.init(document.getElementById('topology')); chart.setOption({ tooltip: {}, animation: false, series: [{ type: 'graph', layout: 'force', force: { repulsion: 1000 }, data: [], // 节点数组 links: [], // 边数组 categories: [ { name: 'Normal' }, { name: 'Suspicious' }, { name: 'Attacker' } ] }] }); // 每 5 秒拉取新数据并更新 function updateChart() { fetch('/api/alerts/latest') .then(r => r.json()) .then(alerts => { const nodes = []; const links = []; const attackerIPs = new Set(); // Step 1: 收集所有告警中的 attacker IP alerts.forEach(alert => { if (alert.ip && alert.severity === 'high') { attackerIPs.add(alert.ip); } }); // Step 2: 构造节点(attacker 为红色,normal 为绿色) attackerIPs.forEach(ip => { nodes.push({ name: ip, value: 10, category: 2, // Attacker symbolSize: 30, itemStyle: { color: '#e74c3c' } }); }); // Step 3: 构造边(attacker → target,target 来自 referer 或 uri) alerts.forEach(alert => { if (attackerIPs.has(alert.ip) && alert.target) { links.push({ source: alert.ip, target: alert.target, lineStyle: { width: 2, curveness: 0.2 } }); } }); chart.setOption({ series: [{ data: nodes, links: links }] }); }); } setInterval(updateChart, 5000); updateChart();

关键逻辑:

  • attackerIPs是 Set,自动去重,避免同一 IP 多次渲染;
  • links构造时,alert.target来自日志中的request_uri(如/admin/login.php)或referer(如http://evil.com/exploit.js),形成“攻击者 → 受害资源”关系;
  • symbolSize: 30和itemStyle.color强化视觉区分,答辩时老师一眼看出攻击链。

5.3 Flask 部署避坑:flask run --host=0.0.0.0为何是玄学,以及gunicorn的必要性

现象:本地开发时flask run正常,部署到 master 服务器后,浏览器访问http://192.168.56.101:5000显示空白,curl http://localhost:5000却返回 HTML。

原因:Flask 默认绑定127.0.0.1:5000,只接受本机回环请求。--host=0.0.0.0强制监听所有网卡,但存在两大风险:
①生产环境不安全:flask run是 Werkzeug 开发服务器,无并发能力,高并发下直接挂;
②端口冲突:若iptables或云平台安全组未放行 5000 端口,外部无法访问。

正确做法(生产部署):

# 安装 gunicorn pip install gunicorn # 启动(workers 数 = CPU 核数 * 2 + 1) gunicorn -w 4 -b 0.0.0.0:5000 --timeout 120 app:app

参数说明:

  • -w 4:启动 4 个 worker 进程,处理并发请求;
  • -b 0.0.0.0:5000:绑定所有地址的 5000 端口;
  • --timeout 120:请求超时 120 秒,避免长连接阻塞;
  • app:app:第一个app是文件名app.py,第二个app是 Flask 实例名。

注意:启动前确保app.py中if __name__ == '__main__':块已删除或注释,否则 gunicorn 会执行两次。


6. 真实应急响应验证:用ab模拟 CC 攻击、sqlmap触发 IOC、以及从日志到告警的端到端耗时测量技巧

6.1 模拟攻击链:三步复现“扫描 → 注入 → 告警”的完整闭环

Step 1:用 Apache Bench 模拟 CC 攻击(高频请求)
在 client 机器(非集群节点)执行:

# 持续 60 秒,每秒 100 个并发请求,访问 /test.php(不存在路径,触发 404) ab -n 6000 -c 100 "http://192.168.56.101/test.php"

此时 Flume 的tail -F会捕获大量 404 日志,Spark Streaming 的windowed_counts将在 30 秒窗口内统计该 IP 请求量,若 >50 则进入suspicious_ips。

Step 2:用 sqlmap 触发 SQL 注入规则

# 向存在漏洞的测试接口注入 sqlmap -u "http://192.168.56.101/vuln?id=1" --technique=U --batch # 生成日志如:GET /vuln?id=1%20UNION%20SELECT%20NULL%2Cpassword%20FROM%20users

Flume 采集该行后,Spark 的check_ioc_broadcast函数会匹配(?i)union\s+select,生成告警事件并写入 Redis。

Step 3:验证告警是否到达前端
打开浏览器http://192.168.56.101:5000,观察:

  • 右上角“最新告警”列表应在 15 秒内(3 个 batch)出现SQLi_UNION_SELECT条目;
  • 拓扑图中应出现红色节点(攻击者 IP)指向/vuln的边;
  • redis-cli中执行LRANGE alerts 0 0应看到完整 JSON 告警。

6.2 端到端耗时测量:从日志写入磁盘到前端显示的 5 个时间戳

要证明“实时性”,必须量化各环节延迟。我们在关键节点插入时间戳:

环节时间戳位置测量方法典型耗时
T1:日志落盘Nginxaccess.log文件末尾tail -1 /var/log/nginx/access.log | awk '{print $4}'0ms(写入即完成)
T2:Flume 采集Flume Agent 的Log4jLogger输出grep "SINK.k1" $FLUME_HOME/logs/flume.log | tail -1100~300ms(取决于 batch-size)
T3:Spark 接收FlumeUtils.createStream()的foreachRDD在foreachRDD内print("Received at:", time.time())200~500ms(网络+序列化)
T4:规则匹配完成alerts.foreachRDD()内print("Alerted at:", time.time())同上300~800ms(正则匹配+Redis写入)
T5:前端显示浏览器控制台console.log(new Date())在updateChart()开头添加5~10s(5 秒轮询间隔)

结论:从 T1 到 T4 的纯处理链路耗时约1~2 秒,完全满足

本文还有配套的精品资源,点击获取

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/10 4:53:32

Foxnic-EAM固定设备资产管理系统:让工业设备实时‘开口说话’

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/10 4:53:22

100份中文文本批量挖掘:清洗、向量化与可解释聚类实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/10 4:52:33

对话式AI记忆系统设计:从写入、检索到冲突处理的工程实践

1. 从"claude-mem"这个名字说起&#xff1a;它到底想解决什么问题第一次看到claude-mem这个命名&#xff0c;我的直觉是&#xff1a;这是一个围绕对话记忆做文章的项目。拆开来看&#xff0c;"claude" 指向的是对话式 AI 的交互场景&#xff0c;"mem&…

作者头像 李华
网站建设 2026/10/10 4:52:21

PCA9422+PIC18F57Q43硬件协同电源管理方案

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/10/10 4:52:01

作业焦虑自救指南:从“rip作业”到掌控任务的系统方法

1. "rip作业"背后的真实情绪&#xff1a;不是懒&#xff0c;是真的被压垮了先别急着把自己归到"不努力"那一类。我见过太多深夜对着电脑屏幕发呆、对着摊开的教材想把它合上再也不想打开的人——嘴里默念一句"rip作业"&#xff0c;然后继续熬夜、…

作者头像 李华