1. 项目概述:为什么需要自动化巡检工具?

在运维、安全、内容审核乃至日常业务监控的领域里,我们常常面临一个重复且枯燥的任务:定期检查大量数据或系统状态,并从中识别出异常、风险或特定模式。传统的人工巡检不仅效率低下,容易因疲劳导致疏漏,而且难以应对海量、高速产生的数据。想象一下,一个安全工程师需要每天手动分析上千条日志,或者一个内容运营要实时监控多个渠道的用户反馈,这几乎是不可能完成的任务。

这正是自动化巡检工具的价值所在。它像一个不知疲倦的哨兵,7x24小时地执行预设的检查规则,将我们从重复劳动中解放出来,专注于更高阶的分析和决策。而随着大语言模型(LLM)能力的爆发,尤其是像 SecGPT-14B 这类专注于安全领域的模型出现,自动化巡检的“大脑”得到了质的飞跃。我们不再仅仅依赖简单的正则表达式或规则引擎,而是可以让 AI 模型理解上下文、进行逻辑推理、甚至生成初步的分析报告。

这个项目,就是教你如何用 Python 中最基础、最强大的 requests 库,去封装 SecGPT-14B 的 API,搭建一个属于你自己的、可定制化的智能巡检工具。无论你是想监控服务器日志里的异常行为,还是想自动审核用户生成内容中的风险,亦或是想从业务数据中自动提炼出关键洞察,这个框架都能为你提供一个坚实的起点。它不只是一个脚本,而是一个可扩展、可维护的工程化解决方案的雏形。

2. 核心思路与架构设计

在动手写代码之前,我们先要把整个工具的骨架搭起来。一个健壮的自动化巡检工具,其核心思路可以概括为:“ 数据获取 -> 智能分析 -> 结果处理 ”的闭环。我们的设计需要围绕这个闭环展开,确保每个环节都可靠、高效且易于扩展。

2.1 整体工作流设计

整个工具的工作流可以分解为以下几个步骤,我将用一个监控 Web 应用错误日志的场景来举例说明:

  1. 数据源对接 :工具需要从某个地方获取待检查的数据。这可能是:

    • 日志文件(如 Nginx access.log, 应用 error.log)。
    • 数据库(定期查询最新的记录)。
    • 消息队列(如 Kafka, RabbitMQ,实时消费消息)。
    • 第三方 API(如云监控平台的指标接口)。
    • 对于我们的教程,为了简化,我们可以从一个本地的文本文件或一个模拟的 API 端点开始。
  2. 数据预处理与分片 :原始数据往往很“脏”或者很大。直接扔给 API 可能超出其上下文长度限制(这正是热词中提到的 maximum context length 错误)。因此,我们需要:

    • 清洗 :去除无关字符、标准化格式。
    • 分片 :如果单次数据量过大,需要将其切割成适合模型处理的“块”。SecGPT-14B 可能有自己的 token 限制(比如热词中暗示的 1048565 tokens,但这通常是指模型的总容量,单次请求的上下文窗口会小得多,常见的是 4K, 8K, 32K等)。我们必须根据 API 文档确定单次请求的上下文上限。
    • 构造提示词(Prompt) :这是与模型沟通的“指令”。我们需要精心设计一个提示词,告诉模型我们给它的是什么数据,希望它做什么分析。例如:“你是一个安全分析专家。请分析以下服务器错误日志片段,找出所有可能表示遭受攻击(如 SQL 注入、路径遍历、暴力破解)的条目,并按风险等级(高、中、低)分类列出。”
  3. 调用 SecGPT-14B API :这是工具的核心。我们将使用 requests 库,按照 SecGPT-14B API 的规范,构造 HTTP 请求,发送我们预处理好的数据和提示词,并接收模型的回复。

  4. 解析与后处理 :模型返回的通常是 JSON 格式的文本。我们需要从中提取出结构化的分析结果(如列表、JSON 对象)。然后,根据业务逻辑进行后处理,比如:

    • 将高风险事件立即触发告警。
    • 将分析结果存储到数据库或文件中,用于后续审计和报表生成。
    • 对结果进行汇总统计。
  5. 调度与执行 :让整个流程自动、周期性地运行。这可以通过操作系统的定时任务(如 Linux 的 crontab, Windows 的任务计划程序),或者在 Python 脚本内部使用 schedule APScheduler 等库来实现。

