async_mysql_client.py 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170
  1. """
  2. AsyncMySQLClient (基于 asyncmy) + 项目统一日志正式上线版
  3. - 自动管理连接池,避免重复初始化和内存泄漏
  4. - 支持 async with 上下文管理
  5. - 高并发协程安全
  6. """
  7. import traceback
  8. from typing import List, Dict, Any, Optional
  9. import asyncio
  10. import asyncmy
  11. from core.utils.log.logger_manager import LoggerManager
  12. class AsyncMySQLClient:
  13. def __init__(self,
  14. host: str,
  15. port: int,
  16. user: str,
  17. password: str,
  18. db: str,
  19. charset: str,
  20. minsize: int = 1,
  21. maxsize: int = 5,
  22. pool_recycle: int = 3600,
  23. logger: Optional[LoggerManager] = None,
  24. aliyun_logr:Optional[LoggerManager] = None,):
  25. """
  26. 初始化配置,延迟创建连接池
  27. """
  28. self._db_settings = {
  29. "host": host,
  30. "port": port,
  31. "user": user,
  32. "password": password,
  33. "db": db,
  34. "autocommit": True,
  35. "charset": charset,
  36. "connect_timeout": 5,
  37. }
  38. self._minsize = minsize
  39. self._maxsize = maxsize
  40. self._pool_recycle = pool_recycle
  41. self._pool: Optional[asyncmy.Pool] = None
  42. self._lock = asyncio.Lock() # 防止并发初始化
  43. self.logger = logger
  44. self.aliyun_logger = aliyun_logr
  45. async def __aenter__(self):
  46. """支持 async with 自动初始化连接池"""
  47. await self.init_pool()
  48. return self
  49. async def __aexit__(self, exc_type, exc_val, exc_tb):
  50. """支持 async with 自动关闭连接池"""
  51. await self.close()
  52. async def init_pool(self):
  53. """
  54. 初始化连接池(懒加载 + 并发锁保护)
  55. """
  56. if self._pool:
  57. return
  58. async with self._lock:
  59. if self._pool:
  60. return
  61. try:
  62. self._pool = await asyncmy.create_pool(
  63. **self._db_settings,
  64. minsize=self._minsize,
  65. maxsize=self._maxsize,
  66. pool_recycle=self._pool_recycle,
  67. )
  68. except Exception as e:
  69. msg = f"[AsyncMySQLClient] 连接池初始化失败: {e} \n {traceback.format_exc()}"
  70. if self.logger:
  71. self.logger.error(msg)
  72. if self.aliyun_logger:
  73. self.aliyun_logger.logging(code="9001", message=msg)
  74. raise
  75. async def close(self):
  76. """
  77. 关闭连接池
  78. """
  79. if self._pool:
  80. self._pool.close()
  81. await self._pool.wait_closed()
  82. self._pool = None
  83. async def fetch_all(self, sql: str, params: Optional[List[Any]] = None) -> List[Dict[str, Any]]:
  84. """
  85. 查询多行,返回字典列表
  86. """
  87. await self.init_pool()
  88. try:
  89. async with self._pool.acquire() as conn:
  90. async with conn.cursor() as cur:
  91. await cur.execute(sql, params or [])
  92. rows = await cur.fetchall()
  93. columns = [desc[0] for desc in cur.description]
  94. result = [dict(zip(columns, row)) for row in rows]
  95. return result
  96. except Exception as e:
  97. msg = f"[AsyncMySQLClient] fetch_all 执行失败: {e} | SQL: {sql}"
  98. if self.logger:
  99. self.logger.error(msg)
  100. if self.aliyun_logger:
  101. self.aliyun_logger.logging(code="9002", message=msg)
  102. raise
  103. async def fetch_one(self, sql: str, params: Optional[List[Any]] = None) -> Optional[Dict[str, Any]]:
  104. """
  105. 查询单行,返回字典
  106. """
  107. await self.init_pool()
  108. try:
  109. async with self._pool.acquire() as conn:
  110. async with conn.cursor() as cur:
  111. await cur.execute(sql, params or [])
  112. row = await cur.fetchone()
  113. if row is None:
  114. return None
  115. columns = [desc[0] for desc in cur.description]
  116. result = dict(zip(columns, row))
  117. return result
  118. except Exception as e:
  119. msg = f"[AsyncMySQLClient] fetch_one 执行失败: {e} | SQL: {sql}"
  120. if self.logger:
  121. self.logger.error(msg)
  122. if self.aliyun_logger:
  123. self.aliyun_logger.logging(code="9003", message=msg)
  124. raise
  125. async def execute(self, sql: str, params: Optional[List[Any]] = None) -> int:
  126. """
  127. 执行单条写操作
  128. """
  129. await self.init_pool()
  130. try:
  131. async with self._pool.acquire() as conn:
  132. async with conn.cursor() as cur:
  133. await cur.execute(sql, params or [])
  134. return cur.rowcount
  135. except Exception as e:
  136. msg = f"[AsyncMySQLClient] execute 执行失败: {e} | SQL: {sql}"
  137. if self.logger:
  138. self.logger.error(msg)
  139. if self.aliyun_logger:
  140. self.aliyun_logger.logging(code="9004", message=msg)
  141. raise
  142. async def executemany(self, sql: str, params_list: List[List[Any]]) -> int:
  143. """
  144. 批量执行写操作
  145. """
  146. await self.init_pool()
  147. try:
  148. async with self._pool.acquire() as conn:
  149. async with conn.cursor() as cur:
  150. await cur.executemany(sql, params_list)
  151. return cur.rowcount
  152. except Exception as e:
  153. msg = f"[AsyncMySQLClient] executemany 执行失败: {e} | SQL: {sql}"
  154. if self.logger:
  155. self.logger.error(msg)
  156. if self.aliyun_logger:
  157. self.aliyun_logger.logging(code="9005", message=msg)
  158. raise