[Feat] (concurrency) 增加并发处理

Merge pull request #16 from SXP-Simon/feat/concurrency
This commit is contained in:
Helian Nuits
2025-09-24 00:00:46 +08:00
committed by GitHub
3 changed files with 148 additions and 67 deletions
+38 -15
View File
@@ -79,24 +79,47 @@ class QQGroupDailyAnalysis(Star):
logger.info(f"Bot管理器状态: {status}")
except Exception as e:
logger.error(f"延迟启动调度器失败: {e}")
logger.debug(f"延迟启动调度器失败,可能由于短时间内多次更新插件配置: {e}")
async def _reload_config_and_restart_scheduler(self):
"""重新加载配置并重启调度器"""
async def terminate(self):
"""插件被卸载/停用时调用,清理资源"""
try:
# 重新加载配置
config_manager.reload_config()
logger.info(f"重新加载配置: 自动分析={config_manager.get_enable_auto_analysis()}")
# 重启调度器
await auto_scheduler.restart_scheduler()
logger.info("配置重载和调度器重启完成")
logger.info("开始清理QQ群日常分析插件资源...")
global auto_scheduler, bot_manager, message_analyzer, report_generator, config_manager
# 停止自动调度器
if auto_scheduler:
logger.info("正在停止自动调度器...")
await auto_scheduler.stop_scheduler()
logger.info("自动调度器已停止")
# 清理bot管理器资源
# if bot_manager:
# logger.info("正在清理bot管理器资源...")
# # 如果有其他需要清理的资源,可以在这里添加
# # 清理消息分析器资源
# if message_analyzer:
# logger.info("正在清理消息分析器资源...")
# # 如果有其他需要清理的资源,可以在这里添加
# # 清理报告生成器资源
# if report_generator:
# logger.info("正在清理报告生成器资源...")
# # 如果有其他需要清理的资源,可以在这里添加
# 重置全局变量
auto_scheduler = None
bot_manager = None
message_analyzer = None
report_generator = None
config_manager = None
logger.info("QQ群日常分析插件资源清理完成")
except Exception as e:
logger.error(f"重新加载配置失败: {e}")
logger.error(f"插件资源清理失败: {e}")
@filter.command("群分析")
@filter.permission_type(PermissionType.ADMIN)
+3 -3
View File
@@ -253,7 +253,7 @@ class LLMAnalyzer:
else:
logger.warning(f"话题分析响应中未找到JSON格式,响应内容: {result_text[:200]}...")
except json.JSONDecodeError as e:
logger.error(f"话题分析JSON解析失败: {e}")
logger.debug(f"话题分析JSON解析失败: {e}")
logger.debug(f"修复后的JSON: {json_text if 'json_text' in locals() else 'N/A'}")
logger.debug(f"原始响应: {result_text}")
@@ -467,7 +467,7 @@ class LLMAnalyzer:
titles_data = json.loads(json_match.group())
return [UserTitle(**title) for title in titles_data], token_usage
except Exception as e:
logger.error(f"用户称号分析JSON解析失败: {e}")
logger.debug(f"用户称号分析JSON解析失败: {e}")
logger.debug(f"原始响应: {result_text}")
return [], token_usage
@@ -570,7 +570,7 @@ class LLMAnalyzer:
quotes_data = json.loads(json_match.group())
return [GoldenQuote(**quote) for quote in quotes_data[:max_golden_quotes]], token_usage
except Exception as e:
logger.error(f"金句分析JSON解析失败: {e}")
logger.debug(f"金句分析JSON解析失败: {e}")
logger.debug(f"原始响应: {result_text}")
return [], token_usage
+107 -49
View File
@@ -114,64 +114,122 @@ class AutoScheduler:
await asyncio.sleep(300)
async def _run_auto_analysis(self):
"""执行自动分析"""
"""执行自动分析 - 并发处理所有群聊"""
try:
logger.info("开始执行自动群聊分析")
logger.info("开始执行自动群聊分析(并发模式)")
# 为每个启用的群执行分析
enabled_groups = self.config_manager.get_enabled_groups()
if not enabled_groups:
logger.info("没有启用的群聊需要分析")
return
logger.info(f"将为 {len(enabled_groups)} 个群聊并发执行分析: {enabled_groups}")
# 创建并发任务 - 为每个群聊创建独立的分析任务
analysis_tasks = []
for group_id in enabled_groups:
try:
logger.info(f"为群 {group_id} 执行自动分析")
await self._perform_auto_analysis_for_group(group_id)
except Exception as e:
logger.error(f"{group_id} 自动分析失败: {e}")
task = asyncio.create_task(
self._perform_auto_analysis_for_group_with_timeout(group_id),
name=f"analysis_group_{group_id}"
)
analysis_tasks.append(task)
# 并发执行所有分析任务,使用 return_exceptions=True 确保单个任务失败不影响其他任务
results = await asyncio.gather(*analysis_tasks, return_exceptions=True)
# 统计执行结果
success_count = 0
error_count = 0
for i, result in enumerate(results):
group_id = enabled_groups[i]
if isinstance(result, Exception):
logger.error(f"{group_id} 分析任务异常: {result}")
error_count += 1
else:
success_count += 1
logger.info(f"并发分析完成 - 成功: {success_count}, 失败: {error_count}, 总计: {len(enabled_groups)}")
except Exception as e:
logger.error(f"自动分析执行失败: {e}")
logger.error(f"自动分析执行失败: {e}", exc_info=True)
async def _perform_auto_analysis_for_group(self, group_id: str):
"""为指定群执行自动分析"""
async def _perform_auto_analysis_for_group_with_timeout(self, group_id: str):
"""为指定群执行自动分析(带超时控制)"""
try:
# 检查bot管理器状态
if not self.bot_manager.is_ready_for_auto_analysis():
status = self.bot_manager.get_status_info()
logger.warning(f"{group_id} 自动分析跳过:bot管理器未就绪 - {status}")
return
logger.info(f"开始为群 {group_id} 执行自动分析")
# 获取群聊消息
analysis_days = self.config_manager.get_analysis_days()
bot_instance = self.bot_manager.get_bot_instance()
messages = await self.message_handler.fetch_group_messages(bot_instance, group_id, analysis_days)
if not messages:
logger.warning(f"{group_id} 未获取到足够的消息记录")
return
# 检查消息数量
min_threshold = self.config_manager.get_min_messages_threshold()
if len(messages) < min_threshold:
logger.warning(f"{group_id} 消息数量不足({len(messages)}条),跳过分析")
return
logger.info(f"{group_id} 获取到 {len(messages)} 条消息,开始分析")
# 进行分析 - 构造正确的 unified_msg_origin
platform_id = self._get_platform_id()
umo = f"{platform_id}:group:{group_id}" if platform_id else None
analysis_result = await self.analyzer.analyze_messages(messages, group_id, umo)
if not analysis_result:
logger.error(f"{group_id} 分析失败")
return
# 生成并发送报告
await self._send_analysis_report(group_id, analysis_result)
# 为每个群聊设置独立的超时时间(20分钟)- 使用 asyncio.wait_for 兼容所有 Python 版本
await asyncio.wait_for(self._perform_auto_analysis_for_group(group_id), timeout=1200)
except asyncio.TimeoutError:
logger.error(f"{group_id} 分析超时(20分钟),跳过该群分析")
except Exception as e:
logger.error(f"{group_id} 自动分析执行失败: {e}", exc_info=True)
logger.error(f"{group_id} 分析任务执行失败: {e}")
async def _perform_auto_analysis_for_group(self, group_id: str):
"""为指定群执行自动分析(核心逻辑)"""
# 为每个群聊使用独立的锁,避免全局锁导致串行化
group_lock_key = f"analysis_{group_id}"
if not hasattr(self, '_group_locks'):
self._group_locks = {}
if group_lock_key not in self._group_locks:
self._group_locks[group_lock_key] = asyncio.Lock()
async with self._group_locks[group_lock_key]:
try:
start_time = asyncio.get_event_loop().time()
# 检查bot管理器状态
if not self.bot_manager.is_ready_for_auto_analysis():
status = self.bot_manager.get_status_info()
logger.warning(f"{group_id} 自动分析跳过:bot管理器未就绪 - {status}")
return
logger.info(f"开始为群 {group_id} 执行自动分析(并发任务)")
# 获取群聊消息
analysis_days = self.config_manager.get_analysis_days()
bot_instance = self.bot_manager.get_bot_instance()
messages = await self.message_handler.fetch_group_messages(bot_instance, group_id, analysis_days)
if not messages:
logger.warning(f"{group_id} 未获取到足够的消息记录")
return
# 检查消息数量
min_threshold = self.config_manager.get_min_messages_threshold()
if len(messages) < min_threshold:
logger.warning(f"{group_id} 消息数量不足({len(messages)}条),跳过分析")
return
logger.info(f"{group_id} 获取到 {len(messages)} 条消息,开始分析")
# 进行分析 - 构造正确的 unified_msg_origin
platform_id = self._get_platform_id()
umo = f"{platform_id}:GroupMessage:{group_id}" if platform_id else None
analysis_result = await self.analyzer.analyze_messages(messages, group_id, umo)
if not analysis_result:
logger.error(f"{group_id} 分析失败")
return
# 生成并发送报告
await self._send_analysis_report(group_id, analysis_result)
# 记录执行时间
end_time = asyncio.get_event_loop().time()
execution_time = end_time - start_time
logger.info(f"{group_id} 分析完成,耗时: {execution_time:.2f}")
except Exception as e:
logger.error(f"{group_id} 自动分析执行失败: {e}", exc_info=True)
finally:
# 清理群聊锁资源(可选,防止内存泄漏)
if hasattr(self, '_group_locks') and len(self._group_locks) > 50:
old_locks = list(self._group_locks.keys())[:10]
for lock_key in old_locks:
if not self._group_locks[lock_key].locked():
del self._group_locks[lock_key]
async def _send_analysis_report(self, group_id: str, analysis_result: dict):
"""发送分析报告到群"""