基于FastAPI+Pydantic实现嵌套列表转指定JSON结构的优化问询
解决方案:高效转换嵌套列表至指定API响应结构
核心思路
针对大数据量场景,采用单次遍历+字典分组+集合去重的方案,避免加载全量数据到DataFrame等重型结构,严格控制内存占用。同时通过Pydantic保证输出结构的合法性,并行处理可根据任务类型(CPU/IO密集)选择线程池或进程池优化速度。
1. 定义Pydantic Schema
先明确输出结构的校验模型,确保转换后的数据符合API要求:
from pydantic import BaseModel class ItemDetail(BaseModel): seqno: str itemcd: str class DiscountDetail(BaseModel): seqno: str dname: str class TenderDetail(BaseModel): seqno: str type: str class TransactionResponse(BaseModel): id: str custname: str tranname: str itemdetail: list[ItemDetail] discountdetail: list[DiscountDetail] tenderdetail: list[TenderDetail]
2. 单线程高效转换实现
通过一次遍历完成分组、去重和规则处理,时间复杂度O(n),内存仅存储去重后的分组数据:
def transform_data(raw_data: list[list[str]]) -> list[TransactionResponse]: # 用(id, custname, tranname)作为分组键,确保同一交易的唯一性 groups = {} for row in raw_data: id_val, custname, tranname = row[0], row[1], row[2] item_seq, item_cd = row[3], row[4] disc_seq, disc_name = row[5], row[6] tender_seq, tender_type = row[7], row[8] key = (id_val, custname, tranname) if key not in groups: # 初始化分组,用集合存储明细元组实现自动去重 groups[key] = { 'item_set': set(), 'disc_set': set(), 'tender_set': set(), 'has_none_disc': False, 'has_none_tender': False } group = groups[key] # 处理商品明细:非None值加入集合 if item_seq != 'None' and item_cd != 'None': group['item_set'].add((item_seq, item_cd)) # 处理折扣明细:存在None则标记为空,否则加入集合 if disc_seq == 'None' or disc_name == 'None': group['has_none_disc'] = True else: group['disc_set'].add((disc_seq, disc_name)) # 处理支付明细:存在None则标记为空,否则加入集合 if tender_seq == 'None' or tender_type == 'None': group['has_none_tender'] = True else: group['tender_set'].add((tender_seq, tender_type)) # 转换分组数据为Pydantic模型 result = [] for (id_val, custname, tranname), group in groups.items(): item_list = [ItemDetail(seqno=seq, itemcd=cd) for seq, cd in group['item_set']] disc_list = [] if group['has_none_disc'] else [DiscountDetail(seqno=seq, dname=name) for seq, name in group['disc_set']] tender_list = [] if group['has_none_tender'] else [TenderDetail(seqno=seq, type=t) for seq, t in group['tender_set']] result.append(TransactionResponse( id=id_val, custname=custname, tranname=tranname, itemdetail=item_list, discountdetail=disc_list, tenderdetail=tender_list )) return result
3. 并行处理优化(超大数据量场景)
若数据量极大,可拆分批次用线程/进程池并行处理,再合并结果:
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor def process_batch(batch: list[list[str]]) -> dict: batch_groups = {} for row in batch: id_val, custname, tranname = row[0], row[1], row[2] item_seq, item_cd = row[3], row[4] disc_seq, disc_name = row[5], row[6] tender_seq, tender_type = row[7], row[8] key = (id_val, custname, tranname) if key not in batch_groups: batch_groups[key] = { 'item_set': set(), 'disc_set': set(), 'tender_set': set(), 'has_none_disc': False, 'has_none_tender': False } group = batch_groups[key] if item_seq != 'None' and item_cd != 'None': group['item_set'].add((item_seq, item_cd)) if disc_seq == 'None' or disc_name == 'None': group['has_none_disc'] = True else: group['disc_set'].add((disc_seq, disc_name)) if tender_seq == 'None' or tender_type == 'None': group['has_none_tender'] = True else: group['tender_set'].add((tender_seq, tender_type)) return batch_groups def merge_groups(all_groups: list[dict]) -> dict: merged = {} for batch_group in all_groups: for key, group_data in batch_group.items(): if key not in merged: merged[key] = { 'item_set': set(), 'disc_set': set(), 'tender_set': set(), 'has_none_disc': False, 'has_none_tender': False } merged_group = merged[key] merged_group['item_set'].update(group_data['item_set']) merged_group['disc_set'].update(group_data['disc_set']) merged_group['tender_set'].update(group_data['tender_set']) # 只要有一个批次标记None,最终分组就保持空明细 merged_group['has_none_disc'] = merged_group['has_none_disc'] or group_data['has_none_disc'] merged_group['has_none_tender'] = merged_group['has_none_tender'] or group_data['has_none_tender'] return merged def transform_data_parallel(raw_data: list[list[str]], num_workers: int = 4) -> list[TransactionResponse]: # 拆分数据为均等批次 batch_size = len(raw_data) // num_workers batches = [raw_data[i:i+batch_size] for i in range(0, len(raw_data), batch_size)] # CPU密集型任务建议用ProcessPoolExecutor,IO密集型用ThreadPoolExecutor with ThreadPoolExecutor(max_workers=num_workers) as executor: batch_results = list(executor.map(process_batch, batches)) merged_groups = merge_groups(batch_results) # 转换为Pydantic模型(同单线程逻辑) result = [] for (id_val, custname, tranname), group in merged_groups.items(): item_list = [ItemDetail(seqno=seq, itemcd=cd) for seq, cd in group['item_set']] disc_list = [] if group['has_none_disc'] else [DiscountDetail(seqno=seq, dname=name) for seq, name in group['disc_set']] tender_list = [] if group['has_none_tender'] else [TenderDetail(seqno=seq, type=t) for seq, t in group['tender_set']] result.append(TransactionResponse( id=id_val, custname=custname, tranname=tranname, itemdetail=item_list, discountdetail=disc_list, tenderdetail=tender_list )) return result
4. FastAPI集成示例
from fastapi import FastAPI app = FastAPI() # 单线程接口 @app.post("/transform", response_model=list[TransactionResponse]) async def transform_endpoint(raw_data: list[list[str]]): return transform_data(raw_data) # 并行接口(适合超大数据量) @app.post("/transform-parallel", response_model=list[TransactionResponse]) def transform_parallel_endpoint(raw_data: list[list[str]]): return transform_data_parallel(raw_data, num_workers=4)
关键优化点说明
- 内存控制:用字典+集合替代DataFrame,仅存储去重后的核心数据,避免全量数据加载。
- 去重效率:利用集合的O(1)去重特性,比列表遍历去重效率提升数倍。
- 并行选择:CPU密集型任务用进程池(规避GIL限制),IO密集型用线程池(开销更低)。
- 流式扩展:若数据来自文件/流,可改为逐行读取处理,进一步降低内存占用。
内容的提问来源于stack exchange,提问作者Pranathi
相关产品推荐
相关产品推荐

