hkocamandev/iot-intrusion-detection
GitHub: hkocamandev/iot-intrusion-detection
基于 Kafka 流处理的 IoT 网络入侵检测系统,结合规则引擎与 XGBoost 机器学习双引擎,并集成 LLM 实现告警智能增强与自然语言交互。
Stars: 0 | Forks: 0
(group: storage)"] --> PG[("PostgreSQL
network_flows")] T1 --> RC["RuleBasedConsumer
(group: rule-engine)"] T1 --> MC["MlBasedConsumer
(group: ml-engine)"] MC -->|"REST /predict"| ML["ML service
FastAPI + XGBoost"] RC -->|RULE| T2[("iot.alerts")] MC -->|ML| T2 T2 --> AC["AlertConsumer
(group: alert-processor)"] --> AL[("PostgreSQL
alerts")] T2 --> WS["AlertWebSocketConsumer
(group: dashboard)"] --> UI["live dashboard"] AC -.->|"retry + DLQ"| DLT[("iot.alerts-dlt")] AL --> EN["AlertEnrichmentScheduler"] -->|LLM| LLMP["Groq / Gemini / Anthropic"] --> AL ``` 每个检测器都有独立的 Kafka consumer group,因此存储、基于规则和基于 ML 的检测 会各自独立地接收每条流量消息;告警投递使用了重试 topic 和死信 topic。  ### ML 检测服务 一个 FastAPI + XGBoost 微服务对每个流进行分类。  ``` curl -s -X POST localhost:8000/predict -H "Content-Type: application/json" \ -d '{"features": {"fwd_pkts_tot": 1500, "flow_duration": 2.3}}' # {"attack_type": "PORT_SCAN", "score": 0.826} ``` ### LLM 增强与自然语言聊天(免费 — 默认在 Groq 上运行) 定时轮询器会使用通俗易懂的解释和建议来增强每条告警: ``` { "attackType": "DDOS", "severity": "MEDIUM", "detectionSource": "ML", "llmExplanation": "A medium-severity DDoS attack is detected, targeting DNS services over UDP, indicating potential reconnaissance activity.", "llmRecommendation": "Verify DNS server logs and implement rate limiting on UDP traffic to prevent abuse." } ``` 使用自然语言提问 — 助手会通过工具调用告警查询来进行回答: ``` curl -s -X POST localhost:8080/api/chat -H "Content-Type: application/json" \ -d '{"question": "Give me an overall threat summary."}' # {"answer": "There have been a total of 24 alerts detected, with 12 from rule-based detection # and 12 from machine learning-based detection ... 4 alerts of each type ... 6 of each severity."} ``` 提供商是可配置的(`app.llm.provider = groq | gemini | anthropic`);Groq 的免费 层级无需信用卡。请参阅[配置](#configuration)。 ## 架构(阶段 1 — 摄取与存储) ``` CSV -> Kafka (iot.traffic.raw) -> StorageConsumer -> PostgreSQL ``` `CsvReplaySimulator` 以 schema 无关的方式读取 RT-IoT2022 CSV,并将每个 流作为 `TrafficEvent` 发布到 `iot.traffic.raw` topic(以源作为 key)。`StorageConsumer` 消费该 topic 并将每个流持久化到 `network_flows` 表中。 ## 阶段 2 — 基于规则的检测 ``` ┌─ StorageConsumer (group: storage) -> network_flows iot.traffic.raw ────────┤ └─ RuleBasedConsumer (group: rule-engine) -> iot.alerts │ AlertConsumer (group: alert-processor) -> alerts (retry + dead-letter: iot.alerts-dlt) ``` 规则是在 `application.yml` 中的 `app.rules` 下定义的配置驱动阈值,因此添加新的 检测不需要修改代码。由于 `rule-engine` 是独立于 `storage` 的 consumer group,这两个消费者会各自独立地接收每条流量消息。 ## 阶段 3 — 基于 ML 的检测 ``` ┌─ StorageConsumer (group: storage) -> network_flows iot.traffic.raw ┼─ RuleBasedConsumer (group: rule-engine) -> iot.alerts (RULE) └─ MlBasedConsumer (group: ml-engine) -> ml-service /predict │ attack_type != NORMAL ▼ iot.alerts (ML) ``` Python FastAPI 服务(`ml-service/`)提供离线训练的 XGBoost 分类器 (`train.py`,如果存在则使用 RT-IoT2022,否则使用合成的标注数据)。`MlBasedConsumer` 在其独立的 consumer group 下读取 `iot.traffic.raw`,通过 REST 带超时地调用模型, 并发布 `detectionSource=ML` 告警。如果 ML 服务不可用,调用会静默回退 — 存储和基于规则的检测不受影响。可以通过 `app.ml.enabled` 切换 ML。 使用 `docker compose up -d --build ml-service` 运行 ML 服务(或在本地运行:`cd ml-service && pip install -r requirements.txt && python train.py && uvicorn app:app --port 8000`)。 ## 阶段 4 — REST API、WebSocket 与仪表板 ``` iot.alerts ─┬─ AlertConsumer (group: alert-processor) -> alerts table └─ AlertWebSocketConsumer (group: dashboard) -> STOMP /topic/alerts -> browser ``` 只读 REST API 暴露了最近告警(`GET /api/alerts`,可选 `?severity=`)和 汇总统计(`GET /api/stats`)。新的 `dashboard` consumer group 通过 STOMP/WebSocket(`/ws` → `/topic/alerts`)中继每个告警。静态仪表板(`src/main/resources/static/`, Chart.js)展示了由 WebSocket 推送的实时告警表以及从 `/api/stats` 刷新的图表。 运行程序(`./mvnw spring-boot:run`)并打开 `http://localhost:8080/`。将流量发布到 `iot.traffic.raw`(见阶段 1)即可看到基于规则和 ML 的告警实时出现。 ## 阶段 5 — LLM 增强与自然语言聊天 ``` alerts table --> AlertEnrichmentScheduler --> SecurityAnalystService --> alerts.llm_explanation / llm_recommendation ``` **定时增强轮询器:** 后台调度器会轮询尚无 LLM 解释的告警(按创建时间排序,批处理大小可配置),并为每一个告警调用 `SecurityAnalystService`。该服务将告警元数据发送给 LLM,然后通过 `applyLlmEnrichment` 将返回的 `explanation` 和 `recommendation` 写回到 `alerts` 表中。 单条告警的错误会被捕获并记录日志,因此单次失败不会 中断整个批处理。 **自然语言聊天(`POST /api/chat`):** `NlQueryService` 基于现有的 REST 查询层(`AlertQueryTools`) 为 LLM 助手配置了工具调用(tool-calling)。助手可以通过调用正确的查询工具并 返回自然语言回答,来解答诸如“过去一小时内有多少 HIGH 级别的告警?”之类的问题。 当禁用 LLM 时,该端点将返回 `503 Service Unavailable`。 **配置(`application.yml`):** | 键 | 用途 | |-----|---------| | `app.llm.enabled` | 启用/禁用所有 LLM 功能(默认 `false`) | | `app.llm.provider` | LLM 提供商:`groq`(默认)、`gemini` 或 `anthropic` | | `app.llm.groq.model` | Groq 模型 ID | | `app.llm.gemini.model` | Gemini 模型 ID | | `app.llm.anthropic.model` | Anthropic 模型 ID | | `app.llm.timeout` | 单次调用超时时间 | | `app.llm.enrichment.batch-size` | 每个调度周期处理的告警数量 | | `app.llm.enrichment.poll-interval` | 调度器轮询间隔 | **运行时要求:** 要使用 LLM 功能,请在环境中为所选的 `app.llm.provider` 设置匹配的 API key(Gemini 对应 `GEMINI_API_KEY`,Anthropic 对应 `ANTHROPIC_API_KEY`)。 测试套件在 `app.llm.enabled=false` 下运行,且不需要 API key。 ## 技术栈 - Java 21, Spring Boot 4.1 - Apache Kafka (Confluent, KRaft 模式) - PostgreSQL 16 (JSONB 特性存储) - 用于集成测试的 Testcontainers - LangChain4j - 通过 Python 和 FastApi 提供的机器学习与 LLM 支持 - 流处理与安全加固 ## 配置 运行时配置通过环境变量提供。已提交的 [`.env.example`](.env.example) 文档记录了每个受支持的变量及其默认值。 1. 将模板复制到本地的 `.env`(已被 gitignored 忽略): cp .env.example .env 2. 填入相应的值。所有配置都有与内置的 `docker compose` 栈相匹配的可用默认值。要使用 LLM 功能,请通过 `APP_LLM_PROVIDER` 选择提供商(默认为 `groq`,或 `gemini` / `anthropic`),设置该提供商的 API key(`GROQ_API_KEY` — 可在 https://console.groq.com 免费获取 — `GEMINI_API_KEY` 或 `ANTHROPIC_API_KEY`),并设置 `APP_LLM_ENABLED=true`。 3. 在启动程序之前将 `.env` 加载到你的 shell 中 — Spring Boot **不会**自动读取 `.env`: set -a; source .env; set +a **切勿提交 `.env` 或真实的 key。** Groq 和 Gemini 都提供免费层级(无需 信用卡);`ANTHROPIC_API_KEY` 在 Anthropic Developer Platform 上按 token 计费,与任何 Claude Pro 订阅分开计算。 ## 本地运行 ``` docker compose up -d # kafka, kafka-ui (:8085), postgres (:5432) set -a; source .env; set +a # load your environment (see Configuration above) ./mvnw spring-boot:run # application (:8080) ``` 要重放数据集,请参阅 `data/README.md`,然后设置 `app.simulator.enabled=true`。 ## 测试 ``` ./mvnw test ``` 集成测试会通过 Testcontainers 启动真实的 Kafka 和 PostgreSQL 容器, 因此必须运行 Docker。 ## 数据来源 https://www.kaggle.com/datasets/supplejade/rt-iot2022real-time-internet-of-things
标签:DLL 劫持, IP 地址批量处理, Kafka, SonarQube插件, Spring Boot, XGBoost, 域名枚举, 大语言模型, 流处理, 测试用例, 物联网安全, 请求拦截, 软件成分分析, 逆向工具, 配置错误, 防御绕过