数据处理端到端实践
这一页不是“跑一个 CSV 报表”的博客示例,而是按商业项目的方式,把一个数据处理脚本从需求、字段口径、目录结构、清洗函数、异常数据、日志、测试、运行参数到排查链路完整串起来。
你可以把它理解成一个小型离线任务模板:业务每天上传订单或资产文件,系统自动清洗、校验、聚合并输出结果。如果结果不对,必须能知道错在哪一层,而不是只看到一个失败堆栈。
学习目标
学完本页,你应该能做到:
- 明确数据处理任务的输入、输出和业务口径。
- 设计一个可维护的数据脚本目录,而不是把所有逻辑写进一个文件。
- 解释读取、校验、转换、异常拆分、去重、聚合、导出的完整流程。
- 写出可运行的 Pandas 商业 Demo。
- 给清洗逻辑写测试,保证口径变更时能及时发现问题。
- 使用命令行参数、日志和处理报告,让脚本适合定时任务。
- 排查报表金额对不上、异常行过多、内存过高、重复订单等生产问题。
- 面试时能说清“数据处理项目怎么做才可靠”。
商业需求
假设电商或医疗资产平台每天会产生一份订单文件 orders.csv,字段如下:
| 字段 | 含义 | 是否核心 | 处理要求 |
|---|---|---|---|
order_id | 订单号 | 是 | 不能为空,按字符串处理,重复时保留更新时间最新的一条 |
city | 城市 | 是 | 不能为空,去掉前后空格 |
amount | 订单金额 | 是 | 转数字,不能转换进入异常文件 |
status | 订单状态 | 是 | 只能是 paid、refund、cancelled |
created_at | 创建时间 | 是 | 转日期,不能转换进入异常文件 |
updated_at | 更新时间 | 是 | 转日期,重复订单按它保留最新 |
channel | 渠道 | 否 | 缺失时标记为 unknown |
目标输出:
clean_orders.csv:清洗后的有效订单。invalid_orders.csv:异常订单,并说明异常原因。city_report.csv:按城市统计已支付订单数、总金额、平均金额。run_summary.json:本次任务行数、异常数、重复数、输出路径和处理状态。
为什么不能只写一个简单脚本
简单写法通常是这样:
import pandas as pd
df = pd.read_csv("orders.csv")
df["amount"] = pd.to_numeric(df["amount"])
print(df.groupby("city")["amount"].sum())这种脚本在学习 API 时可以,在商业项目里风险很大:
| 问题 | 后果 |
|---|---|
| 不指定订单号类型 | 000123 可能变成 123,订单追踪失败 |
| 不处理异常金额 | 一个 abc 就让任务失败,或者被错误忽略 |
| 不输出异常原因 | 业务不知道该修哪一行 |
| 不按业务主键去重 | 重复导入导致销售额翻倍 |
| 没有日志和摘要 | 定时任务失败后不知道处理了多少数据 |
| 没有测试 | 改规则后可能悄悄算错 |
端到端处理流程
flowchart TD
A["读取输入文件"] --> B["检查字段是否齐全"]
B --> C["标准化文本、金额、日期"]
C --> D["生成异常原因"]
D --> E["拆分有效数据和异常数据"]
E --> F["按 order_id 去重"]
F --> G["按 paid 口径聚合城市报表"]
G --> H["导出清洗结果、异常文件、报表"]
H --> I["写入运行摘要和日志"]这条链路的核心思想是:每一步都可解释、可复跑、可排查。
推荐文件结构
order_report/
report_orders.py
requirements.txt
data/
orders.csv
output/
tests/
test_report_orders.py职责划分:
| 文件 | 职责 |
|---|---|
report_orders.py | 主脚本,包含读取、清洗、聚合、导出 |
requirements.txt | 依赖版本,保证别人能复现环境 |
data/orders.csv | 示例输入文件 |
output/ | 运行输出目录 |
tests/test_report_orders.py | 清洗和聚合规则测试 |
依赖文件:
pandas>=2.0
pytest>=8.0示例输入文件
保存为 data/orders.csv:
order_id,city,amount,status,created_at,updated_at,channel
0001,北京,99.5,paid,2026-06-01,2026-06-01 10:00:00,app
0002,上海,120,paid,2026-06-01,2026-06-01 10:10:00,web
0003,北京,,refund,2026-06-02,2026-06-02 09:00:00,app
0004,广州,88,paid,2026-06-02,2026-06-02 11:00:00,mini
0004,广州,99,paid,2026-06-02,2026-06-02 12:00:00,mini
0005,上海,abc,paid,wrong-date,2026-06-03 08:00:00,web
0006,深圳,-10,paid,2026-06-03,2026-06-03 08:30:00,app
,杭州,50,paid,2026-06-03,2026-06-03 09:00:00,
0007,成都,66,unknown,2026-06-03,2026-06-03 09:30:00,web这份数据故意包含几类问题:
| 行为 | 目的 |
|---|---|
| 订单号带前导 0 | 验证订单号必须按字符串处理 |
| 金额为空 | 验证核心字段缺失 |
| 订单号重复 | 验证按更新时间保留最新 |
金额为 abc | 验证数值转换失败 |
日期为 wrong-date | 验证日期转换失败 |
| 金额为负 | 验证业务异常 |
状态为 unknown | 验证枚举校验 |
| 渠道为空 | 验证非核心字段默认值 |
完整可运行代码
保存为 report_orders.py:
from __future__ import annotations
import argparse
import json
import logging
from pathlib import Path
import pandas as pd
REQUIRED_COLUMNS = [
"order_id",
"city",
"amount",
"status",
"created_at",
"updated_at",
"channel",
]
VALID_STATUS = {"paid", "refund", "cancelled"}
def configure_logging() -> None:
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s %(levelname)s %(message)s",
)
def parse_args() -> argparse.Namespace:
parser = argparse.ArgumentParser(description="Clean orders and build city report.")
parser.add_argument("--input", required=True, help="Path to orders csv file.")
parser.add_argument("--output", required=True, help="Output directory.")
return parser.parse_args()
def load_orders(path: Path) -> pd.DataFrame:
logging.info("loading file: %s", path)
return pd.read_csv(
path,
encoding="utf-8",
dtype={"order_id": "string"},
)
def validate_columns(df: pd.DataFrame) -> None:
missing = [column for column in REQUIRED_COLUMNS if column not in df.columns]
if missing:
raise ValueError(f"missing required columns: {missing}")
def normalize_orders(df: pd.DataFrame) -> pd.DataFrame:
result = df.copy()
text_columns = ["order_id", "city", "status", "channel"]
for column in text_columns:
result[column] = result[column].astype("string").str.strip()
result["channel"] = result["channel"].replace("", pd.NA).fillna("unknown")
result["amount"] = pd.to_numeric(result["amount"], errors="coerce")
result["created_at"] = pd.to_datetime(result["created_at"], errors="coerce")
result["updated_at"] = pd.to_datetime(result["updated_at"], errors="coerce")
return result
def build_invalid_reason(row: pd.Series) -> str:
reasons: list[str] = []
if pd.isna(row["order_id"]) or row["order_id"] == "":
reasons.append("order_id is empty")
if pd.isna(row["city"]) or row["city"] == "":
reasons.append("city is empty")
if pd.isna(row["amount"]):
reasons.append("amount is not numeric")
elif row["amount"] < 0:
reasons.append("amount is negative")
if row["status"] not in VALID_STATUS:
reasons.append("status is invalid")
if pd.isna(row["created_at"]):
reasons.append("created_at is invalid")
if pd.isna(row["updated_at"]):
reasons.append("updated_at is invalid")
return "; ".join(reasons)
def split_invalid_rows(df: pd.DataFrame) -> tuple[pd.DataFrame, pd.DataFrame]:
checked = df.copy()
checked["invalid_reason"] = checked.apply(build_invalid_reason, axis=1)
invalid_mask = checked["invalid_reason"] != ""
invalid = checked[invalid_mask].copy()
valid = checked[~invalid_mask].drop(columns=["invalid_reason"]).copy()
return valid, invalid
def deduplicate_orders(df: pd.DataFrame) -> tuple[pd.DataFrame, int]:
duplicate_count = int(df.duplicated(subset=["order_id"]).sum())
deduped = (
df.sort_values(["order_id", "updated_at"])
.drop_duplicates(subset=["order_id"], keep="last")
.copy()
)
return deduped, duplicate_count
def build_city_report(df: pd.DataFrame) -> pd.DataFrame:
paid = df[df["status"] == "paid"].copy()
report = (
paid.groupby("city", as_index=False)
.agg(
order_count=("order_id", "nunique"),
total_amount=("amount", "sum"),
avg_amount=("amount", "mean"),
)
.sort_values(["total_amount", "city"], ascending=[False, True])
)
return report
def export_outputs(
output_dir: Path,
clean_orders: pd.DataFrame,
invalid_orders: pd.DataFrame,
city_report: pd.DataFrame,
summary: dict,
) -> None:
output_dir.mkdir(parents=True, exist_ok=True)
clean_orders.to_csv(output_dir / "clean_orders.csv", index=False, encoding="utf-8-sig")
invalid_orders.to_csv(output_dir / "invalid_orders.csv", index=False, encoding="utf-8-sig")
city_report.to_csv(output_dir / "city_report.csv", index=False, encoding="utf-8-sig")
with (output_dir / "run_summary.json").open("w", encoding="utf-8") as file:
json.dump(summary, file, ensure_ascii=False, indent=2, default=str)
def run(input_path: Path, output_dir: Path) -> dict:
raw = load_orders(input_path)
validate_columns(raw)
normalized = normalize_orders(raw)
valid, invalid = split_invalid_rows(normalized)
clean_orders, duplicate_count = deduplicate_orders(valid)
city_report = build_city_report(clean_orders)
summary = {
"status": "success",
"input": str(input_path),
"output": str(output_dir),
"raw_rows": int(len(raw)),
"valid_rows_before_deduplicate": int(len(valid)),
"invalid_rows": int(len(invalid)),
"duplicate_rows": duplicate_count,
"clean_rows": int(len(clean_orders)),
"report_rows": int(len(city_report)),
}
export_outputs(output_dir, clean_orders, invalid, city_report, summary)
logging.info("summary: %s", summary)
return summary
def main() -> None:
configure_logging()
args = parse_args()
run(Path(args.input), Path(args.output))
if __name__ == "__main__":
main()运行:
python report_orders.py --input data/orders.csv --output output运行后会生成:
output/
clean_orders.csv
invalid_orders.csv
city_report.csv
run_summary.json核心过程逐步解释
第一步:读取时为什么指定 order_id 为字符串
pd.read_csv(path, dtype={"order_id": "string"})订单号、资产编码、手机号、身份证号都不是数学数字。它们虽然由数字组成,但本质是标识符。
如果不指定字符串,可能出现:
| 问题 | 示例 |
|---|---|
| 前导 0 丢失 | 0001 变成 1 |
| 大数字精度丢失 | 长订单号被当成浮点数 |
| 科学计数法显示 | 1.23E+18 |
| JOIN 对不上 | 字符串订单号和数字订单号无法匹配 |
第二步:为什么先校验字段
missing = [column for column in REQUIRED_COLUMNS if column not in df.columns]字段缺失是结构性错误,不能继续往下算。比如没有 updated_at,重复订单就不知道该保留哪一条;没有 status,就无法只统计已支付订单。
这类错误应该直接让任务失败,而不是生成一个看似成功但实际错误的报表。
第三步:为什么清洗要先 copy
result = df.copy()copy 的作用是让函数不修改传入的原始 DataFrame。这样每个函数更像一个独立加工步骤,排查时可以比较原始数据、标准化数据、有效数据之间的差异。
如果函数偷偷修改入参,后续调试会很痛苦:你以为拿的是原始数据,其实已经被前面某一步改过了。
第四步:为什么异常原因要写清楚
异常文件只告诉业务“这一行错了”还不够,必须告诉他“为什么错”。
checked["invalid_reason"] = checked.apply(build_invalid_reason, axis=1)异常原因的价值:
| 价值 | 说明 |
|---|---|
| 业务可修复 | 业务知道是金额错、日期错还是状态错 |
| 技术可排查 | 能统计哪类错误最多 |
| 数据质量可治理 | 反推上游系统字段质量 |
| 可审计 | 后续追问为什么某行没统计,有证据 |
第五步:为什么去重在异常拆分之后
本例先拆异常,再对有效数据去重。原因是异常数据应该完整保留给业务修复。
如果先去重,可能会把一条异常重复数据丢掉,导致业务不知道上游到底传了多少问题数据。
不过真实项目要看业务规则:
| 规则 | 适合场景 |
|---|---|
| 先拆异常再去重 | 需要完整反馈源数据问题 |
| 先按主键取最新再校验 | 上游多次更新同一订单,只关心最新状态 |
| 重复且内容冲突全部异常 | 金额、状态冲突必须人工确认 |
第六步:为什么只统计 paid
商业报表一定要讲口径。销售额通常只统计支付成功订单,不统计退款、取消、无效订单。
paid = df[df["status"] == "paid"].copy()如果口径不清,会出现这些争议:
| 问题 | 结果 |
|---|---|
| 退款算不算负销售额 | 财务报表和运营报表可能不同 |
| 取消订单是否计数 | 订单量和成交量混淆 |
| 部分退款怎么算 | 需要明细级退款金额 |
| 支付后关闭怎么算 | 取决于业务状态机 |
测试代码
保存为 tests/test_report_orders.py:
from io import StringIO
import pandas as pd
from report_orders import (
build_city_report,
deduplicate_orders,
normalize_orders,
split_invalid_rows,
)
def make_df() -> pd.DataFrame:
raw = """order_id,city,amount,status,created_at,updated_at,channel
0001,北京,99.5,paid,2026-06-01,2026-06-01 10:00:00,app
0002,上海,abc,paid,wrong-date,2026-06-01 10:10:00,web
0003,广州,-1,paid,2026-06-02,2026-06-02 11:00:00,mini
0004,深圳,20,refund,2026-06-02,2026-06-02 12:00:00,
"""
return pd.read_csv(StringIO(raw), dtype={"order_id": "string"})
def test_split_invalid_rows() -> None:
normalized = normalize_orders(make_df())
valid, invalid = split_invalid_rows(normalized)
assert len(valid) == 2
assert len(invalid) == 2
assert invalid["invalid_reason"].str.contains("amount").any()
assert invalid["invalid_reason"].str.contains("created_at").any()
def test_deduplicate_keep_latest() -> None:
df = pd.DataFrame([
{"order_id": "1", "city": "北京", "amount": 10, "status": "paid", "created_at": "2026-06-01", "updated_at": "2026-06-01 10:00:00", "channel": "app"},
{"order_id": "1", "city": "北京", "amount": 20, "status": "paid", "created_at": "2026-06-01", "updated_at": "2026-06-01 11:00:00", "channel": "app"},
])
df = normalize_orders(df)
deduped, duplicate_count = deduplicate_orders(df)
assert duplicate_count == 1
assert len(deduped) == 1
assert deduped.iloc[0]["amount"] == 20
def test_city_report_only_paid() -> None:
df = pd.DataFrame([
{"order_id": "1", "city": "北京", "amount": 10, "status": "paid"},
{"order_id": "2", "city": "北京", "amount": 20, "status": "refund"},
{"order_id": "3", "city": "上海", "amount": 30, "status": "paid"},
])
report = build_city_report(df)
result = dict(zip(report["city"], report["total_amount"]))
assert result["北京"] == 10
assert result["上海"] == 30运行测试:
pytest -q测试的意义不是追求形式,而是保护业务口径。比如以后有人把 refund 也算进城市报表,这个测试就会失败。
大文件处理
如果文件很大,不能一次性读入内存,可以用分块:
parts = []
invalid_parts = []
for chunk in pd.read_csv("big_orders.csv", chunksize=100_000, dtype={"order_id": "string"}):
validate_columns(chunk)
normalized = normalize_orders(chunk)
valid, invalid = split_invalid_rows(normalized)
parts.append(valid)
invalid_parts.append(invalid)
all_valid = pd.concat(parts, ignore_index=True)
all_invalid = pd.concat(invalid_parts, ignore_index=True)
clean_orders, duplicate_count = deduplicate_orders(all_valid)
city_report = build_city_report(clean_orders)注意:分块读取不等于分块后直接聚合。因为同一个订单可能出现在不同块里,跨块去重必须在全量有效数据合并后做。
常见大文件策略:
| 问题 | 解决方向 |
|---|---|
| 文件几百 MB | chunksize 分块读取 |
| 文件数十 GB | 先入库,用 SQL 或 Spark 处理 |
| 需要本地高性能分析 | 考虑 DuckDB 或 Polars |
| 跨块去重很大 | 用数据库唯一键或外部排序 |
| 内存持续上涨 | 不要把所有明细放列表,尽早落盘或聚合 |
生产排查流程
flowchart TD
A["报表结果不对"] --> B["确认输入文件版本和行数"]
B --> C["查看 run_summary.json"]
C --> D["检查 invalid_orders.csv"]
D --> E["检查重复订单数量"]
E --> F["抽查 clean_orders.csv 明细"]
F --> G["核对 paid 过滤口径"]
G --> H["核对 groupby 聚合结果"]
H --> I["和历史报表或数据库对账"]排查清单:
| 现象 | 可能原因 | 怎么查 |
|---|---|---|
| 销售额突然变少 | 异常行变多、状态口径变了、文件没读全 | 看 invalid_rows、状态分布、文件行数 |
| 销售额突然翻倍 | 重复订单没去重、重复文件导入 | 看 duplicate_rows 和订单号重复 |
| 某城市消失 | 城市为空、城市名带空格、维表映射失败 | 看 city 的 value_counts() |
| 任务内存高 | 文件过大、结果列表无限累积 | 看文件大小、chunk 逻辑、中间对象 |
| 任务成功但报表为空 | paid 口径不匹配、状态大小写没标准化 | 看 status 分布 |
| Excel 打开乱码 | 编码不适配 | CSV 导出用 utf-8-sig |
商业扩展场景
| 场景 | 改造点 |
|---|---|
| 订单日报 | 按日期、城市、渠道聚合 |
| 库存对账 | 按商品、仓库、批次对齐系统库存 |
| 医疗资产导入 | 校验资产编码、科室、负责人、访问次数 |
| 异步入库前校验 | 先清洗成有效文件,再发 MQ 入库 |
| 定时批处理 | 接入 XXL-JOB、Airflow 或系统定时任务 |
| 数据质量平台 | 把异常原因分类统计,形成质量看板 |
面试标准回答
数据处理脚本怎么设计才可靠?
可以这样答:
我不会把所有逻辑写在一个脚本里直接 groupby,而是先明确输入输出、字段口径和异常规则。实现上会拆成读取、字段校验、类型标准化、异常拆分、业务去重、聚合、导出和运行摘要几个步骤。核心字段转换失败不会静默丢弃,而是输出异常文件和原因。对金额、状态、重复订单、时间范围这类关键规则会写测试。生产运行时会记录原始行数、有效行数、异常行数、重复数和输出路径,方便定时任务失败或报表不一致时排查。为什么异常数据不能直接删除?
异常数据代表上游数据质量问题。直接删除会让报表看似正常,但业务不知道哪些数据没有被统计,也无法修复源头。正确做法是拆分有效数据和异常数据,异常数据带上原因单独输出,必要时阻断任务或通知负责人。
大文件怎么处理?
大文件不能无脑一次性读入内存。可以先用 usecols 减少列,用 chunksize 分块读取,清洗后尽早聚合或落盘。但如果涉及跨块去重、全局排序、全局 TopN,就必须在最终阶段统一处理。数据量继续增大时,应考虑数据库、DuckDB、Polars 或 Spark。
报表金额对不上怎么排查?
先确认输入文件版本、行数和时间范围,再看运行摘要中的有效行、异常行、重复行;接着检查金额和日期转换失败、状态分布、重复订单、过滤条件和聚合口径;最后用明细抽样和聚合前后总金额与数据库或历史报表对账。
关联知识点
| 知识点 | 作用 |
|---|---|
| Python 数据处理总览 | 建立清洗、校验、聚合、导出的主流程 |
| Pandas 清洗分析 | 深入理解 DataFrame、缺失、去重、groupby |
| NumPy 数组计算 | 理解向量化、dtype、axis 和底层数组思维 |
| 异常与文件 | 理解文件编码、异常处理和资源关闭 |
| 日志与调试 | 让脚本在生产环境可排查 |
| Python 测试 | 给清洗规则和报表口径写测试 |
| Python 面试题 | 查看数据处理相关标准回答 |
本章小结
一个能在商业项目里长期运行的数据处理任务,关键不是会不会写 read_csv,而是能不能把数据变成可信结果。可信结果来自清晰的字段口径、严格的类型转换、可解释的异常文件、明确的去重规则、可验证的聚合逻辑、可追踪的日志摘要和自动化测试。
