04_SQLite 数据持久化与 DAO 模式实践:构建高性能本地数据层
SQLite 是桌面应用的首选数据库,但直接使用容易陷入性能陷阱和并发问题。本文基于 GPFX 项目,系统讲解 SQLite 在 Python 桌面应用中的最佳实践,包括连接管理、事务处理、DAO 模式、批量操作优化等核心技术。
一、为什么选择 SQLite?
对于桌面量化分析系统,SQLite 具有独特优势:
| 特性 | SQLite | MySQL/PostgreSQL |
|---|---|---|
| 部署 | 零配置,单文件 | 需安装服务端 |
| 性能 | 本地读写极快 | 网络开销 |
| 并发 | 单写多读 | 高并发 |
| 事务 | 完整 ACID | 完整 ACID |
| 适用场景 | 桌面应用、嵌入式 | 服务端、分布式 |
关键指标:
- 5000 只股票 × 250 交易日/年 × 10 年 = 1250 万条日线数据
- SQLite 本地查询毫秒级响应
- 单文件便于备份和迁移
二、数据库连接管理
2.1 连接池与线程安全
SQLite 连接不能跨线程使用,需要精心设计:
importsqlite3importthreadingfromcontextlibimportcontextmanagerfrompathlibimportPathclassDatabase:"""SQLite 数据库连接管理器"""def__init__(self,db_path:str|Path):self._db_path=str(db_path)self._conn:sqlite3.Connection|None=Noneself._lock=threading.RLock()# 可重入锁,支持嵌套调用defconnect(self)->sqlite3.Connection:"""获取连接(懒加载 + 双重检查)"""ifself._connisNone:withself._lock:ifself._connisNone:# 确保目录存在Path(self._db_path).parent.mkdir(parents=True,exist_ok=True)# 创建连接self._conn=sqlite3.connect(self._db_path,check_same_thread=False,# 允许跨线程使用(需配合锁)timeout=30# 锁等待超时 30 秒)# 返回字典式结果self._conn.row_factory=sqlite3.Row# 性能优化 PRAGMAself._conn.execute("PRAGMA journal_mode=WAL")# 写前日志self._conn.execute("PRAGMA foreign_keys=ON")# 外键约束self._conn.execute("PRAGMA synchronous=NORMAL")# 平衡安全与性能returnself._conndefclose(self):"""关闭连接"""withself._lock:ifself._conn:self._conn.close()self._conn=None关键配置解析:
journal_mode=WAL:写前日志模式- 读操作不阻塞写操作
- 崩溃恢复更安全
- 性能提升 2-3 倍
synchronous=NORMAL:同步级别FULL:每次事务都刷盘(最安全,最慢)NORMAL:WAL 模式下安全且快速OFF:不保证持久性(最快,风险高)
check_same_thread=False:允许跨线程- 必须配合锁使用
- 由上层保证线程安全
2.2 上下文管理器支持
classDatabase:def__enter__(self):"""支持 with 语句"""self.connect()returnselfdef__exit__(self,exc_type,exc_val,exc_tb):"""自动关闭连接"""self.close()returnFalse# 不吞掉异常# 使用示例withDatabase("data/stocks.db")asdb:db.execute("CREATE TABLE IF NOT EXISTS ...")# 退出时自动关闭2.3 事务管理
@contextmanagerdeftransaction(self):"""事务上下文管理器 自动提交/回滚: - 正常结束:commit - 发生异常:rollback """withself._lock:conn=self.connect()try:yieldconn conn.commit()exceptException:try:conn.rollback()exceptException:pass# 回滚失败也要抛出原异常raise# 使用示例defbatch_insert(self,records:list[dict]):withself.transaction()asconn:conn.executemany("INSERT OR REPLACE INTO daily_none VALUES (?, ?, ...)",records)# 所有插入在一个事务中,要么全成功要么全失败三、DAO 模式实现
3.1 通用 SQL 生成器
避免手写重复 SQL,通过元编程自动生成:
classStockDAO:"""股票数据访问对象"""# 字段定义:与 Tushare API 输出参数一致STOCK_BASIC_FIELDS=["ts_code","symbol","name","area","industry","market","exchange","list_status","list_date",]DAILY_DATA_FIELDS=["ts_code","trade_date","open","high","low","close","pre_close","change","pct_chg","vol","amount",]@staticmethoddef_build_upsert_sql(table:str,fields:list[str],conflict_cols:list[str])->str:"""生成 UPSERT SQL(INSERT OR REPLACE) SQLite 语法: INSERT INTO table (cols) VALUES (vals) ON CONFLICT(conflict_cols) DO UPDATE SET col=excluded.col """col_str=", ".join(fields)param_str=", ".join(f":{f}"forfinfields)# 命名参数# 冲突时要更新的字段(排除主键)update_fields=[fforfinfieldsiffnotinconflict_cols]conflict_str=", ".join(conflict_cols)ifnotupdate_fields:# 没有要更新的字段,冲突时什么都不做return(f"INSERT INTO{table}({col_str}) VALUES ({param_str}) "f"ON CONFLICT({conflict_str}) DO NOTHING")# 冲突时更新所有非主键字段update_str=", ".join(f"{f}=excluded.{f}"forfinupdate_fields)return(f"INSERT INTO{table}({col_str}) VALUES ({param_str}) "f"ON CONFLICT({conflict_str}) DO UPDATE SET{update_str}")为什么用命名参数:field而非??
- 可读性更强,参数顺序不易错
- 支持字典传参
{"ts_code": "000001.SZ", ...} - 便于日志调试
3.2 数据规范化
处理 pandas 的 NaN 与 SQLite 的 NULL 转换:
@staticmethoddef_normalize_records(records:list[dict],fields:list[str])->list[dict]:"""规范化记录 - NaN → None(SQLite NULL) - 确保所有字段存在 """importmath result=[]forrinrecords:row={}forfieldinfields:v=r.get(field)# pandas NaN 是 float 类型,需要特殊处理ifisinstance(v,float)andmath.isnan(v):row[field]=Noneelse:row[field]=v result.append(row)returnresult3.3 表名白名单防注入
动态表名是 SQL 注入的高危点,必须严格校验:
classDatabase:# 合法表名白名单_VALID_TABLES={"users","stock_basic","daily_none","daily_qfq","daily_hfq","adj_factor","trade_cal","selection_runs","selection_results","backtest_runs","backtest_daily","watchlist","intraday_min","index_daily",}def_validate_table(self,name:str)->str:"""校验表名合法性"""ifnamenotinself._VALID_TABLES:raiseValueError(f"非法表名:{name}")returnnamedefquery_table(self,table:str,where:str="",params:tuple=()):"""安全的动态表查询"""table=self._validate_table(table)# 校验表名sql=f"SELECT * FROM{table}"ifwhere:sql+=f" WHERE{where}"# where 子句仍需参数化returnself.query(sql,params)安全原则:
- 表名、列名:白名单校验
- 值:参数化查询(
?或:name) - 绝不拼接用户输入
四、批量操作优化
4.1 批量插入
单条插入与批量插入性能对比:
defupsert_daily_batch(self,records:list[dict],adj_type:str):"""批量插入日线数据 性能对比(5000 条记录): - 单条插入:30 秒 - 批量插入:0.5 秒(60 倍提升) """table=self._daily_table(adj_type)fields=self.DAILY_DATA_FIELDS+["updated_at"]# 生成 UPSERT SQLsql=self._build_upsert_sql(table,fields,["ts_code","trade_date"])# 规范化数据(NaN → None)normalized=self._normalize_records(records,fields)# 添加更新时间戳fromdatetimeimportdatetime now=datetime.now().strftime("%Y-%m-%d %H:%M:%S")forrinnormalized:r["updated_at"]=now# 批量执行(关键:使用事务)withself._db.transaction()asconn:conn.executemany(sql,normalized)性能优化要点:
- executemany vs execute:减少 SQL 解析开销
- 事务包裹:5000 次提交合并为 1 次
- 预处理语句:避免重复解析 SQL
4.2 批量查询
避免 N+1 查询问题:
defget_daily_batch(self,ts_codes:list[str],start_date:str,end_date:str)->pd.DataFrame:"""批量查询多只股票日线数据 反模式(N+1 查询): for code in ts_codes: df = get_daily(code, ...) # 5000 次查询 正确做法: 一次查询返回所有数据 """# 构建 IN 子句的占位符placeholders=",".join("?"*len(ts_codes))sql=f""" SELECT * FROM daily_qfq WHERE ts_code IN ({placeholders}) AND trade_date BETWEEN ? AND ? ORDER BY ts_code, trade_date """params=(*ts_codes,start_date,end_date)rows=self._db.query(sql,params)returnpd.DataFrame([dict(r)forrinrows])五、索引设计
5.1 高频查询索引
根据查询模式设计索引:
-- 日线表:最常按股票代码和日期范围查询CREATEINDEXidx_daily_none_codeONdaily_none(ts_code);CREATEINDEXidx_daily_none_dateONdaily_none(trade_date);-- 联合索引:同时按代码和日期查询(覆盖索引)CREATEINDEXidx_daily_none_code_dateONdaily_none(ts_code,trade_date);-- 选股结果:按运行ID查询明细CREATEINDEXidx_sel_results_runONselection_results(run_id);-- 回测净值:按回测ID和日期查询CREATEINDEXidx_bt_daily_btONbacktest_daily(backtest_id,trade_date);5.2 索引使用验证
使用EXPLAIN QUERY PLAN验证索引是否生效:
defexplain_query(self,sql:str,params:tuple=()):"""分析查询计划"""rows=self.query(f"EXPLAIN QUERY PLAN{sql}",params)forrowinrows:print(row["detail"])# 测试dao.explain_query("SELECT * FROM daily_none WHERE ts_code = ? AND trade_date > ?",("000001.SZ","20240101"))# 输出:SEARCH daily_none USING INDEX idx_daily_none_code_date (ts_code=? AND trade_date>?)六、数据迁移与版本管理
6.1 表结构演进
支持平滑的表结构升级:
def_create_stock_basic_table(self):"""创建/升级 stock_basic 表"""# 检查表是否存在cols=self.query("PRAGMA table_info(stock_basic)")ifnotcols:# 表不存在,创建新表col_defs=["ts_code TEXT PRIMARY KEY"]forname,col_typeinself._STOCK_BASIC_COLUMNS.items():col_defs.append(f"{name}{col_type}")self.execute(f"CREATE TABLE stock_basic ({', '.join(col_defs)})")return# 表已存在,检查是否需要添加新列existing={c["name"]forcincols}forname,col_typeinself._STOCK_BASIC_COLUMNS.items():ifnamenotinexisting:self.execute(f"ALTER TABLE stock_basic ADD COLUMN{name}{col_type}")6.2 数据迁移示例
从旧表迁移到新表结构:
def_migrate_daily_tables(self):"""从旧的 daily_data 表迁移到分表结构 旧结构:单表 daily_data,用 adj_type 字段区分 新结构:daily_none / daily_qfq / daily_hfq 三张表 """# 检查旧表是否存在old_cols=self.query("PRAGMA table_info(daily_data)")ifnotold_cols:# 旧表不存在,直接创建新表forsuffixin("none","qfq","hfq"):self._create_one_daily_table(suffix)return# 旧表存在,执行数据迁移withself.transaction()asconn:forsuffixin("none","qfq","hfq"):# 创建新表self._create_one_daily_table(suffix)# 迁移数据conn.execute(f""" INSERT OR IGNORE INTO daily_{suffix}SELECT * FROM daily_data WHERE adj_type = ? """,(suffix,))# 删除旧表conn.execute("DROP TABLE daily_data")迁移原则:
- 使用
INSERT OR IGNORE避免重复 - 在事务中执行,失败可回滚
- 保留旧表备份(可选)
七、并发控制
7.1 读写分离
SQLite WAL 模式支持单写多读:
classDatabase:def__init__(self,db_path:str):self._db_path=db_path self._write_lock=threading.RLock()# 写锁self._read_lock=threading.Lock()# 读锁(可选)defexecute(self,sql:str,params:tuple=()):"""写操作:独占锁"""withself._write_lock:conn=self.connect()try:cur=conn.execute(sql,params)conn.commit()returncurexceptException:conn.rollback()raisedefquery(self,sql:str,params:tuple=())->list:"""读操作:WAL 模式下无需加锁"""# WAL 模式下读操作不会被写操作阻塞conn=self.connect()returnconn.execute(sql,params).fetchall()7.2 线程局部连接
多线程环境下的连接管理:
classThreadLocalDatabase:"""线程局部数据库连接"""def__init__(self,db_path:str):self._db_path=db_path self._local=threading.local()defget_connection(self)->sqlite3.Connection:"""获取当前线程的连接"""ifnothasattr(self._local,"conn"):self._local.conn=sqlite3.connect(self._db_path,check_same_thread=False)self._local.conn.row_factory=sqlite3.Rowreturnself._local.conndefclose_current(self):"""关闭当前线程的连接"""ifhasattr(self._local,"conn"):self._local.conn.close()delself._local.conn八、性能监控与调优
8.1 慢查询日志
importtimeimportlogging logger=logging.getLogger(__name__)classMonitoredDatabase(Database):"""带性能监控的数据库"""SLOW_QUERY_THRESHOLD=1.0# 慢查询阈值(秒)defquery(self,sql:str,params:tuple=()):start=time.time()result=super().query(sql,params)elapsed=time.time()-startifelapsed>self.SLOW_QUERY_THRESHOLD:logger.warning(f"慢查询 ({elapsed:.2f}s):{sql[:100]}... "f"参数:{params}")returnresult8.2 数据库统计
defget_db_stats(self)->dict:"""获取数据库统计信息"""stats={}# 文件大小importos stats["file_size_mb"]=os.path.getsize(self._db_path)/1024/1024# 各表行数tables=["daily_none","daily_qfq","daily_hfq","stock_basic"]fortableintables:count=self.query_one(f"SELECT COUNT(*) as c FROM{table}")stats[f"{table}_rows"]=count["c"]# 页面统计page_count=self.query_one("PRAGMA page_count")["page_count"]page_size=self.query_one("PRAGMA page_size")["page_size"]stats["page_count"]=page_count stats["page_size"]=page_size stats["total_size_mb"]=page_count*page_size/1024/1024returnstats九、总结
GPFX 的 SQLite 数据层核心实践:
| 技术点 | 方案 | 效果 |
|---|---|---|
| 连接管理 | 懒加载 + 双重检查 + 可重入锁 | 线程安全,性能最优 |
| 事务处理 | 上下文管理器自动提交/回滚 | 数据一致性保证 |
| 批量操作 | executemany + 事务包裹 | 60 倍性能提升 |
| SQL 生成 | 元编程自动生成 UPSERT | 减少重复代码 |
| 安全防护 | 表名白名单 + 参数化查询 | 杜绝 SQL 注入 |
| 索引优化 | 高频字段 + 联合索引 | 查询毫秒级响应 |
| 数据迁移 | 事务包裹 + OR IGNORE | 平滑升级无中断 |
| 并发控制 | WAL 模式 + 写锁 | 单写多读不阻塞 |
| 性能监控 | 慢查询日志 + 统计分析 | 问题可追踪 |
这套数据层支撑了 5000+ 只股票、千万级日线数据的高效存储和查询,为上层量化分析提供了坚实的数据基础。
核心技术:SQLite WAL / Python sqlite3 / 事务管理 / DAO 模式 / 索引优化
性能指标:
- 批量插入:5000 条/0.5s
- 单股查询:< 10ms
- 全市场扫描:< 100ms