# src/certflow/services/query_builder.py
"""动态查询构建器 - 对应 Access 窗体的筛选逻辑
.. deprecated::
Access ``SignTb`` 直连 SQL 旧体系残留,与 ORM 版 ``QueryService`` 平行存在。
仅保留供 ``scripts/sync_databases_cli.py`` 同步 Access 刻印状态/合格证编号专用,
**不通过** ``certflow.services`` 公共接口导出,亦不纳入 ORM 查询体系。
若未来 Access 同步脚本下线,本模块应整体删除(见 BLUEPRINT §2.3 偏差 6)。
"""
from __future__ import annotations
import contextlib
from datetime import datetime, timedelta
from pathlib import Path
from typing import Any
import pandas as pd
from loguru import logger
[文档]
class QueryBuilder:
"""动态查询构建器
对应 Access 窗体 Form_筛选数据到标牌刻印数据源表 的查询逻辑
"""
# 字段映射:界面字段名 -> SignTb 字段名
FIELD_MAPPING = {
"JHRQ": "JiHuaRiQi", # 计划日期
"DHDW": "DingHuoDanWei", # 订货单位
"XMMC": "XiangMuMingCheng", # 项目名称
"YHDH": "YaoHuoDanHaoFromXSB", # 要货单号
"ZL": "ZhongLei", # 种类
"ZXM": "ZiXiangMu", # 子项目(实际是生产令号)
"CPMC": "ChanPinMingCheng", # 产品名称
"CPXH": "ChanPinXingHao", # 产品型号
"DN": "DN", # 公称通径
"PN": "PN", # 公称压力
"SYJZ": "ShiYongJieZhi", # 适用介质
"SYWD": "ShiYongWenDu", # 适用温度
"CCNF": "ChuChangRiQiYear", # 出厂年份
"CCYF": "ChuChangRiQiMonth", # 出厂月份
"SN": "SN", # 出厂编号
"KKS": "KKS", # KKS编码
"FTCZ": "FaTiCaiZhi", # 阀体材质
"FGCZ": "FaGanCaiZhi", # 阀杆材质
"QBJCZ": "QiBiJianCaiZhi", # 启闭件材质
# 新增字段 - 注意实际字段名
"XSDDH": "SheBeiWeiHao", # 销售订单号(实际存储在设备位号字段)
"SCLH": "ZiXiangMu", # 生产令号(实际存储在子项目字段)
"hgz_ShuLiang": "hgz_ShuLiang", # 合格证数量
"hgz_PrintTime": "hgz_PrintTime", # 合格证打印时间
"Sign_ALL": "Sign_ALL", # 刻印完毕标记
}
def __init__(self) -> None:
"""初始化查询构建器
创建空的查询条件容器,默认返回全部字段。
"""
self.conditions: dict[str, Any] = {}
self.dn_range: tuple = (None, None)
self.pn_range: tuple = (None, None)
self.id_range: tuple = (None, None)
self.quantity_range: tuple = (None, None)
self.print_date_range: tuple = (None, None)
self.signed_only: bool | None = None
self.has_certificate: bool | None = None
self.combo_filters: dict[str, Any] = {}
self.full_fields: bool = True # Check_326: 是否返回全部字段
[文档]
def set_condition(self, field: str, value: Any) -> None:
"""设置筛选条件
Args:
field: 字段名(如 JHRQ, DHDW, XMMC)
value: 筛选值(支持 LIKE 模糊匹配)
"""
if value and str(value).strip():
self.conditions[field] = value
elif field in self.conditions:
del self.conditions[field]
[文档]
def set_dn_range(self, min_dn: int, max_dn: int) -> None:
"""设置口径范围筛选
Args:
min_dn: 最小口径值
max_dn: 最大口径值
"""
self.dn_range = (min_dn, max_dn)
[文档]
def set_pn_range(self, min_pn: float, max_pn: float) -> None:
"""设置公称压力范围筛选
Args:
min_pn: 最小公称压力值
max_pn: 最大公称压力值
"""
self.pn_range = (min_pn, max_pn)
[文档]
def set_id_range(self, min_id: int, max_id: int) -> None:
"""设置 ID 范围筛选
Args:
min_id: 最小 ID 值
max_id: 最大 ID 值
"""
self.id_range = (min_id, max_id)
[文档]
def set_quantity_range(self, min_qty: int, max_qty: int) -> None:
"""设置合格证数量范围筛选
Args:
min_qty: 最小数量值
max_qty: 最大数量值
"""
self.quantity_range = (min_qty, max_qty)
[文档]
def set_print_date_range(self, start_date: str, end_date: str) -> None:
"""设置铭牌打印日期范围筛选
Args:
start_date: 开始日期 (YYYY-MM-DD 格式)
end_date: 结束日期 (YYYY-MM-DD 格式)
"""
if start_date and end_date:
self.print_date_range = (start_date, end_date)
[文档]
def set_signed_only(self, signed: bool) -> None:
"""设置是否只查询已刻印的记录
Args:
signed: True 只查询已刻印,False 只查询未刻印
"""
self.signed_only = signed
[文档]
def set_has_certificate(self, has_cert: bool) -> None:
"""设置是否只查询有合格证的记录
Args:
has_cert: True 只查询有合格证的记录,False 只查询无合格证的记录
"""
self.has_certificate = has_cert
[文档]
def set_combo_filter(self, field: str, value: Any) -> None:
"""设置下拉筛选条件(精确匹配)
Args:
field: 字段名
value: 精确匹配值,为 None 则删除该条件
"""
if value is not None:
self.combo_filters[field] = value
elif field in self.combo_filters:
del self.combo_filters[field]
[文档]
def set_full_fields(self, full: bool) -> None:
"""设置是否返回全部字段
Args:
full: True 返回全部字段,False 返回简化字段集
"""
self.full_fields = full
# ========== 拆分后的 WHERE 子句构建方法 ==========
def _build_where_clause(self) -> tuple[str, list]:
"""构建 WHERE 子句
Returns:
(where_sql, params) 元组
"""
where_parts = []
params = []
self._add_text_filters(where_parts, params)
self._add_combo_filters(where_parts, params)
self._add_range_filters(where_parts, params)
self._add_status_filters(where_parts, params)
self._add_date_filters(where_parts, params)
if where_parts:
return " WHERE " + " AND ".join(where_parts), params
return "", []
def _add_text_filters(self, where_parts: list, params: list) -> None:
"""添加文本模糊匹配筛选条件
通过 LIKE 运算符将界面字段映射到数据库字段进行模糊匹配。
Args:
where_parts: WHERE 子句片段列表(原地修改)
params: 参数值列表(原地修改)
"""
for field, value in self.conditions.items():
if field in self.FIELD_MAPPING:
db_field = self.FIELD_MAPPING[field]
where_parts.append(f"[{db_field}] LIKE ?")
params.append(f"%{value}%")
def _add_combo_filters(self, where_parts: list, params: list) -> None:
"""添加下拉精确匹配筛选条件
通过等值匹配运算符将下拉选择映射到数据库字段。
Args:
where_parts: WHERE 子句片段列表(原地修改)
params: 参数值列表(原地修改)
"""
for field, value in self.combo_filters.items():
if field in self.FIELD_MAPPING:
db_field = self.FIELD_MAPPING[field]
where_parts.append(f"[{db_field}] = ?")
params.append(value)
def _add_range_filters(self, where_parts: list, params: list) -> None:
"""添加数值范围筛选条件
处理 DN/PN/ID/合格证数量等数值字段的 BETWEEN 范围查询。
Args:
where_parts: WHERE 子句片段列表(原地修改)
params: 参数值列表(原地修改)
"""
if self.dn_range[0] is not None and self.dn_range[1] is not None:
where_parts.append("Val([DN]) BETWEEN ? AND ?")
params.extend([self.dn_range[0], self.dn_range[1]])
if self.pn_range[0] is not None and self.pn_range[1] is not None:
where_parts.append("Val([PN]) BETWEEN ? AND ?")
params.extend([self.pn_range[0], self.pn_range[1]])
if self.id_range[0] is not None and self.id_range[1] is not None:
where_parts.append("Val([ID]) BETWEEN ? AND ?")
params.extend([self.id_range[0], self.id_range[1]])
if self.quantity_range[0] is not None and self.quantity_range[1] is not None:
where_parts.append("Val([hgz_ShuLiang]) BETWEEN ? AND ?")
params.extend([self.quantity_range[0], self.quantity_range[1]])
def _add_status_filters(self, where_parts: list, params: list) -> None:
"""添加状态筛选条件
处理刻印状态 (Sign_ALL) 和合格证有无 (hgz_ShuLiang) 的布尔筛选。
Args:
where_parts: WHERE 子句片段列表(原地修改)
params: 参数值列表(原地修改)
"""
if self.signed_only is not None:
where_parts.append("[Sign_ALL] = ?")
params.append(-1 if self.signed_only else 0)
if self.has_certificate is not None:
if self.has_certificate:
where_parts.append("Val([hgz_ShuLiang]) > 0")
else:
where_parts.append("(Val([hgz_ShuLiang]) = 0 OR [hgz_ShuLiang] IS NULL)")
def _add_date_filters(self, where_parts: list, params: list) -> None:
"""添加日期范围筛选条件
通过 Access 的 Format 函数格式化日期字段后进行范围比较。
Args:
where_parts: WHERE 子句片段列表(原地修改)
params: 参数值列表(原地修改)
"""
if self.print_date_range[0] and self.print_date_range[1]:
where_parts.append("Format([Sign_PrintTime], 'YYYY-MM-DD') BETWEEN ? AND ?")
params.extend([self.print_date_range[0], self.print_date_range[1]])
# ========== SELECT 字段 ==========
def _get_select_fields(self) -> str:
"""获取 SELECT 字段列表
根据 full_fields 标志决定返回完整字段集还是简化字段集。
Returns:
str: SQL SELECT 子句的字段列表
"""
if self.full_fields:
return self._get_full_select_fields()
return self._get_simple_select_fields()
def _get_full_select_fields(self) -> str:
"""获取全部 SELECT 字段 (26个字段)
Returns:
str: 包含所有字段的 SQL SELECT 子句
"""
return """
SignTb.id AS ID,
SignTb.SheBeiWeiHao AS 销售订单号,
SignTb.ZiXiangMu AS 生产令号,
SignTb.JiHuaRiQi AS 计划日期,
SignTb.DingHuoDanWei AS 订货单位,
SignTb.XiangMuMingCheng AS 项目名称,
SignTb.ChanPinMingCheng AS 产品名称,
SignTb.ChanPinXingHao AS 产品型号,
SignTb.DN AS 公称通径,
SignTb.PN AS 公称压力,
SignTb.ChuChangRiQiYear AS 出厂日期年,
SignTb.ChuChangRiQiMonth AS 出厂日期月,
SignTb.SN AS 出厂编号,
SignTb.KKS AS KKS编码,
SignTb.Sign_ALL AS 刻印完毕标记,
SignTb.hgz_ShuLiang AS 合格证数量,
SignTb.hgz_PrintTime AS 合格证打印时间,
SignTb.Sign_PrintTime AS 铭牌打印时间,
SignTb.ShiYongJieZhi AS 适用介质,
SignTb.ShiYongWenDu AS 适用温度,
SignTb.FaTiCaiZhi AS 阀体材质,
SignTb.FaGanCaiZhi AS 阀杆材质,
SignTb.QiBiJianCaiZhi AS 启闭件材质,
SignTb.SNFromXSB AS 销售产品编号,
SignTb.ZhongLei AS 销售种类,
SignTb.YaoHuoDanHaoFromXSB AS 营销部要货单号
"""
def _get_simple_select_fields(self) -> str:
"""获取简化 SELECT 字段 (14个核心字段)
Returns:
str: 包含核心字段的 SQL SELECT 子句
"""
return """
SignTb.id AS ID,
SignTb.SheBeiWeiHao AS 销售订单号,
SignTb.ZiXiangMu AS 生产令号,
SignTb.JiHuaRiQi AS 计划日期,
SignTb.DingHuoDanWei AS 订货单位,
SignTb.XiangMuMingCheng AS 项目名称,
SignTb.ZhongLei AS 销售种类,
SignTb.ChanPinXingHao AS 产品型号,
SignTb.DN AS 公称通径,
SignTb.hgz_ShuLiang AS 合格证数量,
SignTb.SN AS 出厂编号,
SignTb.KKS AS KKS编码,
SignTb.hgz_PrintTime AS 合格证打印时间,
SignTb.Sign_PrintTime AS 铭牌打印时间
"""
# ========== 查询执行 ==========
[文档]
def build_query(self) -> tuple[str, list]:
"""构建完整查询 SQL
Returns:
(sql, params) 元组
"""
select_fields = self._get_select_fields()
where_clause, params = self._build_where_clause()
sql = f"""
SELECT {select_fields}
FROM SignTb
{where_clause}
ORDER BY SignTb.ChanPinXingHao, SignTb.DN, SignTb.SN
"""
return sql.strip(), params
[文档]
def execute(self, db: Any) -> pd.DataFrame:
"""执行查询并返回结果
构建 SQL 查询,将查询 SQL 记录到日志文件,然后执行查询。
Args:
db: AccessDatabase 实例,需提供 query(sql, params) 方法
Returns:
pd.DataFrame: 查询结果数据
"""
sql, params = self.build_query()
# 写入文件避免截断
from datetime import datetime
from pathlib import Path
log_dir = Path(__file__).parent.parent.parent.parent / "logs"
log_dir.mkdir(exist_ok=True)
self._clean_old_sql_files(log_dir)
timestamp = datetime.now().strftime("%Y%m%d_%H%M%S_%f")[:-3]
sql_file = log_dir / f"query_{timestamp}.sql"
self._save_sql_to_file(sql_file, sql, params)
logger.info(f"SQL 查询已保存到: {sql_file}")
logger.debug(f"执行查询: {sql[:200]}...")
if params:
return db.query(sql, params)
return db.query(sql)
def _clean_old_sql_files(self, log_dir: Path) -> None:
"""清理7天前的 SQL 日志文件
Args:
log_dir: 日志目录路径
"""
cutoff = datetime.now() - timedelta(days=7)
for f in log_dir.glob("query_*.sql"):
if datetime.fromtimestamp(f.stat().st_mtime) < cutoff:
f.unlink()
def _save_sql_to_file(self, file_path: Path, sql: str, params: list) -> None:
"""保存 SQL 查询到日志文件
Args:
file_path: 输出文件路径
sql: SQL 查询语句(含占位符)
params: 参数值列表
"""
with open(file_path, "w", encoding="utf-8") as f:
f.write("=" * 100 + "\n")
f.write(f"查询时间: {datetime.now().strftime('%Y-%m-%d %H:%M:%S.%f')[:-3]}\n")
f.write("=" * 100 + "\n\n")
f.write("SQL QUERY:\n")
f.write("-" * 100 + "\n")
f.write(sql + "\n\n")
if params:
self._write_params_to_file(f, params)
self._write_full_sql_to_file(f, sql, params)
f.write("=" * 100 + "\n")
def _write_params_to_file(self, f: Any, params: list) -> None:
"""写入参数列表到日志文件
Args:
f: 已打开的文件句柄(写入模式)。
params: 参数值列表。
"""
f.write("参数列表:\n")
f.write("-" * 100 + "\n")
for i, p in enumerate(params):
f.write(f" [{i}] {p} (type: {type(p).__name__})\n")
f.write("\n")
def _write_full_sql_to_file(self, f: Any, sql: str, params: list) -> None:
"""写入参数替换后的完整 SQL 到日志文件
Args:
f: 已打开的文件句柄(写入模式)。
sql: 含 ? 占位符的 SQL 语句。
params: 参数值列表(按顺序替换占位符)。
"""
f.write("完整 SQL(可直接在 Access 中执行):\n")
f.write("-" * 100 + "\n")
full_sql = sql
for p in params:
if isinstance(p, str):
full_sql = full_sql.replace("?", f"'{p}'", 1)
elif isinstance(p, (int, float)):
full_sql = full_sql.replace("?", str(p), 1)
elif p is None:
full_sql = full_sql.replace("?", "NULL", 1)
else:
full_sql = full_sql.replace("?", str(p), 1)
f.write(full_sql + "\n\n")
[文档]
class TempTableManager:
"""临时表管理器 - 对应 Access 的分表查询和合并功能"""
def __init__(self, max_tables: int = 50) -> None:
"""初始化临时表管理器
Args:
max_tables: 最大允许创建的临时表数量(预留,未强制限制)
"""
self.max_tables = max_tables
self.temp_tables: list[str] = []
self.temp_count = 0
[文档]
def create_temp_table(
self, db: Any, query_builder: QueryBuilder, table_prefix: str = "查询结果表_Path"
) -> str:
"""创建临时表存储查询结果
Args:
db: AccessDatabase 实例
query_builder: QueryBuilder 实例
table_prefix: 临时表前缀
Returns:
创建的临时表名
"""
self.temp_count += 1
temp_table_name = f"{table_prefix}{self.temp_count:03d}"
sql, params = query_builder.build_query()
# 创建临时表
create_sql = f"""
SELECT * INTO [{temp_table_name}]
FROM ({sql}) AS T
"""
# 执行创建
if params:
db.cursor.execute(create_sql, params)
else:
db.cursor.execute(create_sql)
db.conn.commit()
self.temp_tables.append(temp_table_name)
logger.info(f"创建临时表 {temp_table_name}")
return temp_table_name
[文档]
def merge_tables(
self, db: Any, target_table: str = "FormResumeDataTb", distinct: bool = True
) -> int:
"""合并所有临时表到目标表
Args:
db: AccessDatabase 实例
target_table: 目标表名
distinct: 是否去重
Returns:
插入的记录数
"""
if not self.temp_tables:
logger.warning("没有临时表可合并")
return 0
# 构建 UNION 查询
union_operator = "UNION" if distinct else "UNION ALL"
union_parts = [f"SELECT * FROM [{t}]" for t in self.temp_tables]
union_sql = f" {union_operator} ".join(union_parts)
# 删除已存在的目标表
with contextlib.suppress(Exception):
db.cursor.execute(f"DROP TABLE IF EXISTS [{target_table}]")
# 创建新表
insert_sql = f"""
SELECT * INTO [{target_table}]
FROM ({union_sql}) AS T
ORDER BY ID
"""
db.cursor.execute(insert_sql)
db.conn.commit()
inserted = db.cursor.rowcount
logger.info(f"合并 {len(self.temp_tables)} 个临时表到 {target_table},共 {inserted} 条记录")
return inserted
[文档]
def clean_temp_tables(self, db: Any) -> None:
"""清理所有临时表
Args:
db: AccessDatabase 实例
"""
for table in self.temp_tables:
try:
db.cursor.execute(f"DROP TABLE IF EXISTS [{table}]")
logger.debug(f"删除临时表 {table}")
except Exception as e:
logger.warning(f"删除临时表 {table} 失败: {e}")
self.temp_tables.clear()
self.temp_count = 0
db.conn.commit()
[文档]
class BatchUpdater:
"""批量更新管理器 - 对应 Access 的批量刻印状态更新"""
def __init__(self, db: Any) -> None:
"""初始化统计报表生成器
Args:
db: AccessDatabase 实例。
"""
self.db = db
[文档]
def update_print_status(
self, query_builder: QueryBuilder, signed: bool, dry_run: bool = True
) -> dict[str, int]:
"""批量更新刻印状态
Args:
query_builder: 查询条件构建器
signed: 是否标记为已刻印
dry_run: 是否试运行
Returns:
更新结果统计
"""
# 先查询符合条件的记录数
sql, params = query_builder.build_query()
count_sql = f"SELECT COUNT(*) AS N FROM ({sql}) AS T"
if params:
self.db.cursor.execute(count_sql, params)
else:
self.db.cursor.execute(count_sql)
total = self.db.cursor.fetchone()[0]
if total == 0:
logger.info("没有符合条件的记录")
return {"total": 0, "updated": 0}
if dry_run:
logger.info(f"[试运行] 将更新 {total} 条记录为 {'已刻印' if signed else '未刻印'}")
return {"total": total, "updated": 0}
# 构建更新 SQL
sign_value = -1 if signed else 0
sign_text = "已刻印" if signed else "未刻印"
# 提取 WHERE 条件
where_clause, where_params = query_builder._build_where_clause()
update_sql = f"""
UPDATE SignTb SET
Sign_XingHao = ?,
Sign_DN = ?,
Sign_PN = ?,
Sign_JieZhi = ?,
Sign_WenDu = ?,
Sign_Year = ?,
Sign_Month = ?,
Sign_SN = ?,
Sign_KKS = ?,
Sign_ALL = ?,
Sign_FTCZ = ?,
Sign_FGCZ = ?,
Sign_QBJCZ = ?,
Sign_PrintTime = Now(),
hgz_xlMuBanName = hgz_xlMuBanName & '--手动标记铭牌刻印状态为:{sign_text}'
{where_clause}
"""
# 13 个布尔字段 + 2 个扩展 = 15 个参数
params = [sign_value] * 13 + where_params
self.db.cursor.execute(update_sql, params)
self.db.conn.commit()
updated = self.db.cursor.rowcount
logger.info(f"更新 {updated} 条记录为 {sign_text}")
return {"total": total, "updated": updated}
[文档]
class StatisticsReporter:
"""统计报表生成器 - 对应 Access 的按小时统计功能"""
def __init__(self, db: Any) -> None:
"""初始化统计报表生成器
Args:
db: AccessDatabase 实例。
"""
self.db = db
[文档]
def get_hourly_stats(
self, start_date: str | None = None, end_date: str | None = None
) -> pd.DataFrame:
"""按小时统计合格证和铭牌打印数量
Args:
start_date: 开始日期 (YYYY-MM-DD)
end_date: 结束日期 (YYYY-MM-DD)
Returns:
统计结果 DataFrame
"""
where_clause = ""
params = []
if start_date and end_date:
where_clause = "WHERE Format(Sign_PrintTime, 'YYYY-MM-DD') BETWEEN ? AND ?"
params = [start_date, end_date]
sql = f"""
SELECT
Format(Sign_PrintTime, 'YYYY-MM-DD') AS 日期,
Format(Sign_PrintTime, 'HH') AS 小时,
COUNT(*) AS 铭牌打印数量,
SUM(IIF(hgz_PrintTime IS NOT NULL, 1, 0)) AS 合格证打印数量
FROM SignTb
{where_clause}
GROUP BY Format(Sign_PrintTime, 'YYYY-MM-DD'), Format(Sign_PrintTime, 'HH')
ORDER BY 日期 DESC, 小时
"""
if params:
return self.db.query(sql, params)
return self.db.query(sql)
[文档]
def get_customer_summary(self) -> pd.DataFrame:
"""按订货单位汇总统计。
Returns:
pd.DataFrame: 各订货单位的记录数、已/未刻印数量及最早/最新刻印时间。
"""
sql = """
SELECT
DingHuoDanWei AS 订货单位,
COUNT(*) AS 总数量,
SUM(IIF(Sign_ALL = -1, 1, 0)) AS 已刻印数量,
SUM(IIF(Sign_ALL = 0, 1, 0)) AS 未刻印数量,
MIN(Sign_PrintTime) AS 最早刻印时间,
MAX(Sign_PrintTime) AS 最新刻印时间
FROM SignTb
GROUP BY DingHuoDanWei
ORDER BY COUNT(*) DESC
"""
return self.db.query(sql)