数据来源五花八门是几乎所有数据团队的常态:业务系统导出的是Excel,第三方接口返回的是JSON,历史数据存的是CSV,大数据平台用的是Parquet。如果每次都靠人工打开文件逐个调整格式,不仅效率低下,还极易出错。编写一套通用的数据转换脚本,把所有异构数据源统一成标准格式,是解决这个问题的最佳途径。本文将从格式差异分析、脚本设计思路、完整代码实现和性能优化四个方面,详细讲解如何构建自己的数据转换工具。

一、常见数据格式的差异与转换难点
在动手写脚本之前,必须先弄清楚各种格式的特点。CSV是纯文本格式,简单通用但缺乏类型信息,所有数据都以字符串存储,日期和数字需要手动解析;JSON是层级结构,天然支持嵌套字段,转换成表格时需要展平处理;Excel支持多工作表,可能包含合并单元格、公式等复杂内容;Parquet是列式存储格式,读取速度快、占用空间小,但不适合人工查看。
这些差异带来的转换难点主要集中在三个方面。第一是字段命名不一致,比如同一个概念有的数据源叫user_id,有的叫userId,还有的叫用户ID,必须建立字段映射表。第二是数据类型混乱,同一个日期字段可能是2024-01-01、1704067200(时间戳)、2024/1/1等多种形式。第三是脏数据问题,包括缺失值、空字符串、重复行、异常字符等,如果不在转换阶段清洗,会污染下游的数据集。
建议在项目初期就定义一份统一的数据规范文档,明确目标格式的字段名、数据类型、编码方式(推荐UTF-8)和缺失值表示方法。这份文档就是转换脚本的编写依据,后续所有数据源都向这份规范看齐。
二、转换脚本的设计思路与核心代码
一个健壮的转换脚本应该具备三个特性:可配置、可批处理、可容错。可配置指字段映射规则和清洗规则写在外部配置中,新增数据源时只需修改配置而非代码;可批处理指脚本能遍历整个目录自动识别文件类型并批量转换;可容错指单个文件转换失败时记录日志并继续处理其他文件,而不是让整个任务中断。
下面是基于Python和Pandas实现的核心转换模块,它实现了文件类型自动识别、字段映射、类型规范化和缺失值处理:
import pandas as pd
import json
import os
from pathlib import Path
# 字段映射表:统一目标字段名为 user_id, order_time, amount
FIELD_MAPPING = {
"userId": "user_id",
"用户ID": "user_id",
"created_at": "order_time",
"下单时间": "order_time",
"total_price": "amount",
"总价": "amount",
}
def load_data(file_path: str) -> pd.DataFrame:
"""根据扩展名自动选择读取方式"""
suffix = Path(file_path).suffix.lower()
if suffix == ".csv":
return pd.read_csv(file_path, encoding="utf-8-sig")
elif suffix in (".xls", ".xlsx"):
return pd.read_excel(file_path)
elif suffix == ".json":
# 展平嵌套的JSON结构
with open(file_path, "r", encoding="utf-8") as f:
records = json.load(f)
return pd.json_normalize(records)
elif suffix == ".parquet":
return pd.read_parquet(file_path)
else:
raise ValueError(f"不支持的文件格式: {suffix}")
def normalize_columns(df: pd.DataFrame) -> pd.DataFrame:
"""统一列名并规范化数据类型"""
df = df.rename(columns=FIELD_MAPPING)
# 只保留映射表中定义的目标字段
keep_cols = [c for c in set(FIELD_MAPPING.values()) if c in df.columns]
df = df[keep_cols]
if "order_time" in df.columns:
# 自动解析多种日期格式,errors参数保证异常值变成NaT而不报错
df["order_time"] = pd.to_datetime(df["order_time"], errors="coerce")
if "amount" in df.columns:
df["amount"] = pd.to_numeric(df["amount"], errors="coerce")
return df
def clean_data(df: pd.DataFrame) -> pd.DataFrame:
"""基础数据清洗:去重、处理缺失值"""
df = df.drop_duplicates()
# 金额缺失的行直接丢弃,时间缺失则填充为空
df = df.dropna(subset=["amount"])
return df.reset_index(drop=True)
这段代码的关键点在于errors="coerce"参数的使用,它让类型转换失败的数据自动变成缺失值而不是抛出异常,这对处理来源不明的脏数据非常重要。另外json_normalize函数可以自动把嵌套的JSON对象展平成表格结构,比如把user.info.id展开成一列,省去了手写解析逻辑的工作。
三、批量自动化处理与日志记录
单个文件的转换只是基础,实际场景中往往需要一次性处理几百个文件。这时需要一个调度函数遍历目录、调用转换模块并汇总结果。完善的日志记录能让你清楚知道哪些文件转换成功、哪些失败以及失败原因,方便后续排查。
import logging
logging.basicConfig(
filename="convert.log",
level=logging.INFO,
format="%(asctime)s - %(levelname)s - %(message)s",
)
def batch_convert(input_dir: str, output_dir: str):
"""批量转换目录下的所有数据文件并输出为统一的Parquet格式"""
os.makedirs(output_dir, exist_ok=True)
success, failed = 0, 0
for file in Path(input_dir).iterdir():
if file.suffix.lower() not in (".csv", ".json", ".xls", ".xlsx"):
continue
try:
df = load_data(str(file))
df = clean_data(normalize_columns(df))
out_path = Path(output_dir) / (file.stem + ".parquet")
df.to_parquet(out_path, index=False)
success += 1
logging.info(f"转换成功: {file.name}, 共 {len(df)} 行")
except Exception as e:
failed += 1
logging.error(f"转换失败: {file.name}, 原因: {e}")
print(f"处理完成,成功 {success} 个,失败 {failed} 个,详情见 convert.log")
if __name__ == "__main__":
batch_convert("./raw_data", "./clean_data")
输出格式选择Parquet是有原因的:它保留了完整的数据类型信息,读取速度比CSV快数倍,文件体积通常只有CSV的一半左右,非常适合作为统一的目标格式。如果下游工具不支持Parquet,也可以按需改为CSV输出,只需把to_parquet换成to_csv即可。日志文件中记录了每个文件的处理结果,失败的文件会附上具体原因,比如编码错误、JSON解析失败等,据此可以针对性修复源数据。
四、大数据量场景下的性能优化技巧
当单个文件达到几百MB甚至几个GB时,直接用read_csv一次性读入内存可能会撑爆机器。Pandas提供了chunksize参数支持分块读取,每次只加载一部分数据到内存,处理完再读取下一块,内存占用可以控制在很低的水平。
def convert_large_csv(file_path: str, output_path: str):
"""分块读取大CSV文件并逐块写出"""
chunk_size = 100000 # 每次读取10万行
writer = None
for chunk in pd.read_csv(file_path, chunksize=chunk_size, encoding="utf-8-sig"):
chunk = clean_data(normalize_columns(chunk))
if writer is None:
# 第一块数据创建文件并写入表头
chunk.to_parquet(output_path, index=False)
import pyarrow.parquet as pq
table = pq.read_table(output_path)
writer = pq.ParquetWriter(
output_path, table.schema, mode="w"
)
writer.write_table(table)
else:
import pyarrow as pq
writer.write_table(pq.Table.from_pandas(chunk, preserve_index=False))
if writer:
writer.close()
除了分块读取,还有几个实用优化手段。只读取需要的列可以大幅减少内存占用,使用usecols参数指定字段列表即可;对于CSV文件,显式指定dtype参数能避免Pandas反复推断类型带来的开销;如果机器配置允许,可以使用Polars或Dask等替代库,它们在多核并行和内存管理上比Pandas更出色,处理千万行级数据的速度可以提升数倍。
最后提醒一点,转换脚本写好后不要直接投入生产环境,先用一批有代表性的样本数据做小规模测试,验证字段映射、类型转换和清洗逻辑都符合预期,再放量执行。同时把脚本纳入版本管理,每次数据规范变更时同步更新字段映射表,这样转换流程才能长期稳定地服务业务。