You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python实现蜜罐文本日志解析转CSV及去重方案

蜜罐日志解析转CSV及去重方案

日志解析实现

直接用Python标准库就能实现,不需要安装第三方依赖,完全适配你给出的固定日志格式:

import csv
import re
import os
from hashlib import md5

# 配置项,根据实际部署路径修改
RAW_LOG_PATH = "/var/log/honeypot/access.log"
CSV_OUTPUT_PATH = "./honeypot_logs.csv"
OFFSET_RECORD_PATH = "./.log_parse_offset"

# 匹配开头告警行的正则,提取来源IP端口、告警时间
ALERT_LINE_PATTERN = re.compile(r"INTRUSION ATTEMPT DETECTED! from (\d+\.\d+\.\d+\.\d+:\d+) \((\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2})\)")
SEPARATOR = "--------------------------"
CSV_FIELDS = [
    "alert_time", "src_addr", "request_line", "host", "x_forwarded_scheme",
    "x_forwarded_proto", "x_forwarded_for", "x_real_ip", "connection",
    "user_agent", "accept", "accept_encoding", "log_hash"
]

def get_last_offset():
    """读取上次解析到的文件偏移量,无记录则从头开始读"""
    if not os.path.exists(OFFSET_RECORD_PATH):
        return 0
    with open(OFFSET_RECORD_PATH, "r") as f:
        return int(f.read().strip())

def save_offset(offset):
    """保存当前解析完成的文件偏移量"""
    with open(OFFSET_RECORD_PATH, "w") as f:
        f.write(str(offset))

def parse_log_block(block_content):
    """解析单条日志块,返回结构化字段字典"""
    lines = [l.strip() for l in block_content.splitlines() if l.strip()]
    if not lines:
        return None
    # 解析首行告警信息
    alert_match = ALERT_LINE_PATTERN.match(lines[0])
    if not alert_match:
        return None
    src_addr, alert_time = alert_match.groups()
    log_data = {
        "alert_time": alert_time,
        "src_addr": src_addr
    }
    # 跳过告警行和分隔线行,从第三行开始解析请求头
    header_start_idx = 2
    request_line_parsed = False
    for line in lines[header_start_idx:]:
        if not request_line_parsed:
            log_data["request_line"] = line
            request_line_parsed = True
            continue
        if ":" not in line:
            continue
        # 统一转小写匹配键名,兼容头字段大小写混乱的情况
        key, value = line.split(":", 1)
        key = key.strip().lower().replace("-", "_")
        value = value.strip()
        if key == "host":
            log_data["host"] = value
        elif key == "x_forwarded_scheme":
            log_data["x_forwarded_scheme"] = value
        elif key == "x_forwarded_proto":
            log_data["x_forwarded_proto"] = value
        elif key == "x_forwarded_for":
            log_data["x_forwarded_for"] = value
        elif key == "x_real_ip":
            log_data["x_real_ip"] = value
        elif key == "connection":
            log_data["connection"] = value
        elif key == "user_agent":
            log_data["user_agent"] = value
        elif key == "accept":
            log_data["accept"] = value
        elif key == "accept_encoding":
            log_data["accept_encoding"] = value
    # 生成条目唯一哈希用于去重
    hash_str = f"{alert_time}{src_addr}{log_data.get('request_line', '')}{log_data.get('user_agent', '')}".encode()
    log_data["log_hash"] = md5(hash_str).hexdigest()
    # 补全缺失字段为空值
    for field in CSV_FIELDS:
        if field not in log_data:
            log_data[field] = ""
    return log_data

def main():
    last_offset = get_last_offset()
    # 加载已有CSV的哈希集合,做去重兜底
    exist_hashes = set()
    csv_exist = os.path.exists(CSV_OUTPUT_PATH)
    if csv_exist:
        with open(CSV_OUTPUT_PATH, "r", newline="", encoding="utf-8") as f:
            reader = csv.DictReader(f)
            for row in reader:
                exist_hashes.add(row["log_hash"])
    # 从上次偏移位置开始读取新增日志
    new_logs = []
    with open(RAW_LOG_PATH, "r", encoding="utf-8") as f:
        f.seek(last_offset)
        content = f.read()
        current_offset = f.tell()
        # 按固定告警头切分单条日志块
        log_blocks = content.split("INTRUSION ATTEMPT DETECTED!")
        for block in log_blocks:
            if not block.strip():
                continue
            full_block = "INTRUSION ATTEMPT DETECTED!" + block
            if SEPARATOR not in full_block:
                # 识别到未写完的半条日志,回退偏移量下次重读
                current_offset -= len(block.encode("utf-8"))
                break
            parsed = parse_log_block(full_block)
            if parsed and parsed["log_hash"] not in exist_hashes:
                new_logs.append(parsed)
                exist_hashes.add(parsed["log_hash"])
    # 写入新解析的结构化日志到CSV
    if new_logs:
        with open(CSV_OUTPUT_PATH, "a", newline="", encoding="utf-8") as f:
            writer = csv.DictWriter(f, fieldnames=CSV_FIELDS)
            if not csv_exist:
                writer.writeheader()
            writer.writerows(new_logs)
    # 保存最新解析偏移量
    save_offset(current_offset)

if __name__ == "__main__":
    main()

脚本可以直接配crontab定时任务,每分钟执行一次即可,遇到蜜罐未写完的半条日志会自动留存到下次解析,不会丢数据。

去重方案选型结论

你提到的两个初始方案都存在明显缺陷:

  • 解析完直接清空原日志:风险极高,脚本运行和蜜罐写入日志是并行操作,清空动作很容易删掉蜜罐刚写入、还没被脚本读取的新日志;一旦脚本解析逻辑出bug,原始日志被清空后没有任何回溯余地。
  • 仅靠写入前检测CSV重复:当日志量增长到十万条以上,每次启动全量读取CSV做比对会拖慢运行速度,而且如果没有读取位置记录,每次都从头解析全量日志,会产生大量无意义的重复操作。

更合理的是三层防护逻辑,上述脚本已经全部实现:

  • 根源避免重复读取:用独立文件记录每次解析结束时原始日志的字节偏移量,下次运行直接从上次结束的位置开始读,不会重复处理已经解析过的内容;遇到未写完的半条日志自动回退偏移量,等蜜罐写完整后下次再解析,不会丢数据。
  • 轻量兜底去重:启动时只加载已有CSV中所有条目的唯一哈希(哈希由告警时间、源地址、请求行、UA生成,同IP不同时间的正常多次攻击不会被误判为重复,只有完全一致的重复解析条目才会被拦截),新解析的条目先查哈希集合,不存在才写入,就算偏移量记录意外损坏,也不会写出重复内容。
  • 原始日志归档而非删除:不要直接清空原始日志,写个简单的定时任务,每周把超过30天的原始日志打包压缩存到备份目录,既不会长期占满磁盘,出问题时也能拿原始日志做校验回溯。

注意:不要把“同IP访问”判定为重复,同一攻击源IP在不同时间发起的多次尝试是需要留存的核心监测数据,去重只需要拦截完全相同的重复解析条目即可。

内容的提问来源于stack exchange,提问作者johnnyBdemented

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.31 22:33:11