多线程实现消息推送并可重试3次以及1小时后重试
zhezhongyun 2025-05-05 20:10 48 浏览
# -*- coding: utf-8 -*-
"""
Created on Tue Apr 22 09:05:46 2025
@author: 1
"""
import requests
import time
import schedule
import sched
from datetime import datetime, timedelta
import threading
import pymysql
from concurrent.futures import ThreadPoolExecutor
import logging
from queue import Queue
import os
from pathlib import Path
# 获取虚拟环境目录
venv_dir = os.environ.get('VIRTUAL_ENV', None)
if venv_dir:
log_path = Path(venv_dir) / "app.log"
else:
log_path = Path.home() / "app.log"
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s [%(levelname)s] %(message)s',
handlers=[
logging.FileHandler(log_path),
logging.StreamHandler()
]
)
# 信息等级定义
INFO_LEVELS = {
'紧急': {'retry_interval': 60, 'max_retries': 3},
'重要': {'retry_interval': 60, 'max_retries': 3},
'一般': {'retry_interval': 60, 'max_retries': 3}
}
# 数据库连接配置
DB_CONFIG = {
'host': '************',
'user': 'root',
'password': '*******',
'database': '***********',
'charset': 'utf8mb4',
'cursorclass': pymysql.cursors.DictCursor
}
# 创建线程池
executor = ThreadPoolExecutor(max_workers=20)
class MessageSender:
def __init__(self):
self.message_queue = Queue()
self.lock = threading.RLock() # 可重入锁
self.scheduled_messages = {}
def send_message(self, group, message, info_level, attempt=1):
"""发送消息到企业微信机器人"""
headers = {'Content-Type': 'application/json'}
payload = {
"msgtype": "text",
"text": {
"content": message
}
}
try:
response = requests.post(group['robot_webhook'], json=payload, headers=headers, timeout=10)
if response.status_code == 200 and response.json().get('errcode') == 0:
logging.info(f"成功发送到群 {group['group_name']}")
return True
else:
logging.warning(f"发送到群 {group['group_name']} 失败,状态码: {response.status_code}")
return False
except Exception as e:
logging.error(f"发送到群 {group['group_name']} 异常: {e}")
return False
def retry_send(self, group, message, info_level, upno, sendplan, plantime, attempt=1):
"""独立的重试逻辑,不影响其他消息"""
max_retries = INFO_LEVELS[info_level]['max_retries']
retry_interval = INFO_LEVELS[info_level]['retry_interval']
final_retry_delay = 3600 # 1小时(秒)
success = self.send_message(group, message, info_level, attempt)
if not success:
if attempt <= max_retries:
# 前3次快速重试
logging.info(f"将在 {retry_interval} 秒后重试 (尝试 {attempt}/{max_retries})")
time.sleep(retry_interval)
self.retry_send(group, message, info_level, upno, sendplan, plantime, attempt + 1)
else:
# 3次失败后,1小时后再试
logging.warning(f"3次重试失败,将在1小时后再次尝试")
time.sleep(final_retry_delay)
# 重置尝试次数
self.retry_send(group, message, info_level, upno, sendplan, plantime, attempt=1)
def send_to_superior_groups(self, group, message, info_level, upno, sendplan, plantime, attempt=1):
"""向上级群组发送消息,每个上级群组独立处理"""
current_group = group
while current_group['parent_group_id'] is not None:
print('retry5')
superior = self.find_group_by_id(current_group['parent_group_id'])
if superior:
# 为每个上级群组创建独立线程
print('retry6')
executor.submit(self.retry_send, superior, message, info_level, upno, sendplan, plantime)
if attempt == 2:
break
print('retry7')
attempt = attempt + 1
if upno == 1:
break
print('retry8')
current_group = superior
else:
break
def find_group_by_id(self, group_id):
"""根据ID查找群组信息"""
try:
connection = pymysql.connect(**DB_CONFIG) # 不加锁
with connection.cursor() as cursor:
with self.lock: # 只锁住查询部分
cursor.execute(
"SELECT id, group_name, level, robot_webhook, parent_group_id FROM wechat_groups WHERE id = %s",
(group_id,)
)
return cursor.fetchone()
except Exception as e:
logging.error(f"查找群组时出错: {e}")
return None
finally:
if 'connection' in locals() and connection:
connection.close() # 确保连接关闭
def find_group_by_name(self, group_name):
"""根据名称查找群组信息"""
try:
with self.lock:
connection = pymysql.connect(**DB_CONFIG)
with connection.cursor() as cursor:
sql = "SELECT id, group_name, level, robot_webhook, parent_group_id FROM wechat_groups WHERE group_name = %s"
cursor.execute(sql, (group_name,))
print('按群名查找,',group_name)
result = cursor.fetchone()
print('执行命令,')
print(f'查询结果: {result}')
print(1111)
if not result:
logging.warning(f"未找到群组: {group_name}")
return result
except Exception as e:
print(f"查找群组时出错: {e}")
logging.error(f"查找群组时出错: {e}")
return None
finally:
if 'connection' in locals() and connection:
connection.close()
def schedule_message(self, row):
"""根据plantime安排消息发送"""
try:
plantime = row['plantime']
sendplan = row['sendplan'] # 1即时发送 2 定时发送
print('sendplan,',sendplan)
target_group2 = row['role']
print('目前消息群,',target_group2)
if sendplan == 1:
print('如果sendplan=1,立即发送')
# 如果sendplan=1,立即发送
self.process_message(row)
return
print('如果sendplan不是1,定时发送')
# 将plantime转换为datetime对象
if isinstance(plantime, str):
send_time = datetime.strptime(plantime, '%Y-%m-%d %H:%M:%S')
else:
send_time = plantime
now = datetime.now()
if send_time <= now:
# 如果发送时间已过,立即发送
print('如果发送时间已过,立即发送,',target_group2)
target_group = self.find_group_by_name(target_group2)
print('hhh')
if not target_group:
print('jjj')
logging.warning(f" 未找到群组 '{row['role']}',跳过此消息")
return # 直接返回,不继续执行
self.process_message(row)
else:
# 计算延迟时间(秒)
delay = (send_time - now).total_seconds()
# 为消息创建定时任务
message_id = row['id']
print('为消息创建定时任务,',message_id)
if message_id not in self.scheduled_messages:
timer = threading.Timer(delay, self.process_message, args=(row,))
timer.start()
print('已安排')
self.scheduled_messages[message_id] = timer
logging.info(f"已安排消息 {message_id} 在 {send_time} 发送")
except Exception as e:
logging.error(f"安排消息发送时出错: {e}")
def process_message(self, row):
"""处理单条消息"""
try:
message = f"通知: {row['result']}"
target_group = self.find_group_by_name(row['role'])
if not target_group:
logging.warning(f" 未找到群组 '{row['role']}',跳过此消息")
return # 直接返回,不继续执行
upno = row['upno']
sendplan = row['sendplan']
plantime = row['plantime']
# 主发送
self.retry_send(target_group, message, row['levelname'], upno, sendplan, plantime)
# 向上级发送
self.send_to_superior_groups(target_group, message, row['levelname'], upno, sendplan, plantime)
# 更新数据库状态
self.update_message_status(row['id'])
# 从预定消息中移除
if row['id'] in self.scheduled_messages:
del self.scheduled_messages[row['id']]
except Exception as e:
logging.error(f" 处理消息时出错: {e}")
def update_message_status(self, message_id):
"""更新消息状态"""
try:
with self.lock:
connection = pymysql.connect(**DB_CONFIG)
with connection.cursor() as cursor:
update_sql = "UPDATE ai_dayairesult SET isend = 1 WHERE id = %s"
cursor.execute(update_sql, (message_id,))
connection.commit()
except Exception as e:
logging.error(f"更新消息状态时出错: {e}")
finally:
if 'connection' in locals() and connection:
connection.close()
def check_and_schedule_messages(self):
"""检查并安排消息发送"""
try:
with self.lock:
connection = pymysql.connect(**DB_CONFIG)
with connection.cursor() as cursor:
sql = """SELECT id, sqlcmd, result, createdate, role, details, levelname, upno, sendplan, plantime
FROM ai_dayairesult
WHERE isend = 0"""
cursor.execute(sql)
results = cursor.fetchall()
for row in results:
try:
self.schedule_message(row) # 即使某条失败,继续下一条
except Exception as e:
logging.error(f" 安排消息 {row['id']} 时出错: {e}")
except Exception as e:
logging.error(f" 检查并安排消息时出错: {e}")
finally:
if 'connection' in locals() and connection:
connection.close()
def run_scheduler(sender):
"""运行定时任务"""
# 每分钟检查一次待发送消息
schedule.every(1).minutes.do(sender.check_and_schedule_messages)
while True:
schedule.run_pending()
time.sleep(1)
if __name__ == "__main__":
sender = MessageSender()
# 启动调度器线程
scheduler_thread = threading.Thread(target=run_scheduler, args=(sender,))
scheduler_thread.daemon = True
scheduler_thread.start()
# 立即执行一次检查
sender.check_and_schedule_messages()
# 主线程保持运行
try:
while True:
time.sleep(1)
except KeyboardInterrupt:
logging.info("程序退出")
# 取消所有预定但未发送的消息
for timer in sender.scheduled_messages.values():
timer.cancel()相关推荐
- Python入门学习记录之一:变量_python怎么用变量
-
写这个,主要是对自己学习python知识的一个总结,也是加深自己的印象。变量(英文:variable),也叫标识符。在python中,变量的命名规则有以下三点:>变量名只能包含字母、数字和下划线...
- python变量命名规则——来自小白的总结
-
python是一个动态编译类编程语言,所以程序在运行前不需要如C语言的先行编译动作,因此也只有在程序运行过程中才能发现程序的问题。基于此,python的变量就有一定的命名规范。python作为当前热门...
- Python入门学习教程:第 2 章 变量与数据类型
-
2.1什么是变量?在编程中,变量就像一个存放数据的容器,它可以存储各种信息,并且这些信息可以被读取和修改。想象一下,变量就如同我们生活中的盒子,你可以把东西放进去,也可以随时拿出来看看,甚至可以换成...
- 绘制学术论文中的“三线表”具体指导
-
在科研过程中,大家用到最多的可能就是“三线表”。“三线表”,一般主要由三条横线构成,当然在变量名栏里也可以拆分单元格,出现更多的线。更重要的是,“三线表”也是一种数据记录规范,以“三线表”形式记录的数...
- Python基础语法知识--变量和数据类型
-
学习Python中的变量和数据类型至关重要,因为它们构成了Python编程的基石。以下是帮助您了解Python中的变量和数据类型的分步指南:1.变量:变量在Python中用于存储数据值。它们充...
- 一文搞懂 Python 中的所有标点符号
-
反引号`无任何作用。传说Python3中它被移除是因为和单引号字符'太相似。波浪号~(按位取反符号)~被称为取反或补码运算符。它放在我们想要取反的对象前面。如果放在一个整数n...
- Python变量类型和运算符_python中变量的含义
-
别再被小名词坑哭了:Python新手常犯的那些隐蔽错误,我用同事的真实bug拆给你看我记得有一次和同事张姐一起追查一个看似随机崩溃的脚本,最后发现罪魁祸首竟然是她把变量命名成了list。说实话...
- 从零开始:深入剖析 Spring Boot3 中配置文件的加载顺序
-
在当今的互联网软件开发领域,SpringBoot无疑是最为热门和广泛应用的框架之一。它以其强大的功能、便捷的开发体验,极大地提升了开发效率,成为众多开发者构建Web应用程序的首选。而在Spr...
- Python中下划线 ‘_’ 的用法,你知道几种
-
Python中下划线()是一个有特殊含义和用途的符号,它可以用来表示以下几种情况:1在解释器中,下划线(_)表示上一个表达式的值,可以用来进行快速计算或测试。例如:>>>2+...
- 解锁Shell编程:变量_shell $变量
-
引言:开启Shell编程大门Shell作为用户与Linux内核之间的桥梁,为我们提供了强大的命令行交互方式。它不仅能执行简单的文件操作、进程管理,还能通过编写脚本实现复杂的自动化任务。无论是...
- 一文学会Python的变量命名规则!_python的变量命名有哪些要求
-
目录1.变量的命名原则3.内置函数尽量不要做变量4.删除变量和垃圾回收机制5.结语1.变量的命名原则①由英文字母、_(下划线)、或中文开头②变量名称只能由英文字母、数字、下画线或中文字所组成。③英文字...
- 更可靠的Rust-语法篇-区分语句/表达式,略览if/loop/while/for
-
src/main.rs://函数定义fnadd(a:i32,b:i32)->i32{a+b//末尾表达式}fnmain(){leta:i3...
- C++第五课:变量的命名规则_c++中变量的命名规则
-
变量的命名不是想怎么起就怎么起的,而是有一套固定的规则的。具体规则:1.名字要合法:变量名必须是由字母、数字或下划线组成。例如:a,a1,a_1。2.开头不能是数字。例如:可以a1,但不能起1a。3....
- Rust编程-核心篇-不安全编程_rust安全性
-
Unsafe的必要性Rust的所有权系统和类型系统为我们提供了强大的安全保障,但在某些情况下,我们需要突破这些限制来:与C代码交互实现底层系统编程优化性能关键代码实现某些编译器无法验证的安全操作Rus...
- 探秘 Python 内存管理:背后的神奇机制
-
在编程的世界里,内存管理就如同幕后的精密操控者,确保程序的高效运行。Python作为一种广泛使用的编程语言,其内存管理机制既巧妙又复杂,为开发者们提供了便利的同时,也展现了强大的底层控制能力。一、P...
- 一周热门
- 最近发表
- 标签列表
-
- HTML 教程 (33)
- HTML 简介 (35)
- HTML 实例/测验 (32)
- HTML 测验 (32)
- JavaScript 和 HTML DOM 参考手册 (32)
- HTML 拓展阅读 (30)
- HTML文本框样式 (31)
- HTML滚动条样式 (34)
- HTML5 浏览器支持 (33)
- HTML5 新元素 (33)
- HTML5 WebSocket (30)
- HTML5 代码规范 (32)
- HTML5 标签 (717)
- HTML5 标签 (已废弃) (75)
- HTML5电子书 (32)
- HTML5开发工具 (34)
- HTML5小游戏源码 (34)
- HTML5模板下载 (30)
- HTTP 状态消息 (33)
- HTTP 方法:GET 对比 POST (33)
- 键盘快捷键 (35)
- 标签 (226)
- opacity 属性 (32)
- transition 属性 (33)
- 1-1. 变量声明 (31)
