Python脚本如何清理超时任务数据

wen python案例 23

本文目录导读:

Python脚本如何清理超时任务数据

  1. 基础时间戳清理方案
  2. 数据库清理方案
  3. 定时清理脚本
  4. 文件系统清理
  5. 综合清理管理器
  6. 使用建议

我来介绍几种清理超时任务数据的Python脚本方案:

基础时间戳清理方案

import time
import datetime
import json
from typing import Dict, Any
class TaskCleaner:
    def __init__(self, timeout_seconds: int = 3600):
        self.timeout = timeout_seconds
    def clean_timeout_tasks(self, tasks: Dict[str, Any]) -> Dict[str, Any]:
        """
        清理超时任务
        :param tasks: 任务字典 {task_id: {'create_time': timestamp, ...}}
        :return: 清理后的任务字典
        """
        current_time = time.time()
        cleaned_tasks = {}
        for task_id, task_data in tasks.items():
            create_time = task_data.get('create_time', 0)
            # 检查是否超时
            if current_time - create_time < self.timeout:
                cleaned_tasks[task_id] = task_data
            else:
                print(f"移除超时任务: {task_id}, 创建时间: {datetime.datetime.fromtimestamp(create_time)}")
        return cleaned_tasks
    def clean_with_date_compare(self, tasks: Dict[str, Any]) -> Dict[str, Any]:
        """使用datetime进行比较"""
        cutoff_time = datetime.datetime.now() - datetime.timedelta(seconds=self.timeout)
        cleaned_tasks = {}
        for task_id, task_data in tasks.items():
            create_time = task_data.get('create_time')
            if isinstance(create_time, (int, float)):
                task_datetime = datetime.datetime.fromtimestamp(create_time)
            else:
                task_datetime = create_time
            if task_datetime > cutoff_time:
                cleaned_tasks[task_id] = task_data
        return cleaned_tasks
# 使用示例
if __name__ == "__main__":
    # 模拟任务数据
    tasks = {
        "task_1": {"create_time": time.time() - 100, "status": "running"},
        "task_2": {"create_time": time.time() - 7200, "status": "completed"},
        "task_3": {"create_time": time.time() - 5000, "status": "pending"},
        "task_4": {"create_time": time.time() - 500, "status": "running"}
    }
    cleaner = TaskCleaner(timeout_seconds=3600)  # 1小时超时
    cleaned = cleaner.clean_timeout_tasks(tasks)
    print(f"清理后剩余任务数: {len(cleaned)}")

数据库清理方案

import sqlite3
import pymongo
from datetime import datetime, timedelta
import logging
class DatabaseTaskCleaner:
    def __init__(self, db_type='sqlite', timeout_hours=24):
        self.db_type = db_type
        self.timeout = timedelta(hours=timeout_hours)
        logging.basicConfig(level=logging.INFO)
        self.logger = logging.getLogger(__name__)
    def clean_sqlite(self, db_path: str, table: str):
        """清理SQLite数据库中的超时任务"""
        conn = sqlite3.connect(db_path)
        cursor = conn.cursor()
        try:
            cutoff_time = datetime.now() - self.timeout
            cursor.execute(f"""
                DELETE FROM {table} 
                WHERE create_time < ?
            """, (cutoff_time,))
            deleted_count = cursor.rowcount
            conn.commit()
            self.logger.info(f"已清理 {deleted_count} 条超时任务")
        except Exception as e:
            conn.rollback()
            self.logger.error(f"清理失败: {e}")
        finally:
            conn.close()
    def clean_mongodb(self, collection, query_filter=None):
        """清理MongoDB中的超时任务"""
        if query_filter is None:
            query_filter = {}
        cutoff_time = datetime.now() - self.timeout
        query_filter['create_time'] = {'$lt': cutoff_time}
        result = collection.delete_many(query_filter)
        self.logger.info(f"已清理 {result.deleted_count} 条超时任务")
        return result.deleted_count
    def clean_redis(self, redis_client, task_prefix: str):
        """清理Redis中的超时任务"""
        import redis
        cutoff_time = datetime.now().timestamp() - self.timeout.total_seconds()
        deleted_count = 0
        # 查找所有匹配的任务键
        cursor = '0'
        while cursor != 0:
            cursor, keys = redis_client.scan(cursor=cursor, 
                                            match=f"{task_prefix}*")
            for key in keys:
                # 获取任务的创建时间
                task_data = redis_client.hgetall(key)
                if task_data and b'create_time' in task_data:
                    create_time = float(task_data[b'create_time'])
                    if create_time < cutoff_time:
                        redis_client.delete(key)
                        deleted_count += 1
        self.logger.info(f"已清理 {deleted_count} 条超时任务")
        return deleted_count
