如何用Python SAX Parser按指定记录数分块解析大型XML
大型XML分块解析批量写入方案
问题背景
我当前使用SAX解析器处理超大型XML文件,之前采用XMLGenerator拆分XML后再解析的方式,现在希望实现分块解析+批量写入的方案:每解析完成10000条记录,就将数据加载到CSV或DataFrame(最终写入数据库),避免重复解析分块内容,同时解决内存占用过高的问题。
原解析代码如下:
import xml.sax from collections import defaultdict import pandas as pd from sqlalchemy import create_engine class EmplyeeData(xml.sax.ContentHandler): def __init__(self): self.employee_dict = defaultdict(list) def startElement(self, tag, attr): self.tag = tag if tag == 'Emp1oyee': self.employee_dict['Emp1oyee_ID'].append(attr['id']) def characters(self, content): if content.strip(): if self.tag == 'FName': self.FName = content elif self.tag == 'LName': self.LName = content elif self.tag == 'City': self.City = content def endElement(self, tag): if tag == 'FName': self.employee_dict['FName'].append(self.FName) elif tag == 'LName': self.employee_dict['LName'].append(self.LName) elif tag == 'City': self.employee_dict['City'].append(self.City) handler = EmployeeData() parser = xml.sax.make_parser() parser.setContentHandler(handler) parser.parse('employee_xml.xml') EmployeeDetails = parser.getContentHandler() EmployeeData_out = EmployeeDetails.employee_dict df = pd.DataFrame(EmployeeData_out, columns=EmployeeData_out.keys()).set_index('Emp1oyee_ID') # for example I am writing in the csv file, actually I will be loading the data into database table. #I want to load the data incrementaly by parsing certain count of records at a time for example 10000 records at a time. ##con_eng = create_engine('oracle://[user]:[pass]@[host]:[port]/[schema]', echo=False) ##df.to_sql(name='target_table',con=con_eng ,if_exists = 'append', index=False) df.to_csv('employee_details.csv', sep=',', encoding='utf-8')
示例XML结构:
<?xml version="1.0" ?> <Emp1oyees> <Emp1oyee id=1> <FName>SAM</FName> <LName>MARK</LName> <City>NewJersy</City> </Emp1oyee> <Emp1oyee id=2> <FName>RAJ</FName> <LName>KAMAL</LName> <City>NewYork</City> </Emp1oyee> <Emp1oyee id=3> <FName>Brain</FName> <LName>wood</LName> <City>Buffalo</City> </Emp1oyee> ... <Emp1oyee id=1000000> <FName>Mark</FName> <LName>wood</LName> <City>NewJersy</City> </Emp1oyee> </Emp1oyees>
修改后的分块解析方案
通过改造SAX的ContentHandler,新增批次计数和批量写入逻辑,实现每解析指定数量的记录后立即写入存储,释放内存:
import xml.sax from collections import defaultdict import pandas as pd from sqlalchemy import create_engine class EmployeeData(xml.sax.ContentHandler): def __init__(self, batch_size=10000, output_file='employee_details.csv', db_engine=None): self.batch_size = batch_size # 每批处理的记录数 self.output_file = output_file # CSV输出路径 self.db_engine = db_engine # 数据库引擎(可选) self.current_employee = {} # 临时存储单条员工数据 self.batch_data = defaultdict(list) # 批次数据容器 self.record_count = 0 # 已解析记录计数 self.first_write = True # 标记是否首次写入CSV(控制表头) def startElement(self, tag, attr): self.tag = tag if tag == 'Emp1oyee': # 开始解析新员工节点时,初始化临时存储 self.current_employee = {'Emp1oyee_ID': attr['id']} def characters(self, content): content = content.strip() if not content: return # 收集当前字段的内容到单条员工数据 if self.tag in ['FName', 'LName', 'City']: self.current_employee[self.tag] = content def endElement(self, tag): if tag == 'Emp1oyee': # 员工节点解析完成,将数据加入批次 for key, value in self.current_employee.items(): self.batch_data[key].append(value) self.record_count += 1 # 达到批次大小,触发批量写入 if self.record_count % self.batch_size == 0: self._write_batch() # 清空批次数据,准备下一批解析 self.batch_data = defaultdict(list) def _write_batch(self): # 转换批次数据为DataFrame df = pd.DataFrame(self.batch_data).set_index('Emp1oyee_ID') # 写入CSV:首次写入带表头,后续追加 df.to_csv(self.output_file, sep=',', encoding='utf-8', mode='a', header=self.first_write) if self.first_write: self.first_write = False # 如果配置了数据库引擎,写入数据库(追加模式) if self.db_engine: df.to_sql(name='target_table', con=self.db_engine, if_exists='append', index=False) def endDocument(self): # 解析完成后,处理剩余不足一批的记录 if self.record_count % self.batch_size != 0: self._write_batch() # 使用示例 if __name__ == '__main__': # 可选:初始化数据库连接引擎 # con_eng = create_engine('oracle://[user]:[pass]@[host]:[port]/[schema]', echo=False) # 初始化处理器,设置批次大小为10000 handler = EmployeeData(batch_size=10000 # 若需写入数据库,取消下方注释 # db_engine=con_eng ) parser = xml.sax.make_parser() parser.setContentHandler(handler) parser.parse('employee_xml.xml')
核心改动说明
- 新增
batch_size参数,自由控制每批处理的记录数(默认10000) - 用
current_employee临时存储单条员工数据,避免原代码中字段顺序不匹配的问题 - 每解析完一个
<Emp1oyee>节点就计数,达到批次阈值时立即写入CSV/数据库 _write_batch方法自动处理CSV表头(仅首次写入),数据库写入采用追加模式- 解析结束时自动处理剩余的不足一批的记录,避免数据遗漏
- 每次写入后清空批次数据,避免内存持续累积
内容的提问来源于stack exchange,提问作者iavd
相关产品推荐
相关产品推荐

