Python Requests封装SecGPT-14B API构建智能自动化巡检工具
1. 项目概述:为什么需要自动化巡检工具?
在运维、安全、内容审核乃至日常业务监控的领域里,我们常常面临一个重复且枯燥的任务:定期检查大量数据或系统状态,并从中识别出异常、风险或特定模式。传统的人工巡检不仅效率低下,容易因疲劳导致疏漏,而且难以应对海量、高速产生的数据。想象一下,一个安全工程师需要每天手动分析上千条日志,或者一个内容运营要实时监控多个渠道的用户反馈,这几乎是不可能完成的任务。
这正是自动化巡检工具的价值所在。它像一个不知疲倦的哨兵,7x24小时地执行预设的检查规则,将我们从重复劳动中解放出来,专注于更高阶的分析和决策。而随着大语言模型(LLM)能力的爆发,尤其是像 SecGPT-14B 这类专注于安全领域的模型出现,自动化巡检的“大脑”得到了质的飞跃。我们不再仅仅依赖简单的正则表达式或规则引擎,而是可以让 AI 模型理解上下文、进行逻辑推理、甚至生成初步的分析报告。
这个项目,就是教你如何用 Python 中最基础、最强大的 requests 库,去封装 SecGPT-14B 的 API,搭建一个属于你自己的、可定制化的智能巡检工具。无论你是想监控服务器日志里的异常行为,还是想自动审核用户生成内容中的风险,亦或是想从业务数据中自动提炼出关键洞察,这个框架都能为你提供一个坚实的起点。它不只是一个脚本,而是一个可扩展、可维护的工程化解决方案的雏形。
2. 核心思路与架构设计
在动手写代码之前,我们先要把整个工具的骨架搭起来。一个健壮的自动化巡检工具,其核心思路可以概括为:“ 数据获取 -> 智能分析 -> 结果处理 ”的闭环。我们的设计需要围绕这个闭环展开,确保每个环节都可靠、高效且易于扩展。
2.1 整体工作流设计
整个工具的工作流可以分解为以下几个步骤,我将用一个监控 Web 应用错误日志的场景来举例说明:
-
数据源对接 :工具需要从某个地方获取待检查的数据。这可能是:
- 日志文件(如 Nginx access.log, 应用 error.log)。
- 数据库(定期查询最新的记录)。
- 消息队列(如 Kafka, RabbitMQ,实时消费消息)。
- 第三方 API(如云监控平台的指标接口)。
- 对于我们的教程,为了简化,我们可以从一个本地的文本文件或一个模拟的 API 端点开始。
-
数据预处理与分片 :原始数据往往很“脏”或者很大。直接扔给 API 可能超出其上下文长度限制(这正是热词中提到的
maximum context length错误)。因此,我们需要:- 清洗 :去除无关字符、标准化格式。
- 分片 :如果单次数据量过大,需要将其切割成适合模型处理的“块”。SecGPT-14B 可能有自己的 token 限制(比如热词中暗示的 1048565 tokens,但这通常是指模型的总容量,单次请求的上下文窗口会小得多,常见的是 4K, 8K, 32K等)。我们必须根据 API 文档确定单次请求的上下文上限。
- 构造提示词(Prompt) :这是与模型沟通的“指令”。我们需要精心设计一个提示词,告诉模型我们给它的是什么数据,希望它做什么分析。例如:“你是一个安全分析专家。请分析以下服务器错误日志片段,找出所有可能表示遭受攻击(如 SQL 注入、路径遍历、暴力破解)的条目,并按风险等级(高、中、低)分类列出。”
-
调用 SecGPT-14B API :这是工具的核心。我们将使用
requests库,按照 SecGPT-14B API 的规范,构造 HTTP 请求,发送我们预处理好的数据和提示词,并接收模型的回复。 -
解析与后处理 :模型返回的通常是 JSON 格式的文本。我们需要从中提取出结构化的分析结果(如列表、JSON 对象)。然后,根据业务逻辑进行后处理,比如:
- 将高风险事件立即触发告警。
- 将分析结果存储到数据库或文件中,用于后续审计和报表生成。
- 对结果进行汇总统计。
-
调度与执行 :让整个流程自动、周期性地运行。这可以通过操作系统的定时任务(如 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 数据分片策略
分片的核心原则是: 在保证语义完整性的前提下,将数据切割成模型能“消化”的大小。 对于日志、文本数据,可以按以下策略:
- 按行数/大小分片 :最简单的方法。例如,每 1000 行或每 50KB 文本作为一个块。缺点是可能切断一个完整的“事件”(比如一个多行错误堆栈)。
- 按时间窗口分片 :对于有时序的数据(如日志),按固定时间间隔(如每5分钟)切割。这能保证时间上的连续性。
- 按语义分片(高级) :利用更小的模型或规则,识别自然边界。例如,对于日志,可以在连续的时间戳间隙过大处切割;对于文章,可以在段落或章节处切割。
实现一个简单的按行数分片的处理器:
# 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。- 应对策略 :
- 指数退避重试 :我们的客户端代码已经实现了。
time.sleep(2 ** attempt)让重试间隔越来越长。 - 识别 Retry-After 头 :有些 API 会在 429 响应中携带
Retry-After头,告诉你要等多久。我们的代码优先使用这个值。 - 降低请求频率 :在循环调用 API 时,主动添加
time.sleep(1)之类的间隔。 - 使用请求队列 :对于大规模任务,实现一个队列来控制并发请求数。
- 指数退避重试 :我们的客户端代码已经实现了。
- 应对策略 :
-
上下文长度超限 (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 扩展方向
这个基础框架可以朝多个方向扩展:
- 支持更多数据源 :继承一个
BaseFetcher类,实现从数据库、Kafka、云存储等获取数据的子类。 - 工作流引擎 :引入像
Prefect或Airflow这样的工作流调度器,替代简单的schedule库,实现更复杂的依赖关系、任务监控和错误恢复。 - 多模型路由与降级 :除了 SecGPT-14B,还可以集成其他开源或商业模型(如 GLM、DeepSeek)。在主模型调用失败或成本过高时,自动切换到备用模型。
- 构建Web界面 :使用
FastAPI或Flask构建一个简单的管理界面,用于查看巡检历史、配置任务、手动触发分析等。 - 向量数据库结合 :将历史分析结果和已知的威胁情报存入向量数据库(如
ChromaDB,Milvus)。新的巡检结果可以先进行向量相似度搜索,如果发现与历史已知威胁高度相似,可以直接给出结论,减少对 LLM 的调用,提高效率并降低成本。
这个用 Python requests 封装 SecGPT-14B API 构建的自动化巡检工具,其核心价值在于提供了一个清晰、可扩展的模式。它证明了,即使不使用复杂的 AI 框架,仅凭标准的网络请求库和清晰的工程思维,我们也能将强大的大模型能力无缝集成到现有的自动化流程中,创造出真正智能的解决方案。从今天开始,试着用它去解决你工作中那个最重复、最耗时的检查任务吧。
更多推荐
所有评论(0)