# 使用示例
def main():
    # SQLite示例
    sqlite_cleaner = DatabaseTaskCleaner(db_type='sqlite', timeout_hours=24)
    sqlite_cleaner.clean_sqlite('tasks.db', 'task_queue')
    # MongoDB示例
    from pymongo import MongoClient
    client = MongoClient('mongodb://localhost:27017/')
    db = client['task_db']
    collection = db['tasks']
    mongo_cleaner = DatabaseTaskCleaner(db_type='mongodb', timeout_hours=24)
    mongo_cleaner.clean_mongodb(collection, {'status': 'pending'})

定时清理脚本

import schedule
import time
import logging
from datetime import datetime, timedelta
class ScheduledTaskCleaner:
    def __init__(self):
        logging.basicConfig(
            level=logging.INFO,
            format='%(asctime)s - %(levelname)s - %(message)s'
        )
        self.logger = logging.getLogger(__name__)
    def clean_tasks(self):
        """清理超时任务的执行函数"""
        try:
            # 这里调用实际的清理逻辑
            self.logger.info("开始清理超时任务...")
            # 模拟清理过程
            current_time = datetime.now()
            cutoff_time = current_time - timedelta(hours=24)
            self.logger.info(f"清理 {cutoff_time} 之前的任务")
            # 实际清理代码
            # deleted_count = your_clean_function()
            deleted_count = 10  # 模拟清理数量
            self.logger.info(f"清理完成,共清理 {deleted_count} 条超时任务")
        except Exception as e:
            self.logger.error(f"清理任务失败: {e}")
    def run_schedule(self):
        """运行定时任务"""
        # 每天凌晨2点执行清理
        schedule.every().day.at("02:00").do(self.clean_tasks)
        # 或者每6小时执行一次
        # schedule.every(6).hours.do(self.clean_tasks)
        self.logger.info("定时清理任务已启动")
        while True:
            schedule.run_pending()
            time.sleep(60)
    def run_once(self):
        """立即执行一次清理"""
        self.clean_tasks()
# 使用示例
if __name__ == "__main__":
    cleaner = ScheduledTaskCleaner()
    cleaner.run_once()  # 立即执行一次
    # cleaner.run_schedule()  # 启动定时任务

文件系统清理

import os
import shutil
import glob
from pathlib import Path
from datetime import datetime, timedelta
import logging
class FileTaskCleaner:
    def __init__(self, task_dir: str, file_pattern: str = "*.task"):
        self.task_dir = Path(task_dir)
        self.file_pattern = file_pattern
        logging.basicConfig(level=logging.INFO)
        self.logger = logging.getLogger(__name__)
    def clean_old_files(self, timeout_hours: int = 24):
        """清理超时的任务文件"""
        cutoff_time = datetime.now() - timedelta(hours=timeout_hours)
        cleaned_count = 0
        for file_path in self.task_dir.glob(self.file_pattern):
            # 获取文件修改时间
            file_mtime = datetime.fromtimestamp(file_path.stat().st_mtime)
            if file_mtime < cutoff_time:
                try:
                    file_path.unlink()  # 删除文件
                    cleaned_count += 1
                    self.logger.info(f"删除超时文件: {file_path}")
                except Exception as e:
                    self.logger.error(f"删除文件失败 {file_path}: {e}")
        self.logger.info(f"共清理 {cleaned_count} 个超时任务文件")
        return cleaned_count
    def clean_old_directories(self, timeout_hours: int = 48):
        """清理超时的任务目录"""
        cutoff_time = datetime.now() - timedelta(hours=timeout_hours)
        cleaned_count = 0
        for dir_path in self.task_dir.iterdir():
            if not dir_path.is_dir():
                continue
            # 获取目录修改时间
            dir_mtime = datetime.fromtimestamp(dir_path.stat().st_mtime)
            if dir_mtime < cutoff_time:
                try:
                    shutil.rmtree(dir_path)  # 删除整个目录
                    cleaned_count += 1
                    self.logger.info(f"删除超时目录: {dir_path}")
                except Exception as e:
                    self.logger.error(f"删除目录失败 {dir_path}: {e}")
        return cleaned_count
