""" 数据库连接配置 针对 Apache Doris 优化的数据库操作模块 """ import os import time import pymysql from contextlib import contextmanager from lib.logger import log_info, log_error, log_warning from lib.doris_optimize import get_pool, get_doris_connection, execute_doris_query, execute_doris_update, execute_doris_batch_insert def get_db_config(): """从环境变量或默认值获取数据库配置""" return { 'host': os.environ.get('DB_HOST', 'haoslm2.xicp.net'), 'port': int(os.environ.get('DB_PORT', 10216)), 'user': os.environ.get('DB_USER', 'root'), 'password': os.environ.get('DB_PASSWORD', 'DsideaL147258369'), 'database': os.environ.get('DB_NAME', 'yltcharge'), 'charset': 'utf8mb4', 'cursorclass': pymysql.cursors.DictCursor, 'autocommit': True, 'connect_timeout': 10, # 优化:缩短连接超时 'read_timeout': 120, 'write_timeout': 120 } DB_CONFIG = get_db_config() @contextmanager def get_connection(): """获取数据库连接(上下文管理器)- 优先使用连接池""" try: # 使用连接池获取连接 with get_doris_connection() as conn: yield conn except Exception: # 连接池失败时回退到直接连接 log_warning('[数据库] 连接池获取失败,回退到直接连接', 'db') try: conn = pymysql.connect(**DB_CONFIG) yield conn except pymysql.Error as e: log_error(f'数据库连接失败: {e}', 'db', exc_info=True) raise finally: try: conn.close() except: pass def execute_query(sql, params=None, retry=2): """执行查询并返回结果 - 使用 Doris 优化版本""" return execute_doris_query(sql, params, retry) def execute_update(sql, params=None, retry=2): """执行更新操作 - 使用 Doris 优化版本""" return execute_doris_update(sql, params, retry) def execute_insert(sql, params=None, retry=2): """执行插入操作并返回插入ID""" for attempt in range(retry + 1): try: start_time = time.time() with get_doris_connection() as conn: with conn.cursor() as cursor: cursor.execute(sql, params) last_id = cursor.lastrowid elapsed = (time.time() - start_time) * 1000 if elapsed > 100: log_info(f"[Doris慢插入] 耗时: {elapsed:.1f}ms, SQL: {sql[:200]}", 'doris') if attempt > 0: log_info(f'[数据库] 插入重试成功,第{attempt+1}次尝试', 'db') return last_id except pymysql.err.OperationalError as e: if attempt < retry and (e.args[0] == 2013 or e.args[0] == 2006): log_warning(f'[数据库] 插入连接断开,正在重试(第{attempt+1}次): {e}', 'db') time.sleep(0.5 * (attempt + 1)) continue log_error(f'插入失败: {e}\nSQL: {sql}\nParams: {params}', 'db', exc_info=True) raise except Exception as e: log_error(f'插入失败: {e}\nSQL: {sql}\nParams: {params}', 'db', exc_info=True) raise def execute_batch_insert(sql, params_list, batch_size=100): """ 批量插入优化(针对 Doris 特性优化) Args: sql: SQL 语句 params_list: 参数列表 batch_size: 每批大小(默认100) Returns: int: 插入的总行数 """ return execute_doris_batch_insert(sql, params_list, batch_size)