一、进阶功能:多文件批量导入与去重优化
在上一篇基础教程中,我们实现了单文件Excel数据导入SQLite数据库的功能。本次教程将在此基础上扩展多文件批量导入、数据去重、异常重试等进阶功能,进一步提升数据处理的健壮性和效率。
(一)多文件批量导入实现
通过遍历指定文件夹下的所有Excel文件,实现批量导入功能:
import os
from openpyxl import load_workbook
import sqlite3
def read_excel_folder(folder_path):
# 遍历文件夹下的所有.xlsx文件
for filename in os.listdir(folder_path):
if filename.endswith('.xlsx'):
file_path = os.path.join(folder_path, filename)
yield from read_excel(file_path)
def read_excel(file_path):
wb = load_workbook(filename=file_path, read_only=True)
ws = wb.active
for row in ws.iter_rows(min_row=2, values_only=True):
if any(cell is not None for cell in row):
yield row
wb.close()
(二)数据去重策略实现
在导入数据前先检查数据库中是否已存在相同记录,避免重复导入:
def check_duplicate(conn, name, hire_date):
cursor = conn.cursor()
cursor.execute('''
SELECT COUNT(*) FROM employee
WHERE name = ? AND hire_date = ?
''', (name, hire_date))
return cursor.fetchone() > 0
def import_with_duplicate_check(data_generator, db_path='company.db'):
conn = sqlite3.connect(db_path)
cursor = conn.cursor()
insert_sql = '''
INSERT INTO employee (name, age, department, hire_date)
VALUES (?, ?, ?, ?)
'''
try:
count = 0
for row in data_generator:
# 以姓名和入职日期作为唯一标识检查重复
if not check_duplicate(conn, row, row):
cursor.execute(insert_sql, row)
count += 1
conn.commit()
print(f"成功插入{count}条新数据")
except Exception as e:
conn.rollback()
print(f"数据导入失败,错误信息:{e}")
finally:
conn.close()
(三)异常重试机制实现
针对可能出现的网络波动、文件损坏等异常情况,实现自动重试功能:
import time
from functools import wraps
def retry(max_retries=3, delay=2):
def decorator(func):
@wraps(func)
def wrapper(*args, **kwargs):
for attempt in range(max_retries):
try:
return func(*args, **kwargs)
except Exception as e:
print(f"第{attempt+1}次尝试失败,错误信息:{e}")
if attempt < max_retries - 1:
time.sleep(delay)
print(f"等待{delay}秒后重试...")
print(f"已尝试{max_retries}次,均失败,放弃操作")
return None
return wrapper
return decorator
@retry(max_retries=3, delay=2)
def import_with_retry(data_generator, db_path='company.db'):
import_with_duplicate_check(data_generator, db_path)
二、性能优化:提升大数据量导入效率
当处理超大Excel文件(百万级数据)时,需要进一步优化导入性能:
(一)批量插入优化
将数据按批次分组插入,减少数据库交互次数:
def batch_insert(data_generator, batch_size=10000):
batch = []
for row in data_generator:
batch.append(row)
if len(batch) >= batch_size:
yield batch
batch = []
if batch:
yield batch
def import_with_batch(data_generator, db_path='company.db'):
conn = sqlite3.connect(db_path)
cursor = conn.cursor()
insert_sql = '''
INSERT INTO employee (name, age, department, hire_date)
VALUES (?, ?, ?, ?)
'''
try:
total_count = 0
for batch in batch_insert(data_generator):
cursor.executemany(insert_sql, batch)
conn.commit()
total_count += len(batch)
print(f"已插入{total_count}条数据")
print(f"数据导入完成,共插入{total_count}条数据")
except Exception as e:
conn.rollback()
print(f"数据导入失败,错误信息:{e}")
finally:
conn.close()
(二)事务提交优化
将事务提交频率从每次插入改为每批次插入后提交,减少事务开销:
def import_with_transaction(data_generator, db_path='company.db', batch_size=10000):
conn = sqlite3.connect(db_path)
cursor = conn.cursor()
insert_sql = '''
INSERT INTO employee (name, age, department, hire_date)
VALUES (?, ?, ?, ?)
'''
try:
total_count = 0
batch = []
for row in data_generator:
batch.append(row)
if len(batch) >= batch_size:
cursor.executemany(insert_sql, batch)
conn.commit()
total_count += len(batch)
batch = []
print(f"已插入{total_count}条数据")
# 插入剩余数据
if batch:
cursor.executemany(insert_sql, batch)
conn.commit()
total_count += len(batch)
print(f"数据导入完成,共插入{total_count}条数据")
except Exception as e:
conn.rollback()
print(f"数据导入失败,错误信息:{e}")
finally:
conn.close()
(三)内存优化
使用生成器逐行读取数据,避免将整个文件加载到内存:
def read_large_excel(file_path):
wb = load_workbook(filename=file_path, read_only=True)
ws = wb.active
for row in ws.iter_rows(min_row=2, values_only=True):
if any(cell is not None for cell in row):
yield row
wb.close()
三、数据校验与清洗
在导入数据前进行校验和清洗,确保数据质量:
(一)数据校验
检查数据格式是否符合要求:
def validate_data(row):
name, age, department, hire_date = row
# 校验姓名非空
if not name or not isinstance(name, str):
return False, "姓名不能为空且必须为字符串"
# 校验年龄为整数且在合理范围内
if not isinstance(age, int) or age < 18 or age > 65:
return False, "年龄必须为18-65之间的整数"
# 校验部门非空
if not department or not isinstance(department, str):
return False, "部门不能为空且必须为字符串"
# 校验入职日期格式
try:
# 假设日期格式为YYYY-MM-DD
year, month, day = map(int, hire_date.split('-'))
if year < 2000 or year > 2026 or month < 1 or month > 12 or day < 1 or day > 31:
return False, "入职日期格式不正确"
except:
return False, "入职日期格式不正确"
return True, "数据校验通过"
(二)数据清洗
对数据进行清洗和转换:
def clean_data(row):
name, age, department, hire_date = row
# 去除姓名前后空格
name = name.strip()
# 转换年龄为整数
age = int(age)
# 去除部门前后空格并转换为大写
department = department.strip().upper()
# 统一日期格式为YYYY-MM-DD
try:
# 处理不同格式的日期
if '/' in hire_date:
month, day, year = hire_date.split('/')
hire_date = f"{year}-{month.zfill(2)}-{day.zfill(2)}"
elif '.' in hire_date:
day, month, year = hire_date.split('.')
hire_date = f"{year}-{month.zfill(2)}-{day.zfill(2)}"
except:
pass
return (name, age, department, hire_date)
(三)完整的数据处理流程
将校验和清洗整合到导入流程中:
def process_data(data_generator):
for row in data_generator:
is_valid, message = validate_data(row)
if not is_valid:
print(f"数据校验失败:{message},跳过该行数据")
continue
cleaned_row = clean_data(row)
yield cleaned_row
四、日志记录与监控
添加日志记录功能,便于监控导入过程和排查问题:
(一)日志记录实现
import logging
def setup_logging():
logging.basicConfig(
level=logging.INFO,
format='%(asctime)s - %(levelname)s - %(message)s',
handlers=[
logging.FileHandler('import_log.log'),
logging.StreamHandler()
]
)
def import_with_logging(data_generator, db_path='company.db'):
setup_logging()
logger = logging.getLogger(__name__)
conn = sqlite3.connect(db_path)
cursor = conn.cursor()
insert_sql = '''
INSERT INTO employee (name, age, department, hire_date)
VALUES (?, ?, ?, ?)
'''
try:
total_count = 0
for row in data_generator:
cursor.execute(insert_sql, row)
total_count += 1
conn.commit()
logger.info(f"数据导入完成,共插入{total_count}条数据")
except Exception as e:
conn.rollback()
logger.error(f"数据导入失败,错误信息:{e}")
finally:
conn.close()
(二)进度监控实现
import tqdm
def import_with_progress(data_generator, db_path='company.db'):
conn = sqlite3.connect(db_path)
cursor = conn.cursor()
insert_sql = '''
INSERT INTO employee (name, age, department, hire_date)
VALUES (?, ?, ?, ?)
'''
try:
# 将生成器转换为列表以获取总条数
data_list = list(data_generator)
total_count = len(data_list)
with tqdm.tqdm(total=total_count, desc="导入进度") as pbar:
for row in data_list:
cursor.execute(insert_sql, row)
pbar.update(1)
conn.commit()
print(f"数据导入完成,共插入{total_count}条数据")
except Exception as e:
conn.rollback()
print(f"数据导入失败,错误信息:{e}")
finally:
conn.close()
五、完整的批量导入解决方案
将以上功能整合,形成完整的批量导入解决方案:
def main():
# 创建数据库表
create_table()
# 指定Excel文件所在文件夹
excel_folder = 'excel_data'
# 读取文件夹下的所有Excel文件数据
data_gen = read_excel_folder(excel_folder)
# 数据校验与清洗
processed_data = process_data(data_gen)
# 批量导入数据到数据库(带进度监控)
import_with_progress(processed_data)
if __name__ == '__main__':
main()
六、总结与扩展
通过本次教程,我们实现了一个功能完善、性能优化的Excel数据批量导入SQLite数据库的解决方案,具备以下特点:
多文件批量导入:支持遍历指定文件夹下的所有Excel文件
数据去重:避免重复导入相同记录
异常重试:提高系统健壮性
性能优化:支持百万级数据高效导入
数据校验与清洗:确保数据质量
日志记录与进度监控:便于监控和排查问题
后续可根据实际需求进一步扩展功能,例如:
支持增量导入(仅导入新增数据)
支持数据更新(根据主键更新已有记录)
支持多线程/多进程导入,进一步提升速度
开发图形界面,方便非技术人员使用
支持更多数据库类型(如MySQL、PostgreSQL等)