2.2 为什么选择 Python + Requests?

  • Python :在自动化、数据处理和 AI 集成领域是事实上的标准语言。生态丰富,从数据处理(Pandas)到 HTTP 客户端(Requests)都有成熟的库,开发效率极高。
  • Requests :是 Python 中用于 HTTP 请求的“瑞士军刀”。它比标准库的 urllib 更简洁、更人性化,能够非常优雅地处理 API 调用中的各种细节,如超时设置、重试机制、会话保持、代理支持等。对于封装一个第三方 API 来说,它是绝佳的选择。

2.3 工具目录结构规划

一个良好的项目结构有助于代码管理和后期扩展。建议如下:

secgpt_inspector/
├── config.yaml (或 .env)       # 配置文件,存放 API密钥、端点URL等敏感信息
├── main.py                     # 主程序入口,负责调度
├── core/                       # 核心模块
│   ├── __init__.py
│   ├── data_fetcher.py        # 数据获取模块
│   ├── data_processor.py      # 数据预处理与分片模块
│   ├── secgpt_client.py       # SecGPT-14B API 封装客户端(核心)
│   └── result_handler.py      # 结果处理与告警模块
├── jobs/                       # 具体的巡检任务定义
│   ├── __init__.py
│   ├── log_inspection_job.py  # 日志巡检任务
│   └── content_audit_job.py   # 内容审核任务(示例)
├── utils/                      # 工具函数
│   ├── __init__.py
│   └── logger.py              # 日志记录工具
└── requirements.txt            # 项目依赖

这个结构将不同职责的代码分离,使得增加一个新的数据源或分析任务变得非常容易。

3. 核心模块实现:封装 SecGPT-14B API 客户端

这是整个工具最核心的部分。我们的目标是创建一个 SecGPTClient 类,它封装了所有与 SecGPT-14B API 交互的细节,对外提供简洁的方法。

3.1 初始化与配置管理

首先,我们不应该将 API 密钥、基础 URL 等硬编码在代码里。最佳实践是使用配置文件或环境变量。

使用 config.yaml 示例:

# config.yaml
secgpt:
  api_base: "https://api.example-secgpt.com/v1" # 替换为实际API地址
  api_key: "your-secret-api-key-here" # 替换为你的密钥
  model: "SecGPT-14B" # 或具体的模型名称
  timeout: 30 # 请求超时时间(秒)
  max_retries: 3 # 失败重试次数

secgpt_client.py 初始化部分:

import yaml
import requests
from typing import Optional, Dict, Any
import logging
from pathlib import Path

class SecGPTClient:
    def __init__(self, config_path: Optional[str] = None):
        """
        初始化 SecGPT-14B API 客户端。
        
        Args:
            config_path: 配置文件路径。如果为None,则尝试从环境变量读取。
        """
        self.logger = logging.getLogger(__name__)
        
        # 加载配置
        if config_path and Path(config_path).exists():
            with open(config_path, 'r', encoding='utf-8') as f:
                config = yaml.safe_load(f)
            api_config = config.get('secgpt', {})
        else:
            # 简易的环境变量读取(也可使用python-dotenv)
            api_config = {
                'api_base': os.getenv('SECGPT_API_BASE'),
                'api_key': os.getenv('SECGPT_API_KEY'),
                'model': os.getenv('SECGPT_MODEL', 'SecGPT-14B'),
                'timeout': int(os.getenv('SECGPT_TIMEOUT', '30')),
                'max_retries': int(os.getenv('SECGPT_MAX_RETRIES', '3')),
            }
        
        self.api_base = api_config['api_base'].rstrip('/')
        self.api_key = api_config['api_key']
        self.model = api_config['model']
        self.timeout = api_config['timeout']
        self.max_retries = api_config['max_retries']
        
        # 验证必要配置
        if not self.api_base or not self.api_key:
            raise ValueError("API Base URL 和 API Key 必须配置。请检查config.yaml或环境变量。")
        
        # 创建持久化会话,有助于连接复用和保持一些配置(如headers)
        self.session = requests.Session()
        self.session.headers.update({
            'Authorization': f'Bearer {self.api_key}',
            'Content-Type': 'application/json',
        })
        
        self.logger.info(f"SecGPTClient 初始化成功,模型: {self.model}, 端点: {self.api_base}")

注意 :这里我使用了 requests.Session() 。这是一个非常重要的技巧。使用 Session 对象可以在多次请求间保持 TCP 连接,避免重复的三次握手,显著提升性能。同时,Session 可以统一管理 headers、cookies、适配器等设置。

