本文目录导读:

我来介绍几种清理超时任务数据的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()
使用建议
- 选择合适的超时时间: 根据业务需求设置合理的超时时间
- 错误处理: 添加完善的异常处理和日志记录
- 性能考虑: 大量数据时考虑分批处理
- 监控告警: 添加清理结果的监控和告警
- 事务安全: 数据库操作时确保事务完整性
- 资源释放: 清理时注意释放相关资源
根据你的具体需求选择合适的方案,如果需要更具体的实现,请提供更多细节。