如何用Python实时解析持续写入的JSON日志并上传至数据库(非ELK)
持续解析JSON日志并上传新增内容到数据库的最佳方案
针对你的需求——无需ELK栈,持续监控/var/log下的JSON日志文件,将新增行解析后上传到数据库,我整理了几个生产环境中实用的方案,按实现复杂度和灵活性排序:
1. 最简易方案:tail -F + 自定义脚本
这是快速上手的首选,利用Linux原生tail命令跟踪文件新增内容,再通过管道传给自定义脚本处理:
核心思路
- 用
tail -F /var/log/your-service.json.log(-F会自动跟随日志轮转,比如logrotate生成的新文件)实时获取新增行 - 每行内容传给脚本,解析JSON格式、验证字段后批量插入数据库,减少频繁请求数据库的开销
示例脚本框架(以Python为例)
import sys import json import psycopg2 # 替换为你使用的数据库驱动,比如pymysql、sqlite3等 DB_CONFIG = { "host": "your-db-host", "user": "your-db-user", "password": "your-db-pass", "dbname": "your-db-name" } def batch_insert(entries): if not entries: return try: conn = psycopg2.connect(**DB_CONFIG) cursor = conn.cursor() # 批量插入语句,根据你的日志字段和表结构调整 insert_sql = """ INSERT INTO service_logs (event_time, level, message, source) VALUES (%s, %s, %s, %s) """ # 转换日志条目为数据库可接受的格式 data = [ (entry["timestamp"], entry["level"], entry["msg"], entry["source"]) for entry in entries ] cursor.executemany(insert_sql, data) conn.commit() except Exception as e: print(f"Database error: {str(e)}", file=sys.stderr) # 将失败的日志写入临时文件,后续可手动重试 with open("/tmp/failed_logs.txt", "a") as f: for entry in entries: f.write(json.dumps(entry) + "\n") finally: if conn: conn.close() def main(): batch = [] batch_size = 100 # 按需调整批量大小,平衡性能和实时性 for line in sys.stdin: line = line.strip() if not line: continue try: log_entry = json.loads(line) batch.append(log_entry) if len(batch) >= batch_size: batch_insert(batch) batch = [] except json.JSONDecodeError as e: print(f"Invalid JSON line: {line}, error: {str(e)}", file=sys.stderr) # 处理剩余的日志条目 batch_insert(batch) if __name__ == "__main__": main()
使用方式
在终端直接运行:
tail -F /var/log/your-service.json.log | python3 log_uploader.py
优缺点
- 优点:实现简单,无额外依赖,
tail -F天然支持日志轮转 - 缺点:脚本本身没有守护能力,需额外配置系统服务保证持续运行
2. 高效监听方案:inotifywait + 脚本
用inotify-tools的inotifywait命令监听文件修改事件,只在文件有新增内容时读取,比tail更节省资源,适合高频率写入的日志:
核心步骤
- 安装依赖:
sudo apt install inotify-tools(Debian/Ubuntu)或sudo yum install inotify-tools(RHEL/CentOS) - 编写脚本,监听文件的
modify事件,记录上次读取的位置,避免重复处理内容
关键逻辑示例(Bash + Python结合)
LOG_FILE="/var/log/your-service.json.log" LAST_POS=0 # 初始记录文件大小作为起始位置 if [ -f "$LOG_FILE" ]; then LAST_POS=$(stat -c %s "$LOG_FILE") fi while true; do # 等待文件修改事件 inotifywait -e modify "$LOG_FILE" > /dev/null # 获取当前文件大小 CURRENT_POS=$(stat -c %s "$LOG_FILE") # 读取新增内容并传给处理脚本 if [ "$CURRENT_POS" -gt "$LAST_POS" ]; then dd if="$LOG_FILE" bs=1 skip="$LAST_POS" count="$((CURRENT_POS - LAST_POS))" 2>/dev/null | python3 log_uploader.py LAST_POS=$CURRENT_POS fi # 处理日志轮转:原文件被删除/替换的情况 if [ ! -f "$LOG_FILE" ]; then inotifywait -e create "$(dirname "$LOG_FILE")" > /dev/null LAST_POS=0 fi done
优缺点
- 优点:仅在文件变化时处理,资源占用更低;精确控制读取位置,避免重复上传
- 缺点:需手动处理日志轮转的边界情况,依赖
inotify-tools
3. 灵活开发方案:Python watchdog库
如果偏好全Python实现,watchdog库可以监听文件系统事件,结合文件位置记录,实现更灵活的日志跟踪逻辑:
核心步骤
- 安装依赖:
pip install watchdog - 编写监听器,监听文件的修改、创建事件,维护读取位置,解析JSON后上传数据库
关键代码片段
from watchdog.observers import Observer from watchdog.events import FileSystemEventHandler import os import json import psycopg2 import time class LogMonitorHandler(FileSystemEventHandler): def __init__(self, log_path, db_config): self.log_path = log_path self.db_config = db_config # 初始化读取位置 self.last_pos = os.path.getsize(log_path) if os.path.exists(log_path) else 0 def on_modified(self, event): if event.src_path == self.log_path: self.process_new_content() def on_created(self, event): if event.src_path == self.log_path: self.last_pos = 0 self.process_new_content() def process_new_content(self): try: with open(self.log_path, "r") as f: f.seek(self.last_pos) lines = f.readlines() self.last_pos = f.tell() entries = [] for line in lines: line = line.strip() if not line: continue try: entries.append(json.loads(line)) except json.JSONDecodeError as e: print(f"Invalid JSON: {line}, error: {str(e)}", file=sys.stderr) if entries: self.batch_insert(entries) except Exception as e: print(f"Error processing log file: {str(e)}", file=sys.stderr) def batch_insert(self, entries): # 同第一个方案的批量插入逻辑 pass if __name__ == "__main__": log_path = "/var/log/your-service.json.log" db_config = {"host": "your-db-host", ...} event_handler = LogMonitorHandler(log_path, db_config) observer = Observer() observer.schedule(event_handler, path=os.path.dirname(log_path), recursive=False) observer.start() try: while True: time.sleep(1) except KeyboardInterrupt: observer.stop() observer.join()
优缺点
- 优点:全Python实现,扩展性强,方便集成日志过滤、字段转换等额外逻辑
- 缺点:需手动维护读取位置状态,日志轮转时要处理新文件的创建事件
通用注意事项
无论选择哪种方案,都要关注以下细节:
- 批量插入:设置合理的批量大小,平衡实时性和数据库性能
- 错误处理:JSON解析失败、数据库连接异常时,要记录错误日志,将失败的日志写入临时文件以便重试
- 日志轮转适配:配合
logrotate的copytruncate或create模式,确保脚本能识别新日志文件 - 守护进程化:用
systemd编写服务配置,保证脚本在系统重启或意外退出后自动重启
示例systemd服务文件(/etc/systemd/system/log-uploader.service):
启用服务:[Unit] Description=Service Log Uploader After=network.target [Service] User=your-user ExecStart=/usr/bin/python3 /path/to/log_uploader.py Restart=always RestartSec=5 [Install] WantedBy=multi-user.targetsudo systemctl daemon-reload && sudo systemctl enable --now log-uploader.service - 性能监控:定期查看脚本的CPU、内存占用,以及数据库插入成功率,按需调整配置
内容的提问来源于stack exchange,提问作者Pau Muñoz
相关产品推荐
相关产品推荐