3.2 核心聊天补全接口封装

大多数 LLM API,包括 SecGPT-14B,都会提供一个类似于 OpenAI ChatCompletion 的接口。我们需要根据其官方文档来构造请求体。这里假设其接口与主流格式兼容。

    def chat_completion(self, 
                       messages: list, 
                       temperature: float = 0.1, 
                       max_tokens: Optional[int] = None,
                       **kwargs) -> Dict[str, Any]:
        """
        调用聊天补全接口。
        
        Args:
            messages: 消息列表,格式为 [{"role": "user", "content": "..."}, ...]
            temperature: 采样温度,控制随机性。越低输出越确定,越高越有创造性。巡检任务建议较低(如0.1-0.3)。
            max_tokens: 生成的最大token数。不指定则由模型决定。
            **kwargs: 其他可能的API参数,如top_p, stream等。
            
        Returns:
            API返回的完整JSON响应字典。
            
        Raises:
            requests.exceptions.RequestException: 网络或请求错误。
            ValueError: API返回业务逻辑错误。
        """
        url = f"{self.api_base}/chat/completions" # 假设端点路径,需按实际文档调整
        
        payload = {
            "model": self.model,
            "messages": messages,
            "temperature": temperature,
            **kwargs
        }
        if max_tokens is not None:
            payload["max_tokens"] = max_tokens
            
        # 添加重试逻辑,应对网络抖动或API限流(429错误)
        for attempt in range(self.max_retries):
            try:
                self.logger.debug(f"发送请求至 {url}, 尝试 {attempt + 1}/{self.max_retries}")
                response = self.session.post(
                    url, 
                    json=payload, 
                    timeout=self.timeout
                )
                response.raise_for_status() # 如果状态码不是200,抛出HTTPError
                
                result = response.json()
                
                # 检查API返回中是否包含错误(部分API即使HTTP 200也可能有业务错误)
                if "error" in result:
                    error_msg = result["error"].get("message", "Unknown API error")
                    error_code = result["error"].get("code")
                    self.logger.error(f"API业务错误: 代码={error_code}, 信息={error_msg}")
                    raise ValueError(f"API Error {error_code}: {error_msg}")
                    
                return result
                
            except requests.exceptions.Timeout:
                self.logger.warning(f"请求超时 (尝试 {attempt + 1}/{self.max_retries})")
                if attempt == self.max_retries - 1:
                    raise
                time.sleep(2 ** attempt) # 指数退避
            except requests.exceptions.HTTPError as e:
                status_code = e.response.status_code
                # 特别处理 429 Too Many Requests 错误(热词中提到)
                if status_code == 429:
                    retry_after = e.response.headers.get('Retry-After')
                    wait_time = int(retry_after) if retry_after else (2 ** attempt)
                    self.logger.warning(f"触发速率限制(429),等待 {wait_time} 秒后重试 (尝试 {attempt + 1}/{self.max_retries})")
                    time.sleep(wait_time)
                elif status_code == 400:
                    # 处理其他400错误,如上下文过长(热词中提到)
                    error_body = e.response.json()
                    self.logger.error(f"请求参数错误(400): {error_body}")
                    # 如果是上下文过长,这个错误应该在预处理阶段避免,这里可以直接抛出
                    raise ValueError(f"Bad Request: {error_body}") from e
                elif status_code == 402:
                    # 处理余额不足(热词中提到)
                    self.logger.error("API调用失败:余额不足(402)。请充值。")
                    raise
                else:
                    # 其他HTTP错误,可能不需要重试(如401认证失败,404接口不存在)
                    self.logger.error(f"HTTP请求失败: {status_code} - {e.response.text}")
                    raise
            except requests.exceptions.ConnectionError as e:
                self.logger.warning(f"连接错误 (尝试 {attempt + 1}/{self.max_retries}): {e}")
                if attempt == self.max_retries - 1:
                    raise
                time.sleep(2 ** attempt)
            except requests.exceptions.JSONDecodeError as e:
                self.logger.error(f"响应不是有效的JSON: {e}, 原始响应: {response.text[:200]}")
                raise ValueError("Invalid JSON response from API") from e

3.3 设计一个便捷的“巡检分析”方法

