diff --git a/Lib/log.py b/Lib/log.py index d430cf5..075719c 100644 --- a/Lib/log.py +++ b/Lib/log.py @@ -1,3 +1,3 @@ import logging -logger = logging.getLogger("ai-soc-framework") +logger = logging.getLogger("django") diff --git a/Lib/montior.py b/Lib/montior.py index 016acbb..877f4af 100644 --- a/Lib/montior.py +++ b/Lib/montior.py @@ -6,8 +6,10 @@ from apscheduler.schedulers.background import BackgroundScheduler +from CONFIG import REDIS_STREAM_STORE_DAYS from Lib.engine import Engine from Lib.log import logger +from Lib.redis_stream_api import RedisStreamAPI class MainMonitor(object): @@ -18,10 +20,22 @@ class MainMonitor(object): _background_threads = {} def __init__(self): - pass + self.engine = Engine() + self.redis_stream_api = RedisStreamAPI() + self.MainScheduler = BackgroundScheduler(timezone='Asia/Shanghai') def start(self): logger.info("后台服务启动") - engine = Engine() - engine.start() + self.MainScheduler.add_job(func=self.subscribe_clean_thread, + max_instances=1, + trigger='interval', + minutes=5, + id='subscribe_clean_thread') + self.MainScheduler.start() + + # engine + self.engine.start() logger.info("后台服务启动成功") + + def subscribe_clean_thread(self): + self.redis_stream_api.clean_redis_stream(max_age_days=REDIS_STREAM_STORE_DAYS) diff --git a/Lib/redis_stream_api.py b/Lib/redis_stream_api.py index f4888c2..52cc1a7 100644 --- a/Lib/redis_stream_api.py +++ b/Lib/redis_stream_api.py @@ -1,3 +1,4 @@ +import datetime import json from typing import Dict, Any, Optional, List @@ -263,3 +264,29 @@ class RedisStreamAPI: self.redis_client.close() except Exception as e: logger.exception(e) + + def clean_redis_stream(self, max_age_days=30): + """ + 清理Redis Stream中超过指定天数的键值对。 + """ + # 计算最老允许的时间戳,单位为毫秒 + # Unix时间戳(秒)* 1000 + logger.info(f"开始清理Redis Stream中超过 {max_age_days} 天的键值对...") + cutoff_timestamp_ms = int((datetime.datetime.now() - datetime.timedelta(days=max_age_days)).timestamp() * 1000) + + try: + for key in self.redis_client.scan_iter(match='*'): + # 检查键的类型是否为 stream + if self.redis_client.type(key) == 'stream': + # 使用 XTRIM 命令删除早于给定ID的条目 + # ID的格式是 `unix_time_ms-sequence_number` + # 我们可以使用 `unix_time_ms-0` 作为删除的上限ID + trim_id = f'{cutoff_timestamp_ms}-0' + + # `XTRIM` 带有 `MINID` 选项,用于删除ID小于指定ID的所有条目 + trimmed_count = self.redis_client.xtrim(key, minid=trim_id) + logger.info(f" 已从Stream '{key}' 中删除 {trimmed_count} 个过期条目。") + except Exception as e: + logger.exception(e) + finally: + logger.info("清理任务完成。") diff --git a/README.md b/README.md index 60e9be3..42a6e19 100644 --- a/README.md +++ b/README.md @@ -69,10 +69,7 @@ 您可以参考 `MODULES` 目录下的现有模块,开发自己的自动化流程。每个模块都是一个独立的 Python 文件,框架会自动加载并运行它。 - ## TODO -- Redis Stream 过期策略 - ## 许可证