時(shí)日志入侵檢測(cè)系統(tǒng)實(shí)戰(zhàn):從環(huán)境搭建到規(guī)則優(yōu)化)
簡介這份資源是面向計(jì)算機(jī)、大數(shù)據(jù)、人工智能等專業(yè)學(xué)生與技術(shù)學(xué)習(xí)者的分布式實(shí)時(shí)日志分析與入侵檢測(cè)系統(tǒng)完整項(xiàng)目包基于Flume采集日志、Spark進(jìn)行流式處理、Flask搭建可視化與接口層適合用作課程設(shè)計(jì)、期末大作業(yè)或畢業(yè)設(shè)計(jì)的參考方案也可作為學(xué)習(xí)分布式日志管道與安全檢測(cè)思路的實(shí)戰(zhàn)素材。壓縮包共108個(gè)文件約18.88MB包含Scala與Java源碼、sbt構(gòu)建配置、properties與conf配置文件、HTML/CSS/JS前端頁面、Python腳本以及日志與數(shù)據(jù)樣本等覆蓋采集、計(jì)算、展示各環(huán)節(jié)目錄結(jié)構(gòu)便于按模塊查閱。目前已有234人學(xué)習(xí)下載。項(xiàng)目代碼經(jīng)過調(diào)試下載后可直接運(yùn)行讀者可據(jù)此理解Flume到Spark再到Flask的完整數(shù)據(jù)鏈路掌握日志解析、實(shí)時(shí)統(tǒng)計(jì)與入侵行為識(shí)別的實(shí)現(xiàn)方式并參考其中的配置與排錯(cuò)思路快速搭建自己的實(shí)驗(yàn)環(huán)境。1. 從一堆.cache文件說起這套 FlumeSparkFlask 日志入侵檢測(cè)系統(tǒng)到底能跑出什么如果你手頭正好有一個(gè)「基于 FlumeSparkFlask 的分布式實(shí)時(shí)日志分析與入侵檢測(cè)系統(tǒng)」的壓縮包解壓后第一眼看到的很可能不是熟悉的.py或.java而是一串像access_log、$3d85af9b26c1a259b49e.cache、$da50ce791668c9ed0f15$.class這樣的文件。別慌這不是打包出錯(cuò)而是 Spark 在本地或集群模式下運(yùn)行時(shí)留下的中間產(chǎn)物——.cache是 RDD 或 DataFrame 被persist()后落盤的塊文件$.class則是 Scala 編譯出的匿名類。能出現(xiàn)這些文件說明這套代碼至少被真實(shí)提交運(yùn)行過不是純靜態(tài)的「骨架工程」。這套資源解決的是一個(gè)很具體的問題把分散在多臺(tái)機(jī)器上的訪問日志通過 Flume 采集匯聚交給 Spark 做實(shí)時(shí)解析和規(guī)則匹配識(shí)別出暴力破解、異常高頻訪問、可疑路徑掃描等入侵特征最后用 Flask 提供一個(gè)能看圖表和告警的 Web 界面。它適合正在做課程設(shè)計(jì)、期末大作業(yè)或畢設(shè)的計(jì)算機(jī)、大數(shù)據(jù)、人工智能方向的學(xué)生也適合想跑通「采集→計(jì)算→展示」完整鏈路的技術(shù)學(xué)習(xí)者。前提是你得有一點(diǎn) Linux、Java 和 Python 基礎(chǔ)否則連 Flume 的配置文件都改不動(dòng)。2. 拆開壓縮包先看什么Flume、Spark、Flask 三層各自的入口與配置2.1 目錄結(jié)構(gòu)與三個(gè)核心入口拿到壓縮包后不要急著pip install先把目錄樹看清楚。這類項(xiàng)目通常按技術(shù)棧分層常見結(jié)構(gòu)是flume-conf/、spark-job/、flask-web/三個(gè)主目錄外加一個(gè)logs/放模擬日志、一個(gè)sql/放建表語句。你要找的第一個(gè)文件是 Flume 的.conf配置第二個(gè)是 Spark 的提交腳本或main函數(shù)第三個(gè)是 Flask 的app.py或run.py。先確認(rèn)三件事Flume 的 source 類型是exec還是taildirSpark 的入口是SparkSession還是老的SparkContextFlask 是直接讀 Spark 寫出的結(jié)果表還是通過 API 再查一次。這三個(gè)選擇決定了你后面要不要裝 Kafka、要不要配 Hive、要不要起 Redis。很多同學(xué)跑不起來不是代碼錯(cuò)而是沒意識(shí)到這套工程默認(rèn)依賴了外部存儲(chǔ)。# 先看目錄層級(jí)確認(rèn)三個(gè)入口文件的位置 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 作業(yè)的提交方式是 spark-submit 還是 python 直接跑 head -50 spark-job/*.py 2/dev/null || head -50 spark-job/*.scala 2/dev/null上面三條命令的作用分別是定位所有可能的入口文件、提取 Flume 的組件聲明、判斷 Spark 作業(yè)的語言和提交方式。參數(shù)上重點(diǎn)看a1.sources.r1.type如果是TAILDIR就支持?jǐn)帱c(diǎn)續(xù)傳如果是EXEC則每次重啟會(huì)從頭讀生產(chǎn)環(huán)境一般選前者。a1.sinks.k1.type如果是logger說明只是調(diào)試用真正落地通常改成hdfs或kafka。2.2 Flume 采集配置source、channel、sink 怎么改才不丟數(shù)據(jù)Flume 這一層最容易翻車的地方是 channel 容量和 batchSize 不匹配。默認(rèn)capacity1000、transactionCapacity100如果日志突發(fā)流量大source 寫入速度超過 sink 消費(fèi)速度channel 滿了就會(huì)拋ChannelException日志直接丟。常見做法是把capacity調(diào)到 10000 以上transactionCapacity調(diào)到 1000同時(shí)把 sink 的batchSize設(shè)成和transactionCapacity一致。# flume-conf/access-log.conf 關(guān)鍵參數(shù) 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按文件組監(jiān)控日志positionFile記錄讀取偏移量重啟后不會(huì)重復(fù)消費(fèi)。channel 容量給到 20000事務(wù)容量 2000sink 的 batchSize 也設(shè) 2000三者形成背壓緩沖。如果不想引入 Kafka把 sink 改成hdfs或logger也能跑但實(shí)時(shí)性會(huì)打折扣。注意positionFile的路徑要有寫權(quán)限否則 Flume 啟動(dòng)時(shí)會(huì)靜默失敗日志里只報(bào)一行Permission denied。2.3 Spark 實(shí)時(shí)解析從日志行到入侵特征的轉(zhuǎn)換邏輯Spark 這一層干的事是把原始日志行拆成字段然后按規(guī)則打標(biāo)簽。典型日志格式是 Nginx 或 Apache 的 combined 格式用正則提取 IP、時(shí)間、方法、路徑、狀態(tài)碼、UA。提取完之后做兩類判斷一類是閾值類比如同一 IP 在 60 秒內(nèi)請(qǐng)求超過 100 次標(biāo)記為brute_force另一類是模式類比如路徑里出現(xiàn)../或union select標(biāo)記為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 讀或從本地文件讀做離線驗(yàn)證 df spark.readStream.format(kafka) \ .option(kafka.bootstrap.servers, localhost:9092) \ .option(subscribe, access_log_topic) \ .load() # 提取字段正則按實(shí)際日志格式調(diào)整 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 秒窗口內(nèi)同 IP 請(qǐng)求計(jì)數(shù)超過閾值標(biāo)記 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()這段代碼的關(guān)鍵參數(shù)有三個(gè)spark.sql.shuffle.partitions控制聚合時(shí)的并行度本地跑設(shè) 4 就夠集群上按核數(shù)調(diào)window的60 seconds是滑動(dòng)窗口長度改小會(huì)更靈敏但誤報(bào)多checkpointLocation必須指定否則流式作業(yè)重啟后無法恢復(fù)狀態(tài)。正則部分是最容易出問題的地方不同日志格式字段順序不一樣建議先用head -5 access.log看一眼真實(shí)行再對(duì)著改正則。2.4 Flask 展示層把檢測(cè)結(jié)果變成能看的頁面Flask 這一層通常不直接連 Spark而是讀 Spark 寫出的結(jié)果表或 Redis 緩存。常見做法是 Spark 把告警寫入 MySQL 或 HiveFlask 用 SQLAlchemy 查出來渲染成表格和 ECharts 圖。如果你看到app.py里有pymysql或sqlalchemy的 import基本就是這個(gè)路子。# 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 賬號(hào)密碼改charsetutf8mb4不能省否則中文路徑會(huì)亂碼。/api/alerts是給前端 ECharts 異步拉數(shù)據(jù)用的返回 JSON 時(shí)force_asciiFalse保證中文可讀。如果 Flask 啟動(dòng)后頁面空白先看瀏覽器控制臺(tái)有沒有 500再看 MySQL 里alerts表是不是空的——Spark 沒寫進(jìn)去前端自然沒東西顯示。3. 從零跑通全鏈路環(huán)境準(zhǔn)備、啟動(dòng)順序與驗(yàn)證方法3.1 環(huán)境版本對(duì)齊JDK、Scala、Spark、Python 的兼容矩陣這套工程跑不起來十有八九是版本打架。Spark 3.x 默認(rèn)綁 Scala 2.12Spark 2.4 綁 Scala 2.11如果你下的包是 2.4 的卻裝了 2.12 的 Scala提交作業(yè)時(shí)會(huì)報(bào)NoSuchMethodError。Python 側(cè)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 對(duì)應(yīng)2.4 用 2.11Spark3.3.x穩(wěn)定且文檔多避免用 4.x 預(yù)覽版Python3.83.103.11 以上部分庫輪子不全Flume1.9 或 1.111.11 對(duì) TAILDIR 支持更好Flask2.x3.x 也可注意 Jinja2 語法差異對(duì)齊版本最省事的辦法是先spark-submit --version看輸出再python -c import pyspark; print(pyspark.__version__)兩個(gè)不一致就重裝。JDK 用java -version確認(rèn)如果是 17 而 Spark 是 2.4直接換 JDK 8別折騰參數(shù)。3.2 啟動(dòng)順序Flume → Kafka → Spark → Flask 的依賴鏈啟動(dòng)順序錯(cuò)了后面全白搭。正確鏈路是先起 Kafka如果 sink 用 Kafka再起 Flume 采集然后提交 Spark 流式作業(yè)最后起 Flask。因?yàn)?Spark 要訂閱 Kafka topictopic 不存在會(huì)直接報(bào)錯(cuò)退出Flask 要查 MySQL表沒建也會(huì) 500。# 1. 起 Kafka單機(jī)快速驗(yàn)證 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 流式作業(yè) 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每一步都有驗(yàn)證點(diǎn)Kafka 起完后jps應(yīng)該看到Kafka和QuorumPeerMainFlume 起完后往access.log追加一行控制臺(tái)應(yīng)打印出該行Spark 提交后控制臺(tái)應(yīng)出現(xiàn)Batch: 0之類的進(jìn)度Flask 起完后瀏覽器訪問http://127.0.0.1:5000能看到頁面。哪一步?jīng)]輸出就停在那一步排查不要往下走。3.3 用模擬日志驗(yàn)證入侵檢測(cè)規(guī)則是否生效工程里一般帶一個(gè)logs/access.log或生成腳本。如果沒有自己造幾條能觸發(fā)規(guī)則的日志。比如同一 IP 連續(xù) 150 次請(qǐng)求或者路徑里帶../etc/passwd。追加日志用echo循環(huán)觀察 Spark 控制臺(tái)是否打出告警。# 模擬暴力破解同一 IP 快速請(qǐng)求 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追加后等一個(gè)窗口周期默認(rèn) 60 秒Spark 控制臺(tái)應(yīng)出現(xiàn)count 100的記錄Flask 頁面刷新后表格里應(yīng)出現(xiàn)192.168.1.100。如果沒出現(xiàn)先確認(rèn) Flume 是否真的讀到了新行看 Flume 日志再確認(rèn) Spark 的window時(shí)間字段解析是否正確——ts字段如果沒轉(zhuǎn)成 timestampwindow函數(shù)會(huì)直接報(bào)錯(cuò)或返回空。4. 避坑與排查這套工程最容易翻車的五個(gè)地方4.1 現(xiàn)象Flume 啟動(dòng)后日志不采集控制臺(tái)無輸出原因通常是TAILDIR的filegroups路徑寫錯(cuò)或者positionFile所在目錄沒有寫權(quán)限。Flume 對(duì)路徑錯(cuò)誤不敏感不會(huì)報(bào)致命錯(cuò)誤只是靜默不讀。解決方法是先用ls -l確認(rèn)日志文件存在且可讀再把positionFile指到/tmp下最后把 Flume 日志級(jí)別調(diào)到DEBUG看TaildirSource有沒有掃描到文件。4.2 現(xiàn)象Spark 提交報(bào)ClassNotFoundException: kafka.serializer.StringDecoder原因是--packages里的 Kafka 連接器版本和 Spark 版本不匹配。Spark 3.3 要用spark-sql-kafka-0-10_2.12:3.3.0如果寫成2.4.0就會(huì)找不到類。解決方法是先spark-submit --version確認(rèn) Spark 版本再把--packages的版本號(hào)改成一致。如果公司內(nèi)網(wǎng)拉不到包提前把 jar 下好放到$SPARK_HOME/jars下。4.3 現(xiàn)象Flask 頁面能打開但表格為空MySQL 里也沒數(shù)據(jù)原因是 Spark 流式作業(yè)沒有把結(jié)果寫入 MySQL或者寫入了但表名不對(duì)。常見做法是 Spark 用foreachBatch寫 JDBC如果foreachBatch里沒調(diào)df.write.jdbc數(shù)據(jù)就只打在控制臺(tái)。解決方法是檢查 Spark 代碼里有沒有writeStream.foreachBatch或write.jdbc并確認(rèn) MySQL 的alerts表已建好字段和 DataFrame 的 schema 對(duì)得上。4.4 現(xiàn)象日志時(shí)間字段解析失敗window函數(shù)報(bào)AnalysisException原因是正則提取出的ts是字符串直接cast(timestamp)時(shí)格式不匹配。Nginx 默認(rèn)格式是10/Oct/2024:10:00:00 0800Spark 的to_timestamp默認(rèn)不認(rèn)這個(gè)格式。解決方法是顯式指定格式to_timestamp(col(ts), dd/MMM/yyyy:HH:mm:ss Z)注意MMM是英文月份縮寫本地化環(huán)境要設(shè)spark.sql.legacy.timeParserPolicyLEGACY。4.5 現(xiàn)象本地跑得好好的換臺(tái)機(jī)器就報(bào)No such file or directory: /tmp/spark_checkpoint原因是 checkpoint 路徑寫死在代碼里換機(jī)器后目錄不存在。Spark 流式作業(yè)的 checkpoint 目錄必須提前創(chuàng)建且要有寫權(quán)限。解決方法是在代碼里加os.makedirs(/tmp/spark_checkpoint, exist_okTrue)或者把路徑改成從環(huán)境變量讀部署時(shí)統(tǒng)一配。另外 checkpoint 目錄不要放在/tmp下長期跑系統(tǒng)清理會(huì)把它刪掉導(dǎo)致作業(yè)恢復(fù)失敗。5. 進(jìn)階技巧把檢測(cè)規(guī)則從硬編碼改成可配置并用歷史日志回放驗(yàn)證5.1 規(guī)則外置用 JSON 配置替代寫死的閾值原始工程里閾值大概率是寫死在 Python 里的比如count 100。這樣改一次規(guī)則就要改代碼、重提交很麻煩。我一般會(huì)把規(guī)則抽成 JSONSpark 啟動(dòng)時(shí)讀一次廣播到各 executor。這樣調(diào)閾值不用動(dòng)代碼改完重啟作業(yè)即可。# rules.json { brute_force: {window_seconds: 60, threshold: 100, field: ip}, path_scan: {patterns: [../, union select, etc/passwd], field: path} }# 讀取規(guī)則并廣播 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 取規(guī)則 threshold bc_rules.value[brute_force][threshold]廣播變量的好處是每個(gè) executor 只存一份不會(huì)因?yàn)橐?guī)則變大而拖慢序列化。參數(shù)上注意window_seconds和threshold要聯(lián)動(dòng)調(diào)窗口越長閾值應(yīng)越高否則誤報(bào)會(huì)淹沒真實(shí)告警。patterns列表里的字符串會(huì)被拼成正則特殊字符要轉(zhuǎn)義比如../里的.要寫成\.。5.2 歷史日志回放用離線模式驗(yàn)證規(guī)則準(zhǔn)確率流式作業(yè)調(diào)試起來慢改一次等一個(gè)窗口。更高效的做法是先用離線模式跑歷史日志把規(guī)則調(diào)準(zhǔn)了再上流式。Spark 讀本地文件生成 DataFrame套用同樣的解析和判斷邏輯輸出告警數(shù)量和樣例人工看一眼誤報(bào)率。# 離線回放驗(yàn)證 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 請(qǐng)求量最高 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)離線跑的好處是秒出結(jié)果不用等窗口。show(10, truncateFalse)不截?cái)嘧侄畏奖憧赐暾窂?。如果發(fā)現(xiàn)某個(gè)正常 IP 被誤判就把閾值調(diào)高或把該 IP 加白名單。白名單同樣可以放進(jìn)rules.json在過濾時(shí)filter(~col(ip).isin(whitelist))。5.3 一個(gè)我踩過的坑checkpoint 和規(guī)則變更的沖突有次我改了rules.json里的閾值重啟 Spark 作業(yè)后告警數(shù)量沒變。排查半天才發(fā)現(xiàn)流式作業(yè)的 checkpoint 里存了舊的查詢計(jì)劃規(guī)則雖然重新讀了但foreachBatch里用的還是廣播前的舊值。從那以后我每次改規(guī)則要么換一個(gè)新的checkpointLocation要么在代碼里加版本號(hào)規(guī)則版本變了就自動(dòng)切目錄。這個(gè)習(xí)慣幫我省了很多「改了沒生效」的玄學(xué)時(shí)間。import hashlib rule_hash hashlib.md5(json.dumps(rules, sort_keysTrue).encode()).hexdigest()[:8] checkpoint_path f/tmp/spark_checkpoint_{rule_hash}這樣規(guī)則一變checkpoint 目錄跟著變Spark 會(huì)當(dāng)成新作業(yè)啟動(dòng)不會(huì)復(fù)用舊狀態(tài)。代價(jià)是歷史狀態(tài)丟失但對(duì)入侵檢測(cè)這種場(chǎng)景重新開始統(tǒng)計(jì)反而更干凈。希望幫到你。本文還有配套的精品資源點(diǎn)擊獲取