为了让调用更简单,我们可以在客户端上再封装一个高级方法,专门用于巡检任务。这个方法负责构造适合巡检的提示词,并解析返回结果。

    def analyze_for_inspection(self, 
                              data_chunk: str, 
                              inspection_prompt: str,
                              system_prompt: Optional[str] = None) -> str:
        """
        执行一次巡检分析。
        
        Args:
            data_chunk: 待分析的数据文本块。
            inspection_prompt: 针对该数据的具体分析指令(用户提示)。
            system_prompt: 系统提示词,定义模型角色。如果为None,使用默认值。
            
        Returns:
            模型生成的纯文本分析结果。
        """
        if system_prompt is None:
            system_prompt = """你是一个专业的自动化巡检助手。你的任务是严格按照用户的要求,对提供的数据进行分析。
            请只输出分析结论本身,不要添加任何解释性前缀(如“根据分析...”)、后缀或Markdown格式。保持输出简洁、结构化。"""
        
        messages = [
            {"role": "system", "content": system_prompt},
            {"role": "user", "content": f"{inspection_prompt}\n\n以下是待分析的数据:\n```\n{data_chunk}\n```"}
        ]
        
        try:
            response = self.chat_completion(
                messages=messages,
                temperature=0.1, # 巡检要求确定性高,温度设低
                max_tokens=1500  # 根据实际需要调整,控制输出长度
            )
            # 解析响应,提取模型回复内容
            # 假设响应格式为 {"choices": [{"message": {"content": "..."}}]}
            content = response['choices'][0]['message']['content'].strip()
            return content
        except Exception as e:
            self.logger.error(f"调用SecGPT-14B进行分析时失败: {e}")
            # 可以返回一个错误标识,或者根据策略重试、跳过等
            return f"[分析失败] 错误: {str(e)}"

4. 数据处理模块:应对大上下文与复杂输入

直接向模型扔一个 100MB 的日志文件是行不通的。我们需要一个智能的数据处理器。

4.1 数据分片策略

分片的核心原则是: 在保证语义完整性的前提下,将数据切割成模型能“消化”的大小。 对于日志、文本数据,可以按以下策略:

  1. 按行数/大小分片 :最简单的方法。例如,每 1000 行或每 50KB 文本作为一个块。缺点是可能切断一个完整的“事件”(比如一个多行错误堆栈)。
  2. 按时间窗口分片 :对于有时序的数据(如日志),按固定时间间隔(如每5分钟)切割。这能保证时间上的连续性。
  3. 按语义分片(高级) :利用更小的模型或规则,识别自然边界。例如,对于日志,可以在连续的时间戳间隙过大处切割;对于文章,可以在段落或章节处切割。

实现一个简单的按行数分片的处理器:

# data_processor.py
import re
from typing import List, Generator

class DataProcessor:
    def __init__(self, max_chunk_size: int = 2000, overlap_lines: int = 10):
        """
        初始化数据处理器。
        
        Args:
            max_chunk_size: 每个数据块的最大行数。
            overlap_lines: 块与块之间重叠的行数,防止关键信息被切断。
        """
        self.max_chunk_size = max_chunk_size
        self.overlap_lines = overlap_lines
        
    def chunk_by_lines(self, data: str) -> Generator[str, None, None]:
        """
        将文本数据按行分片。
        
        Args:
            data: 原始文本数据。
            
        Yields:
            分片后的文本块。
        """
        lines = data.splitlines()
        total_lines = len(lines)
        start = 0
        
        while start < total_lines:
            end = min(start + self.max_chunk_size, total_lines)
            chunk = '\n'.join(lines[start:end])
            yield chunk
            
            # 更新起始位置,考虑重叠
            start = end - self.overlap_lines if end < total_lines else end
            
    def clean_log_data(self, raw_log: str) -> str:
        """
        简单的日志清洗:去除空行、压缩多余空格。
        可根据实际日志格式扩展,如解析时间戳、IP、URL等。
        """
        # 去除行首行尾空格
        lines = [line.strip() for line in raw_log.splitlines()]
        # 过滤掉空行
        lines = [line for line in lines if line]
        # 简单的正则示例:移除 ANSI 颜色代码(如果日志有颜色)
        ansi_escape = re.compile(r'\x1B(?:[@-Z\\-_]|\[[0-?]*[ -/]*[@-~])')
        lines = [ansi_escape.sub('', line) for line in lines]
        return '\n'.join(lines)
        
    def estimate_tokens(self, text: str, method: str = 'approx') -> int:
        """
        粗略估计文本的token数量。
        注意:这是非常粗略的估计!精确计数需要用到模型的tokenizer。
        SecGPT-14B可能使用类似GPT的tokenizer,英文~1个token对应0.75个单词,中文~1-2个字符一个token。
        
        Args:
            text: 待估计的文本。
            method: 估计方法。'approx':使用简单规则;如需精确,可调用API的token计数端点(如果有)。
            
        Returns:
            估计的token数。
        """
        if method == 'approx':
            # 非常粗略的估计:对于中英文混合,按字符数 * 0.4 估算(经验值)
            return int(len(text) * 0.4)
        else:
            # 更精确的方法:如果API提供/tokenize端点,可以调用
            # 或者使用近似tokenizer,如tiktoken(针对OpenAI模型)
            # 这里为了通用性,先使用粗略估计
            return self.estimate_tokens(text, method='approx')

