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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 22:40:01