如何用SQLAlchemy Core高效批量插数?解决循环变量与函数调用问题
SQLAlchemy Core 批量插入优化方案
针对ORM批量插入数十万条数据速度慢的问题,以下是完整的SQLAlchemy Core实现方案,解决循环内生成init_time和多级部门编码的问题,同时保持分批插入的内存友好性:
核心优化点
- 将单条数据生成逻辑封装为独立函数,避免代码臃肿
- 移除不必要的
datetime.strptime转换,直接使用datetime对象提升效率 - 沿用ORM版本的分批插入策略(每10000条提交一次),避免内存溢出
- 直接调用
TableData.__table__.insert()执行批量插入,最大化Core的性能优势
完整实现代码
import random from datetime import datetime, timedelta # 假设get_random_node_code、core是你已实现/导入的模块 # from your_module import get_random_node_code, core def _generate_single_data(counter, username_list, product_list, organization_tree): # 生成随机日期,直接返回datetime对象,无需转字符串再解析 init_time = datetime.now() - timedelta(days=random.randrange(365)) # 生成多级部门编码 level1 = get_random_node_code(organization_tree) level2 = get_random_node_code(organization_tree[level1]) level3 = get_random_node_code(organization_tree[level1][level2]) level4 = get_random_node_code(organization_tree[level1][level2][level3]) return { 'pkey': random.choice(product_list), 'user_name': random.choice(username_list), 'qa_version': init_time.strftime("%Y-%m-%d"), 'updated_time': init_time, 'created_time': init_time, 'build_exec_time': init_time, 'tests_failed': 0, 'tests_skipped': 0, 'tests_total': 0, 'is_incoco': random.choice(['true', 'false']), 'jacoco_branch_coverage_rate': round(random.uniform(0, 48), 4), 'jacoco_line_coverage_rate': round(random.uniform(0, 48), 4), 'blocker_issues': random.randint(0, 6), 'major_issues': random.randint(0, 6), 'critical_issues': random.randint(0, 6), 'project_key': f"project_key_{counter}", 'applicant_sn': random.choice(username_list), 'jira_version': random.choice(['2.10.24', '7.5.13', '4.15.37', '9.8.5', '3.19.42']), 'git_urls': f"https://gitlab.demo.com/git-number{counter}/git", 'quality_scope': random.choice(['overall', 'newcode']), 'dev_lang': random.choice(['java', 'python', 'golang']), 'source': random.choice([0, 1]), 'unit_username': random.choice(username_list), 'pipeline_type': 'cd_pipeline', 'pname': random.choice(product_list), 'jenkins_template': random.choice(['dailyCiJava', 'dailyCiPython', 'dailyCiGolang']), 'jacoco_lines_covered': random.randint(100, 1000), 'LEVEL_1_DEVOPS_DEPT_CODE': level1, 'LEVEL_2_DEVOPS_DEPT_CODE': level2, 'LEVEL_3_DEVOPS_DEPT_CODE': level3, 'LEVEL_4_DEVOPS_DEPT_CODE': level4 } def create_data(sess, engine, count=50, username_list=None, product_list=None, project_list=None, dept_list=None, organization_tree=None): # 检查表是否存在(按需保留) core.table_action.create_table_if_not_exists(DwdDeployBranchData, engine) batch_size = 10000 insert_stmt = TableData.__table__.insert() for batch_start in range(1, count, batch_size): batch_end = min(batch_start + batch_size, count) # 生成当前批次的数据字典列表 batch_data = [ _generate_single_data(counter, username_list, product_list, organization_tree) for counter in range(batch_start, batch_end) ] # 执行批量插入 sess.execute(insert_stmt, batch_data) sess.commit() # 显式清空批次数据,释放内存 del batch_data
注意事项
- 确保传入的
username_list、product_list、organization_tree不为空,否则random.choice会抛出异常 - 可根据数据库性能调整
batch_size(MySQL通常建议1000-10000条/批) - 若使用异步引擎,需改用
await sess.execute()的异步语法 - 移除
datetime.strptime转换是因为datetime.now() - timedelta已返回datetime对象,无需额外解析字符串,能节省CPU开销
内容的提问来源于stack exchange,提问作者boxuan666
相关产品推荐
相关产品推荐