实操心得 overlap_lines (重叠行)是一个防止信息被切断的实用技巧。例如,一个错误堆栈跨越了分片边界,重叠可以确保它在两个块中都出现,虽然增加了少量重复计算,但保证了分析的完整性。对于安全巡检,宁可多查,不可漏查。

4.2 提示词工程:让模型理解你的意图

提示词是驱动模型产出的“方向盘”。一个糟糕的提示词会得到无用甚至误导的结果。对于自动化巡检,提示词需要 清晰、具体、结构化

一个针对 Web 访问日志安全巡检的提示词示例:

INSPECTION_PROMPT_TEMPLATE = """
请严格扮演一个Web安全分析专家的角色。

你的任务是对下面提供的Nginx访问日志片段进行安全分析。

**分析要求:**
1.  **识别潜在攻击**:检查每条日志记录,判断其是否可能属于以下任何一类攻击行为:
    - SQL注入(特征:包含 `'`, `"`, `UNION`, `SELECT`, `--`, `/*` 等SQL相关字符或模式)
    - 路径遍历/目录穿越(特征:包含 `../`, `..\\`, `/etc/passwd`, `win.ini` 等)
    - 跨站脚本(XSS)攻击(特征:包含 `<script>`, `javascript:`, `onerror=` 等)
    - 命令注入(特征:包含 `;`, `|`, `&`, `$(` 等系统命令拼接符)
    - 暴力破解或扫描(特征:短时间内同一IP对多个路径返回401/403,或使用常见漏洞扫描路径如 `/admin`, `/phpmyadmin`)
    - 其他可疑模式(如非常长的User-Agent,异常的HTTP方法如PROPFIND)

2.  **输出格式**:你必须以严格的JSON格式输出,且只输出这个JSON对象,不要有任何其他文字。JSON结构如下:
```json
{
  "summary": {
    "total_entries": <总日志条数>,
    "suspicious_entries": <可疑条数>,
    "risk_level": "<整体风险等级:低/中/高>"
  },
  "details": [
    {
      "line_number": <在原数据中的大致行号或索引>,
      "log_entry": "<可疑的原始日志行>",
      "attack_type": "<推测的攻击类型,如'SQL Injection'>",
      "confidence": "<置信度,低/中/高>",
      "reason": "<简要说明判断理由>"
    },
    // ... 更多可疑条目
  ]
}

日志数据如下:

{data_chunk}

现在,开始你的分析。 """


这个提示词做了几件关键的事:
1.  **明确角色和任务**:让模型进入“安全专家”的上下文。
2.  **给出具体检查项**:列出了要检测的攻击类型和特征,引导模型关注关键点。
3.  **强制结构化输出**:要求返回 JSON。这对于自动化处理至关重要!我们可以用 `json.loads()` 轻松地将模型输出转化为 Python 字典,进而进行后续处理。避免了从自由文本中提取信息的麻烦。
4.  **提供示例格式**:给出了 JSON 的具体结构,减少了模型“编造”格式的可能。

## 5. 组装完整巡检任务

现在,我们将客户端、数据处理器和具体的业务逻辑组装起来,形成一个完整的巡检任务。

### 5.1 定义一个日志巡检任务

