Refactor SQL script creation in ReportScript.sql for improved readability and performance. Update DatabaseManager to enhance connection handling with detailed documentation. Add a new method in DataProcessor to parse SQL scripts, ensuring reliable extraction of valid SQL statements.
This commit is contained in:
+15
-117
@@ -91,14 +91,7 @@ ADD INDEX `TC`(`日期时间`, `NCGI`),
|
|||||||
ADD INDEX `DC`(`日期`, `NCGI`);
|
ADD INDEX `DC`(`日期`, `NCGI`);
|
||||||
|
|
||||||
|
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `CCE_MAX_IDX`(
|
CREATE TEMPORARY TABLE IF NOT EXISTS `CCE_MAX_IDX` AS SELECT `日期`, CGI, MAX(`PDCCH资源利用率`) AS `PDCCH资源利用率_忙时` FROM 4G GROUP BY `日期`, CGI;
|
||||||
SELECT
|
|
||||||
`日期`,
|
|
||||||
CGI,
|
|
||||||
MAX( `PDCCH资源利用率`) AS `PDCCH资源利用率_忙时`
|
|
||||||
FROM 4G
|
|
||||||
GROUP BY `日期`, CGI
|
|
||||||
);
|
|
||||||
|
|
||||||
ALTER TABLE `CCE_MAX_IDX`
|
ALTER TABLE `CCE_MAX_IDX`
|
||||||
ADD COLUMN `日期时间` datetime,
|
ADD COLUMN `日期时间` datetime,
|
||||||
@@ -114,23 +107,11 @@ SET
|
|||||||
`CCE_MAX_IDX`.`日期时间`=4G.`日期时间`;
|
`CCE_MAX_IDX`.`日期时间`=4G.`日期时间`;
|
||||||
|
|
||||||
|
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `CCE_MAX` (
|
CREATE TEMPORARY TABLE IF NOT EXISTS `CCE_MAX` AS SELECT 4G.* FROM `CCE_MAX_IDX` LEFT JOIN 4G ON `CCE_MAX_IDX`.CGI=4G.CGI AND `CCE_MAX_IDX`.`日期时间`=4G.`日期时间`;
|
||||||
SELECT 4G.*
|
|
||||||
FROM `CCE_MAX_IDX` LEFT JOIN 4G
|
|
||||||
ON `CCE_MAX_IDX`.CGI=4G.CGI AND
|
|
||||||
`CCE_MAX_IDX`.`日期时间`=4G.`日期时间`
|
|
||||||
);
|
|
||||||
|
|
||||||
DROP TABLE IF EXISTS `CCE_MAX_IDX`;
|
DROP TABLE IF EXISTS `CCE_MAX_IDX`;
|
||||||
|
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `ULPRB_MAX_IDX`(
|
CREATE TEMPORARY TABLE IF NOT EXISTS `ULPRB_MAX_IDX` AS SELECT `日期`, CGI, MAX(`上行PUSCH利用率`) AS `上行PUSCH利用率_忙时` FROM 4G GROUP BY `日期`, CGI;
|
||||||
SELECT
|
|
||||||
`日期`,
|
|
||||||
CGI,
|
|
||||||
MAX( `上行PUSCH利用率`) AS `上行PUSCH利用率_忙时`
|
|
||||||
FROM 4G
|
|
||||||
GROUP BY `日期`, CGI
|
|
||||||
);
|
|
||||||
|
|
||||||
ALTER TABLE `ULPRB_MAX_IDX`
|
ALTER TABLE `ULPRB_MAX_IDX`
|
||||||
ADD COLUMN `日期时间` datetime,
|
ADD COLUMN `日期时间` datetime,
|
||||||
@@ -154,14 +135,7 @@ CREATE TABLE IF NOT EXISTS `ULPRB_MAX` (
|
|||||||
|
|
||||||
DROP TABLE IF EXISTS `ULPRB_MAX_IDX`;
|
DROP TABLE IF EXISTS `ULPRB_MAX_IDX`;
|
||||||
|
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `DLPRB_MAX_IDX`(
|
CREATE TEMPORARY TABLE IF NOT EXISTS `DLPRB_MAX_IDX` AS SELECT `日期`, CGI, MAX(`下行PDSCH利用率`) AS `下行PDSCH利用率_忙时` FROM 4G GROUP BY `日期`, CGI;
|
||||||
SELECT
|
|
||||||
`日期`,
|
|
||||||
CGI,
|
|
||||||
MAX( `下行PDSCH利用率`) AS `下行PDSCH利用率_忙时`
|
|
||||||
FROM 4G
|
|
||||||
GROUP BY `日期`, CGI
|
|
||||||
);
|
|
||||||
|
|
||||||
ALTER TABLE `DLPRB_MAX_IDX`
|
ALTER TABLE `DLPRB_MAX_IDX`
|
||||||
ADD COLUMN `日期时间` datetime,
|
ADD COLUMN `日期时间` datetime,
|
||||||
@@ -176,12 +150,7 @@ ON
|
|||||||
SET
|
SET
|
||||||
`DLPRB_MAX_IDX`.`日期时间`=4G.`日期时间`;
|
`DLPRB_MAX_IDX`.`日期时间`=4G.`日期时间`;
|
||||||
|
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `DLPRB_MAX` (
|
CREATE TEMPORARY TABLE IF NOT EXISTS `DLPRB_MAX` AS SELECT 4G.* FROM `DLPRB_MAX_IDX` LEFT JOIN 4G ON `DLPRB_MAX_IDX`.CGI=4G.CGI AND `DLPRB_MAX_IDX`.`日期时间`=4G.`日期时间`;
|
||||||
SELECT 4G.*
|
|
||||||
FROM `DLPRB_MAX_IDX` LEFT JOIN 4G
|
|
||||||
ON `DLPRB_MAX_IDX`.CGI=4G.CGI AND
|
|
||||||
`DLPRB_MAX_IDX`.`日期时间`=4G.`日期时间`
|
|
||||||
);
|
|
||||||
|
|
||||||
DROP TABLE IF EXISTS `DLPRB_MAX_IDX`;
|
DROP TABLE IF EXISTS `DLPRB_MAX_IDX`;
|
||||||
|
|
||||||
@@ -191,35 +160,11 @@ ALTER TABLE `ULPRB_MAX` ADD INDEX `DC`(`日期`, `CGI`);
|
|||||||
|
|
||||||
|
|
||||||
|
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `4G_Date` (
|
CREATE TEMPORARY TABLE IF NOT EXISTS `4G_Date` AS SELECT `日期`, CGI FROM 4G GROUP BY `日期`, CGI;
|
||||||
SELECT `日期`, CGI FROM 4G GROUP BY `日期`, CGI
|
|
||||||
);
|
|
||||||
ALTER TABLE `4G_Date` ADD INDEX `DC`(`日期`, CGI);
|
ALTER TABLE `4G_Date` ADD INDEX `DC`(`日期`, CGI);
|
||||||
|
|
||||||
|
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `4G_MAX_IDX` (
|
CREATE TEMPORARY TABLE IF NOT EXISTS `4G_MAX_IDX` AS SELECT 4G_Date.`日期`, 4G_Date.CGI, GREATEST(`ULPRB_MAX`.`上行PUSCH利用率`, DLPRB_MAX.`下行PDSCH利用率`,CCE_MAX.`PDCCH资源利用率`) AS `最大值` FROM 4G_Date INNER JOIN DLPRB_MAX ON 4G_Date.`日期` = DLPRB_MAX.`日期` AND 4G_Date.CGI = DLPRB_MAX.CGI INNER JOIN CCE_MAX ON 4G_Date.`日期` = CCE_MAX.`日期` AND 4G_Date.CGI = CCE_MAX.CGI INNER JOIN ULPRB_MAX ON 4G_Date.`日期` = ULPRB_MAX.`日期` AND 4G_Date.CGI = ULPRB_MAX.CGI;
|
||||||
SELECT
|
|
||||||
4G_Date.`日期`,
|
|
||||||
4G_Date.CGI,
|
|
||||||
GREATEST(`ULPRB_MAX`.`上行PUSCH利用率`, DLPRB_MAX.`下行PDSCH利用率`,CCE_MAX.`PDCCH资源利用率`) AS `最大值`
|
|
||||||
FROM
|
|
||||||
4G_Date
|
|
||||||
INNER JOIN
|
|
||||||
DLPRB_MAX
|
|
||||||
ON
|
|
||||||
4G_Date.`日期` = DLPRB_MAX.`日期` AND
|
|
||||||
4G_Date.CGI = DLPRB_MAX.CGI
|
|
||||||
INNER JOIN
|
|
||||||
CCE_MAX
|
|
||||||
ON
|
|
||||||
4G_Date.`日期` = CCE_MAX.`日期` AND
|
|
||||||
4G_Date.CGI = CCE_MAX.CGI
|
|
||||||
INNER JOIN
|
|
||||||
ULPRB_MAX
|
|
||||||
ON
|
|
||||||
4G_Date.`日期` = ULPRB_MAX.`日期` AND
|
|
||||||
4G_Date.CGI = ULPRB_MAX.CGI
|
|
||||||
);
|
|
||||||
DROP TABLE IF EXISTS `4G_Date`;
|
DROP TABLE IF EXISTS `4G_Date`;
|
||||||
ALTER TABLE `4G_MAX_IDX` ADD INDEX `DC`(`日期`, CGI), ADD INDEX (`最大值`);
|
ALTER TABLE `4G_MAX_IDX` ADD INDEX `DC`(`日期`, CGI), ADD INDEX (`最大值`);
|
||||||
|
|
||||||
@@ -275,9 +220,7 @@ DROP TABLE IF EXISTS `4G_MAX_IDX`;
|
|||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `日流量`(SELECT `日期`,CGI, SUM(`上下行总流量_GB`) AS `日流量` FROM 4G GROUP BY `日期`,CGI);
|
CREATE TEMPORARY TABLE IF NOT EXISTS `日流量`(SELECT `日期`,CGI, SUM(`上下行总流量_GB`) AS `日流量` FROM 4G GROUP BY `日期`,CGI);
|
||||||
ALTER TABLE `日流量` ADD INDEX (`CGI`);
|
ALTER TABLE `日流量` ADD INDEX (`CGI`);
|
||||||
DROP TABLE IF EXISTS `4G日均流量`;
|
DROP TABLE IF EXISTS `4G日均流量`;
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `4G日均流量` (
|
CREATE TEMPORARY TABLE IF NOT EXISTS `4G日均流量` AS SELECT CGI, AVG(`日流量`) AS `日均流量`, COUNT(`日期`) AS `流量有效天数` FROM `日流量` WHERE `日流量` > 0 GROUP BY CGI;
|
||||||
SELECT CGI, AVG(`日流量`) AS `日均流量`, COUNT(`日期`) AS `流量有效天数` FROM `日流量` WHERE `日流量` > 0 GROUP BY CGI
|
|
||||||
);
|
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
@@ -313,14 +256,7 @@ DROP TABLE IF EXISTS `4G`;
|
|||||||
|
|
||||||
# 5G脚本
|
# 5G脚本
|
||||||
|
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `ULPRB_MAX_IDX`(
|
CREATE TEMPORARY TABLE IF NOT EXISTS `ULPRB_MAX_IDX` AS SELECT `日期`, NCGI, MAX(`上行PRB平均利用率`) AS `上行PRB平均利用率_忙时` FROM 5G GROUP BY `日期`, NCGI;
|
||||||
SELECT
|
|
||||||
`日期`,
|
|
||||||
NCGI,
|
|
||||||
MAX(`上行PRB平均利用率`) AS `上行PRB平均利用率_忙时`
|
|
||||||
FROM 5G
|
|
||||||
GROUP BY `日期`, NCGI
|
|
||||||
);
|
|
||||||
|
|
||||||
ALTER TABLE `ULPRB_MAX_IDX`
|
ALTER TABLE `ULPRB_MAX_IDX`
|
||||||
ADD COLUMN `日期时间` datetime,
|
ADD COLUMN `日期时间` datetime,
|
||||||
@@ -335,22 +271,10 @@ ON
|
|||||||
SET
|
SET
|
||||||
`ULPRB_MAX_IDX`.`日期时间`=5G.`日期时间`;
|
`ULPRB_MAX_IDX`.`日期时间`=5G.`日期时间`;
|
||||||
|
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `ULPRB_MAX` (
|
CREATE TEMPORARY TABLE IF NOT EXISTS `ULPRB_MAX` AS SELECT 5G.* FROM `ULPRB_MAX_IDX` LEFT JOIN 5G ON `ULPRB_MAX_IDX`.NCGI=5G.NCGI AND `ULPRB_MAX_IDX`.`日期时间`=5G.`日期时间`;
|
||||||
SELECT 5G.*
|
|
||||||
FROM `ULPRB_MAX_IDX` LEFT JOIN 5G
|
|
||||||
ON `ULPRB_MAX_IDX`.NCGI=5G.NCGI AND
|
|
||||||
`ULPRB_MAX_IDX`.`日期时间`=5G.`日期时间`
|
|
||||||
);
|
|
||||||
|
|
||||||
# 下行PRB MAX表
|
# 下行PRB MAX表
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `DLPRB_MAX_IDX`(
|
CREATE TEMPORARY TABLE IF NOT EXISTS `DLPRB_MAX_IDX` AS SELECT `日期`, NCGI, MAX(`下行PRB平均利用率`) AS `下行PRB平均利用率_忙时` FROM 5G GROUP BY `日期`, NCGI;
|
||||||
SELECT
|
|
||||||
`日期`,
|
|
||||||
NCGI,
|
|
||||||
MAX( `下行PRB平均利用率`) AS `下行PRB平均利用率_忙时`
|
|
||||||
FROM 5G
|
|
||||||
GROUP BY `日期`, NCGI
|
|
||||||
);
|
|
||||||
|
|
||||||
ALTER TABLE `DLPRB_MAX_IDX`
|
ALTER TABLE `DLPRB_MAX_IDX`
|
||||||
ADD COLUMN `日期时间` datetime,
|
ADD COLUMN `日期时间` datetime,
|
||||||
@@ -365,12 +289,7 @@ ON
|
|||||||
SET
|
SET
|
||||||
`DLPRB_MAX_IDX`.`日期时间`=5G.`日期时间`;
|
`DLPRB_MAX_IDX`.`日期时间`=5G.`日期时间`;
|
||||||
|
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `DLPRB_MAX` (
|
CREATE TEMPORARY TABLE IF NOT EXISTS `DLPRB_MAX` AS SELECT 5G.* FROM `DLPRB_MAX_IDX` LEFT JOIN 5G ON `DLPRB_MAX_IDX`.NCGI=5G.NCGI AND `DLPRB_MAX_IDX`.`日期时间`=5G.`日期时间`;
|
||||||
SELECT 5G.*
|
|
||||||
FROM `DLPRB_MAX_IDX` LEFT JOIN 5G
|
|
||||||
ON `DLPRB_MAX_IDX`.NCGI=5G.NCGI AND
|
|
||||||
`DLPRB_MAX_IDX`.`日期时间`=5G.`日期时间`
|
|
||||||
);
|
|
||||||
|
|
||||||
DROP TABLE IF EXISTS `DLPRB_MAX_IDX`;
|
DROP TABLE IF EXISTS `DLPRB_MAX_IDX`;
|
||||||
ALTER TABLE `DLPRB_MAX` ADD INDEX `DC`(`日期`, `NCGI`);
|
ALTER TABLE `DLPRB_MAX` ADD INDEX `DC`(`日期`, `NCGI`);
|
||||||
@@ -379,29 +298,10 @@ ALTER TABLE `ULPRB_MAX` ADD INDEX `DC`(`日期`, `NCGI`);
|
|||||||
# 5G 统计MAX数据
|
# 5G 统计MAX数据
|
||||||
|
|
||||||
DROP TABLE IF EXISTS `5G_Date`;
|
DROP TABLE IF EXISTS `5G_Date`;
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `5G_Date` (
|
CREATE TEMPORARY TABLE IF NOT EXISTS `5G_Date` AS SELECT `日期`, NCGI FROM 5G GROUP BY `日期`, NCGI;
|
||||||
SELECT `日期`, NCGI FROM 5G GROUP BY `日期`, NCGI
|
|
||||||
);
|
|
||||||
ALTER TABLE `5G_Date` ADD INDEX `DC`(`日期`, NCGI);
|
ALTER TABLE `5G_Date` ADD INDEX `DC`(`日期`, NCGI);
|
||||||
|
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `5G_MAX_IDX` (
|
CREATE TEMPORARY TABLE IF NOT EXISTS `5G_MAX_IDX` AS SELECT 5G_Date.`日期`, 5G_Date.NCGI, GREATEST(ULPRB_MAX.`上行PRB平均利用率`, DLPRB_MAX.`下行PRB平均利用率`) AS `最大值` FROM 5G_Date INNER JOIN DLPRB_MAX ON 5G_Date.`日期` = DLPRB_MAX.`日期` AND 5G_Date.NCGI = DLPRB_MAX.NCGI INNER JOIN ULPRB_MAX ON 5G_Date.`日期` = ULPRB_MAX.`日期` AND 5G_Date.NCGI = ULPRB_MAX.NCGI;
|
||||||
SELECT
|
|
||||||
5G_Date.`日期`,
|
|
||||||
5G_Date.NCGI,
|
|
||||||
GREATEST(ULPRB_MAX.`上行PRB平均利用率`, DLPRB_MAX.`下行PRB平均利用率`) AS `最大值`
|
|
||||||
FROM
|
|
||||||
5G_Date
|
|
||||||
INNER JOIN
|
|
||||||
DLPRB_MAX
|
|
||||||
ON
|
|
||||||
5G_Date.`日期` = DLPRB_MAX.`日期` AND
|
|
||||||
5G_Date.NCGI = DLPRB_MAX.NCGI
|
|
||||||
INNER JOIN
|
|
||||||
ULPRB_MAX
|
|
||||||
ON
|
|
||||||
5G_Date.`日期` = ULPRB_MAX.`日期` AND
|
|
||||||
5G_Date.NCGI = ULPRB_MAX.NCGI
|
|
||||||
);
|
|
||||||
|
|
||||||
DROP TABLE IF EXISTS `5G_Date`;
|
DROP TABLE IF EXISTS `5G_Date`;
|
||||||
ALTER TABLE `5G_MAX_IDX` ADD INDEX `DC`(`日期`, NCGI), ADD INDEX (`最大值`);
|
ALTER TABLE `5G_MAX_IDX` ADD INDEX `DC`(`日期`, NCGI), ADD INDEX (`最大值`);
|
||||||
@@ -440,9 +340,7 @@ DROP TABLE IF EXISTS `5G_MAX_IDX`;
|
|||||||
|
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `日流量`(SELECT `日期`,NCGI, SUM(`5G上下行总流量_GB`) AS `日流量` FROM 5G GROUP BY `日期`,NCGI);
|
CREATE TEMPORARY TABLE IF NOT EXISTS `日流量`(SELECT `日期`,NCGI, SUM(`5G上下行总流量_GB`) AS `日流量` FROM 5G GROUP BY `日期`,NCGI);
|
||||||
ALTER TABLE `日流量` ADD INDEX (`NCGI`);
|
ALTER TABLE `日流量` ADD INDEX (`NCGI`);
|
||||||
CREATE TEMPORARY TABLE IF NOT EXISTS `5G日均流量` (
|
CREATE TEMPORARY TABLE IF NOT EXISTS `5G日均流量` AS SELECT NCGI, AVG(`日流量`) AS `日均流量`, COUNT(`日期`) AS `流量有效天数` FROM `日流量` WHERE `日流量` > 0 GROUP BY NCGI;
|
||||||
SELECT NCGI, AVG(`日流量`) AS `日均流量`, COUNT(`日期`) AS `流量有效天数` FROM `日流量` WHERE `日流量` > 0 GROUP BY NCGI
|
|
||||||
);
|
|
||||||
|
|
||||||
CREATE TABLE IF NOT EXISTS `5G_结果表` (
|
CREATE TABLE IF NOT EXISTS `5G_结果表` (
|
||||||
SELECT NCGI, `CU小区配置名称`,
|
SELECT NCGI, `CU小区配置名称`,
|
||||||
|
|||||||
+10
-1
@@ -47,7 +47,16 @@ class DatabaseManager:
|
|||||||
|
|
||||||
@contextmanager
|
@contextmanager
|
||||||
def get_connection(self):
|
def get_connection(self):
|
||||||
"""获取 PyMySQL 连接(上下文管理器)"""
|
"""
|
||||||
|
获取 PyMySQL 连接(上下文管理器)
|
||||||
|
|
||||||
|
注意:此方法创建的是独立连接(非连接池),适用于:
|
||||||
|
- 需要在整个操作过程中保持同一 session 的场景
|
||||||
|
- 使用临时表(TEMPORARY TABLE)的场景(临时表是 session 级别的)
|
||||||
|
- 需要事务一致性的长时间操作
|
||||||
|
|
||||||
|
如果需要高性能的短连接操作,请使用 engine 属性(连接池)
|
||||||
|
"""
|
||||||
mysql = self.config.mysql
|
mysql = self.config.mysql
|
||||||
conn = pymysql.connect(
|
conn = pymysql.connect(
|
||||||
host=mysql.host,
|
host=mysql.host,
|
||||||
|
|||||||
+54
-8
@@ -676,8 +676,57 @@ class DataProcessor:
|
|||||||
speed = round(total_rows / elapsed) if elapsed > 0 else 0
|
speed = round(total_rows / elapsed) if elapsed > 0 else 0
|
||||||
self.logger.success(f"表 {table_name} 导入完成: {total_rows} 行, 耗时 {elapsed}s, 速度 {speed} 行/秒")
|
self.logger.success(f"表 {table_name} 导入完成: {total_rows} 行, 耗时 {elapsed}s, 速度 {speed} 行/秒")
|
||||||
|
|
||||||
|
@staticmethod
|
||||||
|
def parse_sql_script(sql_text: str) -> List[str]:
|
||||||
|
"""
|
||||||
|
解析 SQL 脚本,提取有效的 SQL 语句
|
||||||
|
|
||||||
|
改进的SQL分割逻辑:直接按分号分割,更可靠
|
||||||
|
这样可以确保所有以分号结尾的语句都被正确识别
|
||||||
|
|
||||||
|
Args:
|
||||||
|
sql_text: SQL 脚本文本内容
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
有效的 SQL 语句列表
|
||||||
|
"""
|
||||||
|
if not sql_text or not sql_text.strip():
|
||||||
|
return []
|
||||||
|
|
||||||
|
# 改进的SQL分割逻辑:直接按分号分割,更可靠
|
||||||
|
# 这样可以确保所有以分号结尾的语句都被正确识别
|
||||||
|
parts = sql_text.split(';')
|
||||||
|
valid_sqls = []
|
||||||
|
|
||||||
|
for part in parts:
|
||||||
|
# 处理多行语句:移除以#开头的注释行,但保留SQL语句
|
||||||
|
lines = []
|
||||||
|
for line in part.split('\n'):
|
||||||
|
line = line.strip()
|
||||||
|
# 跳过空行和整行注释
|
||||||
|
if line and not line.startswith('#'):
|
||||||
|
lines.append(line)
|
||||||
|
|
||||||
|
if lines:
|
||||||
|
# 合并多行语句,保留换行符(MySQL支持多行SQL)
|
||||||
|
cleaned_sql = '\n'.join(lines)
|
||||||
|
cleaned_sql = cleaned_sql.strip()
|
||||||
|
# 跳过空语句(可能只剩下注释)
|
||||||
|
if cleaned_sql:
|
||||||
|
valid_sqls.append(cleaned_sql)
|
||||||
|
|
||||||
|
return valid_sqls
|
||||||
|
|
||||||
def _execute_sql_script(self):
|
def _execute_sql_script(self):
|
||||||
"""执行 SQL 脚本"""
|
"""
|
||||||
|
执行 SQL 脚本
|
||||||
|
|
||||||
|
重要说明:
|
||||||
|
- 使用 get_connection() 获取独立连接(非连接池),确保整个脚本在同一 session 中执行
|
||||||
|
- 临时表(TEMPORARY TABLE)是 session 级别的,必须在同一连接中创建和使用
|
||||||
|
- 如果使用连接池,不同 SQL 语句可能分配到不同连接,导致临时表不可见
|
||||||
|
- 因此整个脚本必须在同一个连接中顺序执行,不能使用连接池
|
||||||
|
"""
|
||||||
if not SQL_SCRIPT.exists():
|
if not SQL_SCRIPT.exists():
|
||||||
self.logger.warning("SQL 脚本文件不存在,跳过")
|
self.logger.warning("SQL 脚本文件不存在,跳过")
|
||||||
return
|
return
|
||||||
@@ -691,13 +740,8 @@ class DataProcessor:
|
|||||||
self.logger.warning("SQL 脚本文件为空,跳过执行")
|
self.logger.warning("SQL 脚本文件为空,跳过执行")
|
||||||
return
|
return
|
||||||
|
|
||||||
sqls = sqlparse.split(sql_text)
|
# 使用抽离的解析函数
|
||||||
# 过滤掉空语句和注释
|
valid_sqls = self.parse_sql_script(sql_text)
|
||||||
valid_sqls = []
|
|
||||||
for sql in sqls:
|
|
||||||
sql = sql.strip()
|
|
||||||
if sql and not sql.startswith('#'):
|
|
||||||
valid_sqls.append(sql)
|
|
||||||
|
|
||||||
if not valid_sqls:
|
if not valid_sqls:
|
||||||
self.logger.warning("SQL 脚本中没有有效的 SQL 语句(可能全是注释或空行)")
|
self.logger.warning("SQL 脚本中没有有效的 SQL 语句(可能全是注释或空行)")
|
||||||
@@ -707,6 +751,8 @@ class DataProcessor:
|
|||||||
self.logger.info(f"共找到 {total} 条有效的 SQL 语句")
|
self.logger.info(f"共找到 {total} 条有效的 SQL 语句")
|
||||||
|
|
||||||
executed_count = 0
|
executed_count = 0
|
||||||
|
# 使用独立连接(非连接池),确保整个脚本在同一 session 中执行
|
||||||
|
# 这对于临时表(TEMPORARY TABLE)至关重要,因为临时表是 session 级别的
|
||||||
with self.db.get_connection() as conn:
|
with self.db.get_connection() as conn:
|
||||||
with conn.cursor() as cursor:
|
with conn.cursor() as cursor:
|
||||||
for i, sql in enumerate(valid_sqls, 1):
|
for i, sql in enumerate(valid_sqls, 1):
|
||||||
|
|||||||
Reference in New Issue
Block a user