You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.17 17:10:29