```python
# jobs/log_inspection_job.py
import json
import time
from pathlib import Path
from typing import Dict, Any
from core.data_processor import DataProcessor
from core.secgpt_client import SecGPTClient
from core.result_handler import ResultHandler # 假设有一个结果处理器
import logging

class LogInspectionJob:
    def __init__(self, 
                 client: SecGPTClient, 
                 log_file_path: str,
                 prompt_template: str,
                 result_handler: ResultHandler):
        self.client = client
        self.log_file_path = Path(log_file_path)
        self.prompt_template = prompt_template
        self.result_handler = result_handler
        self.processor = DataProcessor(max_chunk_size=500, overlap_lines=5) # 每块500行日志
        self.logger = logging.getLogger(__name__)
        
    def run(self):
        """执行一次巡检任务。"""
        if not self.log_file_path.exists():
            self.logger.error(f"日志文件不存在: {self.log_file_path}")
            return
            
        self.logger.info(f"开始巡检日志文件: {self.log_file_path}")
        
        try:
            # 1. 读取日志
            with open(self.log_file_path, 'r', encoding='utf-8', errors='ignore') as f:
                raw_log = f.read()
                
            # 2. 数据清洗
            cleaned_log = self.processor.clean_log_data(raw_log)
            if not cleaned_log:
                self.logger.info("日志文件为空,跳过分析。")
                return
                
            # 3. 数据分片
            all_results = []
            for i, chunk in enumerate(self.processor.chunk_by_lines(cleaned_log)):
                self.logger.info(f"正在分析数据块 {i+1}...")
                
                # 4. 构造本次分析的提示词
                current_prompt = self.prompt_template.format(data_chunk=chunk)
                
                # 5. 调用API进行分析
                # 注意:这里直接使用 analyze_for_inspection 方法,它内部处理了system prompt
                analysis_result_text = self.client.analyze_for_inspection(
                    data_chunk=chunk, # 虽然prompt里包含了,这里也传一份给方法备用
                    inspection_prompt=current_prompt
                )
                
                # 6. 解析结果(期望是JSON)
                try:
                    # 尝试从返回文本中提取JSON。模型有时会在JSON外加 Markdown 代码块标记。
                    json_str = analysis_result_text.strip()
                    if json_str.startswith('```json'):
                        json_str = json_str[7:] # 去掉 ```json
                    if json_str.startswith('```'):
                        json_str = json_str[3:]
                    if json_str.endswith('```'):
                        json_str = json_str[:-3]
                        
                    result_dict = json.loads(json_str)
                    all_results.append(result_dict)
                    self.logger.debug(f"数据块 {i+1} 分析成功,发现 {result_dict.get('summary', {}).get('suspicious_entries', 0)} 条可疑记录。")
                    
                except json.JSONDecodeError as e:
                    self.logger.error(f"无法解析模型返回的JSON,数据块 {i+1}。原始输出:\n{analysis_result_text[:500]}")
                    # 可以记录原始输出,或者尝试用正则提取关键信息
                    all_results.append({"error": "JSON解析失败", "raw_output": analysis_result_text[:1000]})
                    
                # 礼貌性暂停,避免对API造成过大压力(尤其是免费或低配额账户)
                time.sleep(0.5)
                
            # 7. 汇总所有分片的结果
            final_summary = self._aggregate_results(all_results)
            
            # 8. 处理结果(存储、告警等)
            self.result_handler.handle(final_summary, source=str(self.log_file_path))
            
            self.logger.info(f"日志巡检完成。总计分析 {final_summary.get('total_entries_aggregated', 0)} 条记录,发现 {final_summary.get('total_suspicious_aggregated', 0)} 条可疑记录。")
            
        except Exception as e:
            self.logger.exception(f"执行日志巡检任务时发生未预期错误: {e}")
            
    def _aggregate_results(self, results: list) -> Dict[str, Any]:
        """聚合多个分片的分析结果。"""
        aggregated = {
            "total_entries_aggregated": 0,
            "total_suspicious_aggregated": 0,
            "risk_distribution": {"high": 0, "medium": 0, "low": 0},
            "all_details": [],
            "chunk_errors": 0
        }
        
        for res in results:
            if "error" in res:
                aggregated["chunk_errors"] += 1
                continue
                
            summary = res.get("summary", {})
            details = res.get("details", [])
            
            aggregated["total_entries_aggregated"] += summary.get("total_entries", 0)
            aggregated["total_suspicious_aggregated"] += summary.get("suspicious_entries", 0)
            
            # 简单风险聚合
            risk = summary.get("risk_level", "low").lower()
            if risk in aggregated["risk_distribution"]:
                aggregated["risk_distribution"][risk] += 1
                
            aggregated["all_details"].extend(details)
            
        # 判断整体风险(简单逻辑:有任何分片高风险则整体高风险)
        if aggregated["risk_distribution"]["high"] > 0:
            aggregated["overall_risk"] = "high"
        elif aggregated["risk_distribution"]["medium"] > 0:
            aggregated["overall_risk"] = "medium"
        else:
            aggregated["overall_risk"] = "low"
            
        return aggregated

