简介:这份资源是面向计算机、大数据、人工智能等专业学生与技术学习者的分布式实时日志分析与入侵检测系统完整项目包,基于Flume采集日志、Spark进行流式处理、Flask搭建可视化与接口层,适合用作课程设计、期末大作业或毕业设计的参考方案,也可作为学习分布式日志管道与安全检测思路的实战素材。压缩包共108个文件,约18.88MB,包含Scala与Java源码、sbt构建配置、properties与conf配置文件、HTML/CSS/JS前端页面、Python脚本以及日志与数据样本等,覆盖采集、计算、展示各环节,目录结构便于按模块查阅。目前已有234人学习下载。项目代码经过调试,下载后可直接运行,读者可据此理解Flume到Spark再到Flask的完整数据链路,掌握日志解析、实时统计与入侵行为识别的实现方式,并参考其中的配置与排错思路快速搭建自己的实验环境。
1. 从一堆.cache文件说起:这套 Flume+Spark+Flask 日志入侵检测系统到底能跑出什么
如果你手头正好有一个「基于 Flume+Spark+Flask 的分布式实时日志分析与入侵检测系统」的压缩包,解压后第一眼看到的很可能不是熟悉的.py或.java,而是一串像access_log、$3d85af9b26c1a259b49e.cache、$da50ce791668c9ed0f15$.class这样的文件。别慌,这不是打包出错,而是 Spark 在本地或集群模式下运行时留下的中间产物——.cache是 RDD 或 DataFrame 被persist()后落盘的块文件,$.class则是 Scala 编译出的匿名类。能出现这些文件,说明这套代码至少被真实提交运行过,不是纯静态的「骨架工程」。
这套资源解决的是一个很具体的问题:把分散在多台机器上的访问日志,通过 Flume 采集汇聚,交给 Spark 做实时解析和规则匹配,识别出暴力破解、异常高频访问、可疑路径扫描等入侵特征,最后用 Flask 提供一个能看图表和告警的 Web 界面。它适合正在做课程设计、期末大作业或毕设的计算机、大数据、人工智能方向的学生,也适合想跑通「采集→计算→展示」完整链路的技术学习者。前提是你得有一点 Linux、Java 和 Python 基础,否则连 Flume 的配置文件都改不动。
2. 拆开压缩包先看什么:Flume、Spark、Flask 三层各自的入口与配置
2.1 目录结构与三个核心入口
拿到压缩包后不要急着pip install,先把目录树看清楚。这类项目通常按技术栈分层,常见结构是flume-conf/、spark-job/、flask-web/三个主目录,外加一个logs/放模拟日志、一个sql/放建表语句。你要找的第一个文件是 Flume 的.conf配置,第二个是 Spark 的提交脚本或main函数,第三个是 Flask 的app.py或run.py。
先确认三件事:Flume 的 source 类型是exec还是taildir,Spark 的入口是SparkSession还是老的SparkContext,Flask 是直接读 Spark 写出的结果表还是通过 API 再查一次。这三个选择决定了你后面要不要装 Kafka、要不要配 Hive、要不要起 Redis。很多同学跑不起来,不是代码错,而是没意识到这套工程默认依赖了外部存储。
# 先看目录层级,确认三个入口文件的位置 find . -maxdepth 3 -type f \( -name "*.conf" -o -name "*.py" -o -name "*.scala" -o -name "*.sql" \) | sort # 看 Flume 配置里 source、channel、sink 分别是什么 grep -E "a1\.(sources|channels|sinks)" flume-conf/*.conf # 看 Spark 作业的提交方式,是 spark-submit 还是 python 直接跑 head -50 spark-job/*.py 2>/dev/null || head -50 spark-job/*.scala 2>/dev/null上面三条命令的作用分别是:定位所有可能的入口文件、提取 Flume 的组件声明、判断 Spark 作业的语言和提交方式。参数上重点看a1.sources.r1.type,如果是TAILDIR就支持断点续传,如果是EXEC则每次重启会从头读,生产环境一般选前者。a1.sinks.k1.type如果是logger说明只是调试用,真正落地通常改成hdfs或kafka。
2.2 Flume 采集配置:source、channel、sink 怎么改才不丢数据
Flume 这一层最容易翻车的地方是 channel 容量和 batchSize 不匹配。默认capacity=1000、transactionCapacity=100,如果日志突发流量大,source 写入速度超过 sink 消费速度,channel 满了就会抛ChannelException,日志直接丢。常见做法是把capacity调到 10000 以上,transactionCapacity调到 1000,同时把 sink 的batchSize设成和transactionCapacity一致。
# flume-conf/access-log.conf 关键参数 a1.sources.r1.type = TAILDIR a1.sources.r1.positionFile = /tmp/flume_taildir_position.json a1.sources.r1.filegroups.f1 = /home/logs/access.log.* a1.sources.r1.batchSize = 1000 a1.channels.c1.type = memory a1.channels.c1.capacity = 20000 a1.channels.c1.transactionCapacity = 2000 a1.sinks.k1.type = org.apache.flume.sink.kafka.KafkaSink a1.sinks.k1.kafka.topic = access_log_topic a1.sinks.k1.kafka.bootstrap.servers = localhost:9092 a1.sinks.k1.batchSize = 2000这段配置的逻辑是:TAILDIR按文件组监控日志,positionFile记录读取偏移量,重启后不会重复消费。channel 容量给到 20000,事务容量 2000,sink 的 batchSize 也设 2000,三者形成背压缓冲。如果不想引入 Kafka,把 sink 改成hdfs或logger也能跑,但实时性会打折扣。注意positionFile的路径要有写权限,否则 Flume 启动时会静默失败,日志里只报一行Permission denied。
2.3 Spark 实时解析:从日志行到入侵特征的转换逻辑
Spark 这一层干的事是把原始日志行拆成字段,然后按规则打标签。典型日志格式是 Nginx 或 Apache 的 combined 格式,用正则提取 IP、时间、方法、路径、状态码、UA。提取完之后做两类判断:一类是阈值类,比如同一 IP 在 60 秒内请求超过 100 次标记为brute_force;另一类是模式类,比如路径里出现../或union select标记为path_scan。
# spark-job/log_analyzer.py 核心片段 from pyspark.sql import SparkSession from pyspark.sql.functions import regexp_extract, window, count, col spark = SparkSession.builder \ .appName("LogIntrusionDetect") \ .config("spark.sql.shuffle.partitions", "4") \ .getOrCreate() # 从 Kafka 读,或从本地文件读做离线验证 df = spark.readStream.format("kafka") \ .option("kafka.bootstrap.servers", "localhost:9092") \ .option("subscribe", "access_log_topic") \ .load() # 提取字段,正则按实际日志格式调整 parsed = df.select( regexp_extract(col("value").cast("string"), r'^(\S+)', 1).alias("ip"), regexp_extract(col("value").cast("string"), r'\[(.*?)\]', 1).alias("ts"), regexp_extract(col("value").cast("string"), r'"(GET|POST) (\S+)', 2).alias("path"), regexp_extract(col("value").cast("string"), r'" (\d{3}) ', 1).alias("status") ) # 60 秒窗口内同 IP 请求计数,超过阈值标记 windowed = parsed.groupBy( window(col("ts").cast("timestamp"), "60 seconds"), col("ip") ).count().filter(col("count") > 100) query = windowed.writeStream \ .outputMode("update") \ .format("console") \ .option("checkpointLocation", "/tmp/spark_checkpoint") \ .start() query.awaitTermination()这段代码的关键参数有三个:spark.sql.shuffle.partitions控制聚合时的并行度,本地跑设 4 就够,集群上按核数调;window的60 seconds是滑动窗口长度,改小会更灵敏但误报多;checkpointLocation必须指定,否则流式作业重启后无法恢复状态。正则部分是最容易出问题的地方,不同日志格式字段顺序不一样,建议先用head -5 access.log看一眼真实行,再对着改正则。
2.4 Flask 展示层:把检测结果变成能看的页面
Flask 这一层通常不直接连 Spark,而是读 Spark 写出的结果表或 Redis 缓存。常见做法是 Spark 把告警写入 MySQL 或 Hive,Flask 用 SQLAlchemy 查出来渲染成表格和 ECharts 图。如果你看到app.py里有pymysql或sqlalchemy的 import,基本就是这个路子。
# flask-web/app.py 核心片段 from flask import Flask, render_template from sqlalchemy import create_engine import pandas as pd app = Flask(__name__) engine = create_engine("mysql+pymysql://root:password@localhost:3306/logdb?charset=utf8mb4") @app.route("/") def index(): df = pd.read_sql("SELECT ip, alert_type, COUNT(*) AS cnt FROM alerts GROUP BY ip, alert_type ORDER BY cnt DESC LIMIT 50", engine) return render_template("index.html", rows=df.to_dict("records")) @app.route("/api/alerts") def api_alerts(): df = pd.read_sql("SELECT * FROM alerts ORDER BY ts DESC LIMIT 200", engine) return df.to_json(orient="records", force_ascii=False)这里create_engine的连接串要按你本地的 MySQL 账号密码改,charset=utf8mb4不能省,否则中文路径会乱码。/api/alerts是给前端 ECharts 异步拉数据用的,返回 JSON 时force_ascii=False保证中文可读。如果 Flask 启动后页面空白,先看浏览器控制台有没有 500,再看 MySQL 里alerts表是不是空的——Spark 没写进去,前端自然没东西显示。
3. 从零跑通全链路:环境准备、启动顺序与验证方法
3.1 环境版本对齐:JDK、Scala、Spark、Python 的兼容矩阵
这套工程跑不起来,十有八九是版本打架。Spark 3.x 默认绑 Scala 2.12,Spark 2.4 绑 Scala 2.11,如果你下的包是 2.4 的却装了 2.12 的 Scala,提交作业时会报NoSuchMethodError。Python 侧,PySpark 的版本必须和 Spark 本体一致,pip install pyspark==3.3.0就要配 Spark 3.3.0 的安装包。
| 组件 | 推荐版本 | 说明 |
|---|---|---|
| JDK | 1.8 或 11 | Spark 3.x 建议 11,Spark 2.4 只能 1.8 |
| Scala | 2.12.x | 与 Spark 3.x 对应,2.4 用 2.11 |
| Spark | 3.3.x | 稳定且文档多,避免用 4.x 预览版 |
| Python | 3.8~3.10 | 3.11 以上部分库轮子不全 |
| Flume | 1.9 或 1.11 | 1.11 对 TAILDIR 支持更好 |
| Flask | 2.x | 3.x 也可,注意 Jinja2 语法差异 |
对齐版本最省事的办法是先spark-submit --version看输出,再python -c "import pyspark; print(pyspark.__version__)",两个不一致就重装。JDK 用java -version确认,如果是 17 而 Spark 是 2.4,直接换 JDK 8,别折腾参数。
3.2 启动顺序:Flume → Kafka → Spark → Flask 的依赖链
启动顺序错了,后面全白搭。正确链路是:先起 Kafka(如果 sink 用 Kafka),再起 Flume 采集,然后提交 Spark 流式作业,最后起 Flask。因为 Spark 要订阅 Kafka topic,topic 不存在会直接报错退出;Flask 要查 MySQL,表没建也会 500。
# 1. 起 Kafka(单机快速验证) bin/zookeeper-server-start.sh -daemon config/zookeeper.properties bin/kafka-server-start.sh -daemon config/server.properties bin/kafka-topics.sh --create --topic access_log_topic --bootstrap-server localhost:9092 --partitions 1 --replication-factor 1 # 2. 起 Flume bin/flume-ng agent --conf conf --conf-file flume-conf/access-log.conf --name a1 -Dflume.root.logger=INFO,console # 3. 提交 Spark 流式作业 spark-submit --master local[2] --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0 spark-job/log_analyzer.py # 4. 起 Flask cd flask-web && python app.py每一步都有验证点:Kafka 起完后jps应该看到Kafka和QuorumPeerMain;Flume 起完后往access.log追加一行,控制台应打印出该行;Spark 提交后控制台应出现Batch: 0之类的进度;Flask 起完后浏览器访问http://127.0.0.1:5000能看到页面。哪一步没输出,就停在那一步排查,不要往下走。
3.3 用模拟日志验证入侵检测规则是否生效
工程里一般带一个logs/access.log或生成脚本。如果没有,自己造几条能触发规则的日志。比如同一 IP 连续 150 次请求,或者路径里带../etc/passwd。追加日志用echo循环,观察 Spark 控制台是否打出告警。
# 模拟暴力破解:同一 IP 快速请求 150 次 for i in $(seq 1 150); do echo '192.168.1.100 - - [10/Oct/2024:10:00:00 +0800] "GET /login HTTP/1.1" 401 0 "-" "curl/7.68"' >> logs/access.log done # 模拟路径扫描 echo '192.168.1.101 - - [10/Oct/2024:10:01:00 +0800] "GET /../../etc/passwd HTTP/1.1" 404 0 "-" "nikto"' >> logs/access.log追加后等一个窗口周期(默认 60 秒),Spark 控制台应出现count > 100的记录,Flask 页面刷新后表格里应出现192.168.1.100。如果没出现,先确认 Flume 是否真的读到了新行(看 Flume 日志),再确认 Spark 的window时间字段解析是否正确——ts字段如果没转成 timestamp,window函数会直接报错或返回空。
4. 避坑与排查:这套工程最容易翻车的五个地方
4.1 现象:Flume 启动后日志不采集,控制台无输出
原因通常是TAILDIR的filegroups路径写错,或者positionFile所在目录没有写权限。Flume 对路径错误不敏感,不会报致命错误,只是静默不读。解决方法是先用ls -l确认日志文件存在且可读,再把positionFile指到/tmp下,最后把 Flume 日志级别调到DEBUG看TaildirSource有没有扫描到文件。
4.2 现象:Spark 提交报ClassNotFoundException: kafka.serializer.StringDecoder
原因是--packages里的 Kafka 连接器版本和 Spark 版本不匹配。Spark 3.3 要用spark-sql-kafka-0-10_2.12:3.3.0,如果写成2.4.0就会找不到类。解决方法是先spark-submit --version确认 Spark 版本,再把--packages的版本号改成一致。如果公司内网拉不到包,提前把 jar 下好放到$SPARK_HOME/jars下。
4.3 现象:Flask 页面能打开但表格为空,MySQL 里也没数据
原因是 Spark 流式作业没有把结果写入 MySQL,或者写入了但表名不对。常见做法是 Spark 用foreachBatch写 JDBC,如果foreachBatch里没调df.write.jdbc,数据就只打在控制台。解决方法是检查 Spark 代码里有没有writeStream.foreachBatch或write.jdbc,并确认 MySQL 的alerts表已建好,字段和 DataFrame 的 schema 对得上。
4.4 现象:日志时间字段解析失败,window函数报AnalysisException
原因是正则提取出的ts是字符串,直接cast("timestamp")时格式不匹配。Nginx 默认格式是10/Oct/2024:10:00:00 +0800,Spark 的to_timestamp默认不认这个格式。解决方法是显式指定格式:to_timestamp(col("ts"), "dd/MMM/yyyy:HH:mm:ss Z"),注意MMM是英文月份缩写,本地化环境要设spark.sql.legacy.timeParserPolicy=LEGACY。
4.5 现象:本地跑得好好的,换台机器就报No such file or directory: /tmp/spark_checkpoint
原因是 checkpoint 路径写死在代码里,换机器后目录不存在。Spark 流式作业的 checkpoint 目录必须提前创建,且要有写权限。解决方法是在代码里加os.makedirs("/tmp/spark_checkpoint", exist_ok=True),或者把路径改成从环境变量读,部署时统一配。另外 checkpoint 目录不要放在/tmp下长期跑,系统清理会把它删掉,导致作业恢复失败。
5. 进阶技巧:把检测规则从硬编码改成可配置,并用历史日志回放验证
5.1 规则外置:用 JSON 配置替代写死的阈值
原始工程里阈值大概率是写死在 Python 里的,比如count > 100。这样改一次规则就要改代码、重提交,很麻烦。我一般会把规则抽成 JSON,Spark 启动时读一次,广播到各 executor。这样调阈值不用动代码,改完重启作业即可。
# rules.json { "brute_force": {"window_seconds": 60, "threshold": 100, "field": "ip"}, "path_scan": {"patterns": ["../", "union select", "etc/passwd"], "field": "path"} }# 读取规则并广播 import json from pyspark.sql import SparkSession spark = SparkSession.builder.appName("LogIntrusionDetect").getOrCreate() with open("rules.json", "r", encoding="utf-8") as f: rules = json.load(f) bc_rules = spark.sparkContext.broadcast(rules) # 在 foreachBatch 或 map 里用 bc_rules.value 取规则 threshold = bc_rules.value["brute_force"]["threshold"]广播变量的好处是每个 executor 只存一份,不会因为规则变大而拖慢序列化。参数上注意window_seconds和threshold要联动调,窗口越长阈值应越高,否则误报会淹没真实告警。patterns列表里的字符串会被拼成正则,特殊字符要转义,比如../里的.要写成\.。
5.2 历史日志回放:用离线模式验证规则准确率
流式作业调试起来慢,改一次等一个窗口。更高效的做法是先用离线模式跑历史日志,把规则调准了再上流式。Spark 读本地文件生成 DataFrame,套用同样的解析和判断逻辑,输出告警数量和样例,人工看一眼误报率。
# 离线回放验证 df = spark.read.text("logs/access.log") parsed = df.select( regexp_extract(col("value"), r'^(\S+)', 1).alias("ip"), regexp_extract(col("value"), r'\[(.*?)\]', 1).alias("ts"), regexp_extract(col("value"), r'"(GET|POST) (\S+)', 2).alias("path") ) # 按 IP 聚合,看哪些 IP 请求量最高 parsed.groupBy("ip").count().orderBy(col("count").desc()).show(10, truncate=False) # 按路径匹配可疑模式 suspicious = parsed.filter(col("path").rlike("(\\.\\./|union select|etc/passwd)")) suspicious.show(20, truncate=False)离线跑的好处是秒出结果,不用等窗口。show(10, truncate=False)不截断字段,方便看完整路径。如果发现某个正常 IP 被误判,就把阈值调高或把该 IP 加白名单。白名单同样可以放进rules.json,在过滤时filter(~col("ip").isin(whitelist))。
5.3 一个我踩过的坑:checkpoint 和规则变更的冲突
有次我改了rules.json里的阈值,重启 Spark 作业后告警数量没变。排查半天才发现,流式作业的 checkpoint 里存了旧的查询计划,规则虽然重新读了,但foreachBatch里用的还是广播前的旧值。从那以后我每次改规则,要么换一个新的checkpointLocation,要么在代码里加版本号,规则版本变了就自动切目录。这个习惯帮我省了很多「改了没生效」的玄学时间。
import hashlib rule_hash = hashlib.md5(json.dumps(rules, sort_keys=True).encode()).hexdigest()[:8] checkpoint_path = f"/tmp/spark_checkpoint_{rule_hash}"这样规则一变,checkpoint 目录跟着变,Spark 会当成新作业启动,不会复用旧状态。代价是历史状态丢失,但对入侵检测这种场景,重新开始统计反而更干净。希望帮到你。
本文还有配套的精品资源,点击获取