# 使用示例
if __name__ == "__main__":
    cleaner = FileTaskCleaner("/path/to/tasks", "*.task")
    cleaner.clean_old_files(timeout_hours=24)  # 清理24小时前的文件
    cleaner.clean_old_directories(timeout_hours=48)  # 清理48小时前的目录

综合清理管理器

import threading
import time
from typing import List, Callable
import logging
class TaskCleanupManager:
    def __init__(self):
        self.cleaners = []
        self.logger = logging.getLogger(__name__)
        self._stop_event = threading.Event()
    def add_cleaner(self, cleaner_func: Callable, interval: int = 3600):
        """添加清理器"""
        self.cleaners.append({
            'func': cleaner_func,
            'interval': interval,
            'last_run': 0
        })
    def run_cleanup(self):
        """执行一次性清理"""
        for cleaner in self.cleaners:
            try:
                self.logger.info(f"执行清理: {cleaner['func'].__name__}")
                cleaner['func']()
                cleaner['last_run'] = time.time()
            except Exception as e:
                self.logger.error(f"清理失败: {e}")
    def start_scheduled(self):
        """启动定时清理"""
        def cleanup_loop():
            while not self._stop_event.is_set():
                current_time = time.time()
                for cleaner in self.cleaners:
                    elapsed = current_time - cleaner['last_run']
                    if elapsed >= cleaner['interval']:
                        try:
                            self.logger.info(f"定时清理: {cleaner['func'].__name__}")
                            cleaner['func']()
                            cleaner['last_run'] = current_time
                        except Exception as e:
                            self.logger.error(f"定时清理失败: {e}")
                # 每分钟检查一次
                time.sleep(60)
        thread = threading.Thread(target=cleanup_loop, daemon=True)
        thread.start()
        self.logger.info("定时清理服务已启动")
        return thread
    def stop(self):
        """停止定时清理"""
        self._stop_event.set()
        self.logger.info("清理服务已停止")
# 使用示例
def clean_memory_tasks():
    """清理内存中的超时任务示例"""
    print("清理内存任务...")
def clean_database_tasks():
    """清理数据库中的超时任务示例"""
    print("清理数据库任务...")
if __name__ == "__main__":
    manager = TaskCleanupManager()
    # 添加不同频率的清理任务
    manager.add_cleaner(clean_memory_tasks, interval=3600)  # 每小时
    manager.add_cleaner(clean_database_tasks, interval=86400)  # 每天
    # 启动定时清理
    manager.start_scheduled()
    try:
        # 保持程序运行
        while True:
            time.sleep(10)
    except KeyboardInterrupt:
        manager.stop()

使用建议

  1. 选择合适的超时时间: 根据业务需求设置合理的超时时间
  2. 错误处理: 添加完善的异常处理和日志记录
  3. 性能考虑: 大量数据时考虑分批处理
  4. 监控告警: 添加清理结果的监控和告警
  5. 事务安全: 数据库操作时确保事务完整性
  6. 资源释放: 清理时注意释放相关资源

根据你的具体需求选择合适的方案,如果需要更具体的实现,请提供更多细节。

抱歉,评论功能暂时关闭!