5.2 主程序入口与调度

最后,我们需要一个 main.py 来把一切串起来,并实现定时调度。

# main.py
import logging
import schedule
import time
from core.secgpt_client import SecGPTClient
from jobs.log_inspection_job import LogInspectionJob
from core.result_handler import ResultHandler # 实现一个结果处理器

# 配置日志
logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
    handlers=[
        logging.FileHandler('inspector.log'),
        logging.StreamHandler()
    ]
)
logger = logging.getLogger(__name__)

def job():
    """定义要执行的巡检任务"""
    logger.info("=== 开始执行定时巡检任务 ===")
    
    # 1. 初始化客户端
    try:
        client = SecGPTClient(config_path='config.yaml')
    except Exception as e:
        logger.error(f"初始化SecGPT客户端失败: {e}")
        return
        
    # 2. 初始化结果处理器(示例:打印到控制台并保存到文件)
    class SimpleResultHandler(ResultHandler):
        def handle(self, result: dict, source: str):
            logger.info(f"巡检结果来自 {source}:")
            logger.info(f"  整体风险等级: {result.get('overall_risk')}")
            logger.info(f"  可疑事件总数: {result.get('total_suspicious_aggregated')}")
            # 保存到JSON文件
            import json
            from datetime import datetime
            filename = f"inspection_result_{datetime.now().strftime('%Y%m%d_%H%M%S')}.json"
            with open(filename, 'w', encoding='utf-8') as f:
                json.dump(result, f, indent=2, ensure_ascii=False)
            logger.info(f"  详细结果已保存至: {filename}")
            # 这里可以添加告警逻辑,如风险高时发送邮件、钉钉消息等
            if result.get('overall_risk') == 'high':
                logger.warning("⚠️  检测到高风险事件,建议立即人工介入检查!")
                
    result_handler = SimpleResultHandler()
    
    # 3. 定义并运行多个巡检任务(示例)
    jobs_to_run = [
        LogInspectionJob(
            client=client,
            log_file_path='/var/log/nginx/access.log', # 你的日志路径
            prompt_template=INSPECTION_PROMPT_TEMPLATE, # 前面定义的提示词模板
            result_handler=result_handler
        ),
        # 可以在这里添加更多任务,例如:
        # ContentAuditJob(client, data_source='...', result_handler=result_handler),
    ]
    
    for job_instance in jobs_to_run:
        try:
            job_instance.run()
        except Exception as e:
            logger.error(f"执行任务 {job_instance.__class__.__name__} 时出错: {e}", exc_info=True)
            
    logger.info("=== 定时巡检任务执行完毕 ===\n")

if __name__ == '__main__':
    logger.info("自动化巡检工具启动...")
    
    # 立即执行一次
    job()
    
    # 设置定时任务(例如每30分钟执行一次)
    schedule.every(30).minutes.do(job)
    # 或者每天凌晨2点执行
    # schedule.every().day.at("02:00").do(job)
    
    logger.info("调度器已启动,将按计划运行。")
    
    # 保持程序运行
    while True:
        schedule.run_pending()
        time.sleep(60) # 每分钟检查一次是否有任务需要执行

6. 避坑指南与高级技巧

在实际部署和运行中,你肯定会遇到各种问题。以下是我从经验中总结出的关键点和解决方案。

6.1 API 调用相关陷阱

  • 速率限制 (429 Too Many Requests) :这是最常遇到的问题。热词中也提到了 exceeded retry limit, last status: 429

    • 应对策略
      1. 指数退避重试 :我们的客户端代码已经实现了。 time.sleep(2 ** attempt) 让重试间隔越来越长。
      2. 识别 Retry-After 头 :有些 API 会在 429 响应中携带 Retry-After 头,告诉你要等多久。我们的代码优先使用这个值。
      3. 降低请求频率 :在循环调用 API 时,主动添加 time.sleep(1) 之类的间隔。
      4. 使用请求队列 :对于大规模任务,实现一个队列来控制并发请求数。
  • 上下文长度超限 (400: context length) :热词中提到了这个错误。模型有单次请求的 token 上限。

    • 根本解决 :在数据预处理阶段做好分片。使用 DataProcessor.estimate_tokens 进行粗略估计,确保每个数据块加上提示词后的总 token 数远低于限制(建议留出 20%-30% 的余量给模型的输出)。
    • 动态调整 :如果某个分片仍然超限,可以在代码中捕获这个特定错误,然后自动将分片大小减半并重试。
  • 余额不足 (402 Insufficient Balance)

    • 预防 :在调用前,如果 API 提供查询余额的接口,可以先查询。或者在配置中设置月度预算和告警。
    • 处理 :在代码中捕获 402 错误,并立即停止任务,发送高优先级告警。
  • 连接中断 (Connection closed mid-response) :网络不稳定或服务器端问题可能导致响应不完整。

    • 应对 :使用 requests timeout 参数,并实现重试逻辑。对于特别重要的分析,可以考虑将接收到的部分响应暂存,并在重试时告知模型“请继续上一次的回复”。

