简介这份资源是面向计算机、大数据、人工智能等专业学生与技术学习者的分布式实时日志分析与入侵检测系统完整项目包基于Flume采集日志、Spark进行流式处理、Flask搭建可视化与接口层适合用作课程设计、期末大作业或毕业设计的参考方案也可作为学习分布式日志管道与安全检测思路的实战素材。压缩包共108个文件约18.88MB包含Scala与Java源码、sbt构建配置、properties与conf配置文件、HTML/CSS/JS前端页面、Python脚本以及日志与数据样本等覆盖采集、计算、展示各环节目录结构便于按模块查阅。目前已有234人学习下载。项目代码经过调试下载后可直接运行读者可据此理解Flume到Spark再到Flask的完整数据链路掌握日志解析、实时统计与入侵行为识别的实现方式并参考其中的配置与排错思路快速搭建自己的实验环境。1. 从一堆.cache文件说起这套 FlumeSparkFlask 日志入侵检测系统到底能跑出什么如果你手头正好有一个「基于 FlumeSparkFlask 的分布式实时日志分析与入侵检测系统」的压缩包解压后第一眼看到的很可能不是熟悉的.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还是taildirSpark 的入口是SparkSession还是老的SparkContextFlask 是直接读 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 不匹配。默认capacity1000、transactionCapacity100如果日志突发流量大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事务容量 2000sink 的 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 或 HiveFlask 用 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(mysqlpymysql://root:passwordlocalhost:3306/logdb?charsetutf8mb4) 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, rowsdf.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(orientrecords, force_asciiFalse)这里create_engine的连接串要按你本地的 MySQL 账号密码改charsetutf8mb4不能省否则中文路径会乱码。/api/alerts是给前端 ECharts 异步拉数据用的返回 JSON 时force_asciiFalse保证中文可读。如果 Flask 启动后页面空白先看浏览器控制台有没有 500再看 MySQL 里alerts表是不是空的——Spark 没写进去前端自然没东西显示。3. 从零跑通全链路环境准备、启动顺序与验证方法3.1 环境版本对齐JDK、Scala、Spark、Python 的兼容矩阵这套工程跑不起来十有八九是版本打架。Spark 3.x 默认绑 Scala 2.12Spark 2.4 绑 Scala 2.11如果你下的包是 2.4 的却装了 2.12 的 Scala提交作业时会报NoSuchMethodError。Python 侧PySpark 的版本必须和 Spark 本体一致pip install pyspark3.3.0就要配 Spark 3.3.0 的安装包。组件推荐版本说明JDK1.8 或 11Spark 3.x 建议 11Spark 2.4 只能 1.8Scala2.12.x与 Spark 3.x 对应2.4 用 2.11Spark3.3.x稳定且文档多避免用 4.x 预览版Python3.83.103.11 以上部分库轮子不全Flume1.9 或 1.111.11 对 TAILDIR 支持更好Flask2.x3.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 topictopic 不存在会直接报错退出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.loggerINFO,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和QuorumPeerMainFlume 起完后往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字段如果没转成 timestampwindow函数会直接报错或返回空。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 0800Spark 的to_timestamp默认不认这个格式。解决方法是显式指定格式to_timestamp(col(ts), dd/MMM/yyyy:HH:mm:ss Z)注意MMM是英文月份缩写本地化环境要设spark.sql.legacy.timeParserPolicyLEGACY。4.5 现象本地跑得好好的换台机器就报No such file or directory: /tmp/spark_checkpoint原因是 checkpoint 路径写死在代码里换机器后目录不存在。Spark 流式作业的 checkpoint 目录必须提前创建且要有写权限。解决方法是在代码里加os.makedirs(/tmp/spark_checkpoint, exist_okTrue)或者把路径改成从环境变量读部署时统一配。另外 checkpoint 目录不要放在/tmp下长期跑系统清理会把它删掉导致作业恢复失败。5. 进阶技巧把检测规则从硬编码改成可配置并用历史日志回放验证5.1 规则外置用 JSON 配置替代写死的阈值原始工程里阈值大概率是写死在 Python 里的比如count 100。这样改一次规则就要改代码、重提交很麻烦。我一般会把规则抽成 JSONSpark 启动时读一次广播到各 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, encodingutf-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, truncateFalse) # 按路径匹配可疑模式 suspicious parsed.filter(col(path).rlike((\\.\\./|union select|etc/passwd))) suspicious.show(20, truncateFalse)离线跑的好处是秒出结果不用等窗口。show(10, truncateFalse)不截断字段方便看完整路径。如果发现某个正常 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_keysTrue).encode()).hexdigest()[:8] checkpoint_path f/tmp/spark_checkpoint_{rule_hash}这样规则一变checkpoint 目录跟着变Spark 会当成新作业启动不会复用旧状态。代价是历史状态丢失但对入侵检测这种场景重新开始统计反而更干净。希望帮到你。本文还有配套的精品资源点击获取
阅读完成 · 觉得有帮助?