<think>我们之前已经编写了一个完整的Python程序,用于连接达梦数据库并扩充表数据。现在用户特别询问如何在main函数中传参来实现这个功能。我们将重点解释参数传递部分,并展示如何将参数传递给数据库连接和扩展函数。
在之前的程序中,我们使用了argparse模块来解析命令行参数。下面详细说明参数传递的实现:
1. 定义命令行参数解析器
2. 添加所需的参数(用户、密码、服务器、端口、模式、表名、目标大小)
3. 解析参数并传递给连接函数和扩展函数
以下是具体的代码部分(与之前程序一致,但重点注释参数传递部分):
</think>
# Python连接达梦数据库实现数据扩展的方法
下面我将详细解释如何在Python程序中通过main函数传参连接达梦(DM)数据库,并根据源表数据随机复制行数据扩充到指定大小(如1MB到1GB),同时保证主键ID不重复。
## 完整实现代码
```python
import dmPython
import random
import argparse
import time
import math
from tqdm import tqdm
def connect_dm_db(user, password, server, port, schema):
"""连接达梦数据库"""
try:
conn = dmPython.connect(
user=user,
password=password,
server=server,
port=port,
schema=schema,
autoCommit=False
)
cursor = conn.cursor()
print(f"✅ 成功连接达梦数据库 [模式: {schema}]")
return conn, cursor
except dmPython.Error as e:
print(f"❌ 数据库连接失败: {e}")
return None, None
def get_table_info(cursor, table_name):
"""获取表结构信息"""
try:
# 获取主键列
cursor.execute(f"""
SELECT cols.COLUMN_NAME
FROM ALL_CONSTRAINTS cons
JOIN ALL_CONS_COLUMNS cols
ON cons.CONSTRAINT_NAME = cols.CONSTRAINT_NAME
WHERE cons.TABLE_NAME = '{table_name.upper()}'
AND cons.CONSTRAINT_TYPE = 'P'
""")
pk_columns = [row[0] for row in cursor.fetchall()]
# 获取所有列信息
cursor.execute(f"""
SELECT COLUMN_NAME, DATA_TYPE, IDENTITY
FROM ALL_TAB_COLUMNS
WHERE TABLE_NAME = '{table_name.upper()}'
""")
columns = []
identity_column = None
for row in cursor.fetchall():
col_name, data_type, identity = row
columns.append(col_name)
if identity == 'YES':
identity_column = col_name
# 获取当前行数
cursor.execute(f"SELECT COUNT(*) FROM {table_name}")
row_count = cursor.fetchone()[0]
return {
"table_name": table_name,
"primary_keys": pk_columns,
"columns": columns,
"identity_column": identity_column,
"current_row_count": row_count
}
except dmPython.Error as e:
print(f"❌ 获取表信息失败: {e}")
return None
def estimate_row_size(cursor, table_name):
"""估算单行数据大小"""
try:
cursor.execute(f"SELECT * FROM {table_name} LIMIT 1")
sample_row = cursor.fetchone()
return len(str(sample_row)) if sample_row else 100 # 默认100字节
except:
return 100 # 保守估计
def expand_table_data(conn, cursor, table_info, target_size_gb):
"""扩展表数据到目标大小"""
table_name = table_info["table_name"]
identity_col = table_info["identity_column"]
pk_columns = table_info["primary_keys"]
# 估算当前表大小
row_size = estimate_row_size(cursor, table_name)
current_size_bytes = table_info["current_row_count"] * row_size
target_size_bytes = target_size_gb * 1024**3 # GB转换为字节
if current_size_bytes >= target_size_bytes:
print(f"⚠️ 表已超过目标大小: {current_size_bytes/(1024**2):.2f}MB > {target_size_bytes/(1024**2):.2f}MB")
return
# 计算需要添加的行数
rows_needed = math.ceil((target_size_bytes - current_size_bytes) / row_size)
print(f"📊 当前表大小: {current_size_bytes/(1024**2):.2f}MB")
print(f"🎯 目标表大小: {target_size_bytes/(1024**2):.2f}MB")
print(f"🔢 需要添加行数: {rows_needed:,}")
# 获取所有列名(排除自增主键)
insert_columns = [col for col in table_info["columns"] if col != identity_col]
insert_columns_sql = ", ".join(insert_columns)
placeholders = ", ".join(["?"] * len(insert_columns))
# 准备插入SQL
insert_sql = f"""
INSERT INTO {table_name} ({insert_columns_sql})
VALUES ({placeholders})
"""
# 批量插入参数
batch_size = 1000
batches = math.ceil(rows_needed / batch_size)
inserted_rows = 0
print("🚀 开始数据扩展...")
start_time = time.time()
# 使用进度条
with tqdm(total=rows_needed, desc="📤 插入进度", unit="行") as pbar:
for _ in range(batches):
# 随机选择一行作为模板
cursor.execute(f"SELECT {insert_columns_sql} FROM {table_name} ORDER BY RAND() LIMIT 1")
template_row = cursor.fetchone()
if not template_row:
print("❌ 无法获取模板行数据")
break
batch_data = []
for _ in range(min(batch_size, rows_needed - inserted_rows)):
# 复制模板行数据
new_row = list(template_row)
# 处理主键列(如果存在且不是自增列)
for i, col in enumerate(insert_columns):
if col in pk_columns and col != identity_col:
# 生成唯一主键值
cursor.execute(f"SELECT MAX({col}) FROM {table_name}")
max_id = cursor.fetchone()[0] or 0
new_row[i] = max_id + random.randint(1, 10000)
batch_data.append(tuple(new_row))
# 批量插入
try:
cursor.executemany(insert_sql, batch_data)
conn.commit()
inserted_rows += len(batch_data)
pbar.update(len(batch_data))
except dmPython.Error as e:
print(f"❌ 插入失败: {e}")
conn.rollback()
break
# 验证最终大小
cursor.execute(f"SELECT COUNT(*) FROM {table_name}")
final_row_count = cursor.fetchone()[0]
final_size_bytes = final_row_count * row_size
elapsed_time = time.time() - start_time
print("\n✅ 数据扩展完成!")
print(f"⏱️ 耗时: {elapsed_time:.2f}秒")
print(f"📈 最终行数: {final_row_count:,}")
print(f"📏 最终大小: {final_size_bytes/(1024**2):.2f}MB")
print(f"🚀 插入速度: {inserted_rows/elapsed_time:.2f} 行/秒")
def main():
"""主函数:通过命令行参数控制数据扩展"""
# 创建参数解析器
parser = argparse.ArgumentParser(
description='达梦数据库表数据扩展工具',
formatter_class=argparse.ArgumentDefaultsHelpFormatter
)
# 添加命令行参数
parser.add_argument('--user', required=True, help='数据库用户名')
parser.add_argument('--password', required=True, help='数据库密码')
parser.add_argument('--server', default='localhost', help='数据库服务器地址')
parser.add_argument('--port', type=int, default=5236, help='数据库端口')
parser.add_argument('--schema', required=True, help='数据库模式')
parser.add_argument('--table', required=True, help='目标表名')
parser.add_argument('--target', type=float, default=1.0, help='目标大小(GB)')
# 解析参数
args = parser.parse_args()
# 连接数据库
conn, cursor = connect_dm_db(
user=args.user,
password=args.password,
server=args.server,
port=args.port,
schema=args.schema
)
if not conn or not cursor:
return
try:
# 获取表信息
table_info = get_table_info(cursor, args.table)
if not table_info:
print(f"❌ 无法获取表 '{args.table}' 的信息")
return
print("\n📋 表结构信息:")
print(f" - 表名: {table_info['table_name']}")
print(f" - 主键列: {', '.join(table_info['primary_keys']) or '无'}")
print(f" - 自增列: {table_info['identity_column'] or '无'}")
print(f" - 总列数: {len(table_info['columns'])}")
print(f" - 当前行数: {table_info['current_row_count']:,}")
# 扩展表数据
expand_table_data(conn, cursor, table_info, args.target)
finally:
# 关闭连接
cursor.close()
conn.close()
print("🔌 数据库连接已关闭")
if __name__ == "__main__":
main()
```
## 使用方法
### 1. 安装依赖
```bash
pip install dmPython tqdm
```
### 2. 运行脚本
```bash
python dm_data_expander.py \
--user SYSDBA \
--password SYSDBA \
--server 127.0.0.1 \
--port 5236 \
--schema YOUR_SCHEMA \
--table YOUR_TABLE \
--target 1.0
```
## 参数传递机制详解
### 1. 参数解析器创建
```python
parser = argparse.ArgumentParser(
description='达梦数据库表数据扩展工具',
formatter_class=argparse.ArgumentDefaultsHelpFormatter
)
```
### 2. 添加命令行参数
```python
# 必需参数
parser.add_argument('--user', required=True, help='数据库用户名')
parser.add_argument('--password', required=True, help='数据库密码')
parser.add_argument('--schema', required=True, help='数据库模式')
parser.add_argument('--table', required=True, help='目标表名')
# 可选参数(带默认值)
parser.add_argument('--server', default='localhost', help='数据库服务器地址')
parser.add_argument('--port', type=int, default=5236, help='数据库端口')
parser.add_argument('--target', type=float, default=1.0, help='目标大小(GB)')
```
### 3. 参数解析与传递
```python
# 解析参数
args = parser.parse_args()
# 将参数传递给数据库连接函数
conn, cursor = connect_dm_db(
user=args.user,
password=args.password,
server=args.server,
port=args.port,
schema=args.schema
)
# 将参数传递给数据扩展函数
expand_table_data(conn, cursor, table_info, args.target)
```
## 关键技术点
### 1. 主键唯一性保证
```python
# 处理非自增主键
if col in pk_columns and col != identity_col:
cursor.execute(f"SELECT MAX({col}) FROM {table_name}")
max_id = cursor.fetchone()[0] or 0
new_row[i] = max_id + random.randint(1, 10000)
```
### 2. 数据大小计算
数据大小计算使用公式:
$$ \text{rows\_needed} = \left\lceil \frac{\text{target\_size} - \text{current\_size}}{\text{row\_size}} \right\rceil $$
### 3. 高效批量插入
```python
# 批量插入1000行/次
batch_size = 1000
cursor.executemany(insert_sql, batch_data)
```
### 4. 随机行选择
```python
# 随机选择一行作为模板
cursor.execute(f"SELECT {insert_columns_sql} FROM {table_name} ORDER BY RAND() LIMIT 1")
```
## 性能优化建议
1. **临时禁用索引**:
```python
# 插入前禁用索引
cursor.execute(f"ALTER INDEX {index_name} INVISIBLE")
# 插入后恢复索引
cursor.execute(f"ALTER INDEX {index_name} VISIBLE")
```
2. **增加批量大小**:
```python
# 根据内存调整批量大小
batch_size = 5000 # 默认1000
```
3. **使用直接路径插入**:
```python
# 达梦特有优化
cursor.execute("SET DIRECT_INSERT ON")
```
4. **并行处理**:
```python
# 启用并行插入
cursor.execute("ALTER SESSION ENABLE PARALLEL DML")
```
## 扩展功能
1. **断点续传**:保存已插入行数,中断后可继续
2. **多线程插入**:使用线程池加速大规模数据插入
3. **数据分布分析**:验证插入后数据分布是否均匀
4. **自动索引维护**:插入完成后自动重建索引