6.2 提示词与结果处理优化

  • 模型不按格式输出 :即使你要求输出 JSON,模型有时还是会加上解释性文字。

    • 技巧1 :在系统提示词中强调“ 只输出JSON,不要有任何其他文字 ”。
    • 技巧2 :在用户提示词中,将 JSON 结构放在最后,并说“ 请用以下JSON格式输出你的分析结果: ”。
    • 技巧3 :使用 输出解析库 ,如 LangChain 的 OutputParser Pydantic ,它们能更鲁棒地处理模型输出。对于简单项目,用 json.loads() 配合字符串清理(如去掉代码块标记)通常也够用。
  • 分析结果不准或遗漏

    • 提供示例 :在提示词中给出1-2个正面和反面的例子(Few-shot Learning),能极大提升模型在特定任务上的表现。
    • 分而治之 :如果任务非常复杂,不要指望一个提示词解决所有问题。可以设计多轮对话,先让模型提取关键信息,再让模型基于提取的信息做判断。
    • 后处理校验 :对模型输出的“置信度”进行过滤,只处理“高”置信度的结果。对于中低置信度的,可以标记出来供人工复核。

6.3 性能与稳定性

  • 异步并发 :如果需要处理大量独立的数据块,同步请求会非常慢。可以使用 asyncio aiohttp 库实现异步并发请求,大幅提升效率。但要注意 API 的并发限制。
    # 简化的异步示例思路
    import aiohttp
    import asyncio
    
    async def analyze_chunk_async(session, chunk, prompt):
        async with session.post(api_url, json=payload) as resp:
            return await resp.json()
    
    # 创建任务列表并并发执行
    tasks = [analyze_chunk_async(session, chunk, prompt) for chunk in chunks]
    results = await asyncio.gather(*tasks, return_exceptions=True)
    
  • 结果持久化与状态管理 :对于长期运行的巡检,需要记录哪些数据已经分析过,避免重复分析。可以将已处理数据的 ID 或哈希值存储到数据库或文件中。
  • 完善的日志 :日志是排查问题的生命线。不仅要记录信息、错误,还要记录关键的请求参数和响应片段(注意脱敏API密钥)。使用 logging 模块,并设置合理的日志级别和轮转策略。

6.4 扩展方向

这个基础框架可以朝多个方向扩展:

  1. 支持更多数据源 :继承一个 BaseFetcher 类,实现从数据库、Kafka、云存储等获取数据的子类。
  2. 工作流引擎 :引入像 Prefect Airflow 这样的工作流调度器,替代简单的 schedule 库,实现更复杂的依赖关系、任务监控和错误恢复。
  3. 多模型路由与降级 :除了 SecGPT-14B,还可以集成其他开源或商业模型(如 GLM、DeepSeek)。在主模型调用失败或成本过高时,自动切换到备用模型。
  4. 构建Web界面 :使用 FastAPI Flask 构建一个简单的管理界面,用于查看巡检历史、配置任务、手动触发分析等。
  5. 向量数据库结合 :将历史分析结果和已知的威胁情报存入向量数据库(如 ChromaDB , Milvus )。新的巡检结果可以先进行向量相似度搜索,如果发现与历史已知威胁高度相似,可以直接给出结论,减少对 LLM 的调用,提高效率并降低成本。

这个用 Python requests 封装 SecGPT-14B API 构建的自动化巡检工具,其核心价值在于提供了一个清晰、可扩展的模式。它证明了,即使不使用复杂的 AI 框架,仅凭标准的网络请求库和清晰的工程思维,我们也能将强大的大模型能力无缝集成到现有的自动化流程中,创造出真正智能的解决方案。从今天开始,试着用它去解决你工作中那个最重复、最耗时的检查任务吧。

更多推荐