基于Python的数据对账:单机器运行内存不足问题求解
跨库对账内存不足的分块比对解决方案
针对单机器拉取全量数据对账触发out of memory error的问题,核心解决思路是基于唯一标识(如示例中的id)分批次拉取两端对应区间的数据,逐块完成比对,避免一次性加载全量数据占用内存。
实现步骤
- 先查询源端/目标端数据的
id范围(最大、最小值),确定分块步长(比如每块10000条) - 循环遍历每个
id区间,分别拉取源端和目标端该区间的数据集 - 对每块数据执行对账逻辑,记录差异
- 所有块处理完成后,汇总差异结果
修改后的代码实现
smartdq_logging('Log', gv_log_level, f"connecting to source db...") cl_connect_sourcedb = lib.environment_connection.DBConnect() cl_connect_sourcedb.connect_with_environment_id(environmentIdSource, "SMARTDQ", cl_connect_smartdq) smartdq_logging('Log', gv_log_level, f"connecting to target db...") cl_connect_target = lib.environment_connection.DBConnect() cl_connect_target.connect_with_environment_id(environmentIdTarget, "SMARTDQ", cl_connect_smartdq) # 1. 获取数据ID范围(以源端为准,假设两端ID范围一致) get_range_query = "SELECT MIN(id) AS min_id, MAX(id) AS max_id FROM (SELECT LEVEL AS id FROM DUAL CONNECT BY LEVEL <= 1000000)" df_range = db_func.read_data_to_dataframe(get_range_query, cl_connect_sourcedb, "original") min_id = df_range['min_id'].iloc[0] max_id = df_range['max_id'].iloc[0] # 2. 设置分块步长,根据内存情况调整 chunk_size = 10000 # 存储所有差异结果 diff_results = [] print('start recon...') # 3. 循环分块拉取并比对 for start_id in range(min_id, max_id + 1, chunk_size): end_id = min(start_id + chunk_size - 1, max_id) # 构造分块查询语句 source_chunk_query = f"""SELECT id, value FROM ( SELECT LEVEL AS id, ROUND(DBMS_RANDOM.VALUE(0, 99999) / 100000, 5) AS value FROM DUAL CONNECT BY LEVEL <= 1000000 ) WHERE id BETWEEN {start_id} AND {end_id} ORDER BY id""" target_chunk_query = f"""SELECT id, value FROM ( SELECT LEVEL AS id, ROUND(DBMS_RANDOM.VALUE(0, 99999) / 100000, 5) AS value FROM DUAL CONNECT BY LEVEL <= 1000000 ) WHERE id BETWEEN {start_id} AND {end_id} ORDER BY id""" # 拉取当前块的源端和目标端数据 df_source_chunk = db_func.read_data_to_dataframe(source_chunk_query, cl_connect_sourcedb, "original") df_target_chunk = db_func.read_data_to_dataframe(target_chunk_query, cl_connect_target, "original") # 4. 执行对账逻辑:合并数据找差异 df_merged = df_source_chunk.merge(df_target_chunk, on='id', how='outer', suffixes=('_source', '_target'), indicator=True) # 筛选出差异行:仅存在源端、仅存在目标端、值不一致 diffs = df_merged[ (df_merged['_merge'] != 'both') | (df_merged['value_source'] != df_merged['value_target']) ].copy() diffs['chunk_range'] = f"{start_id}-{end_id}" # 保存当前块的差异 diff_results.append(diffs) smartdq_logging('Log', gv_log_level, f"Processed chunk {start_id}-{end_id}, found {len(diffs)} differences") # 5. 汇总所有差异 if diff_results: final_diff = pd.concat(diff_results, ignore_index=True) # 可以将差异写入文件或数据库 final_diff.to_csv('reconciliation_diff.csv', index=False) smartdq_logging('Log', gv_log_level, f"Reconciliation completed. Total differences: {len(final_diff)}") else: smartdq_logging('Log', gv_log_level, "Reconciliation completed. No differences found.")
关键说明
- 分块依据:必须基于两端共有的唯一键(如
id),确保每块拉取的是对应的数据,保证比对的准确性 - 步长调整:
chunk_size可根据机器内存情况调整,内存小则调小步长 - 对账逻辑:示例中用
merge+indicator筛选差异,可根据实际需求修改(比如更严格的数值精度比对) - 资源释放:每块处理完成后,当前块的DataFrame会被自动回收,不会占用过多内存
内容的提问来源于stack exchange,提问作者İsmail Uzun
相关产品推荐
相关产品推荐

