Skip to content

数据处理端到端实践

这一页不是“跑一个 CSV 报表”的博客示例,而是按商业项目的方式,把一个数据处理脚本从需求、字段口径、目录结构、清洗函数、异常数据、日志、测试、运行参数到排查链路完整串起来。

你可以把它理解成一个小型离线任务模板:业务每天上传订单或资产文件,系统自动清洗、校验、聚合并输出结果。如果结果不对,必须能知道错在哪一层,而不是只看到一个失败堆栈。

学习目标

学完本页,你应该能做到:

  1. 明确数据处理任务的输入、输出和业务口径。
  2. 设计一个可维护的数据脚本目录,而不是把所有逻辑写进一个文件。
  3. 解释读取、校验、转换、异常拆分、去重、聚合、导出的完整流程。
  4. 写出可运行的 Pandas 商业 Demo。
  5. 给清洗逻辑写测试,保证口径变更时能及时发现问题。
  6. 使用命令行参数、日志和处理报告,让脚本适合定时任务。
  7. 排查报表金额对不上、异常行过多、内存过高、重复订单等生产问题。
  8. 面试时能说清“数据处理项目怎么做才可靠”。

商业需求

假设电商或医疗资产平台每天会产生一份订单文件 orders.csv,字段如下:

字段含义是否核心处理要求
order_id订单号不能为空,按字符串处理,重复时保留更新时间最新的一条
city城市不能为空,去掉前后空格
amount订单金额转数字,不能转换进入异常文件
status订单状态只能是 paidrefundcancelled
created_at创建时间转日期,不能转换进入异常文件
updated_at更新时间转日期,重复订单按它保留最新
channel渠道缺失时标记为 unknown

目标输出:

  1. clean_orders.csv:清洗后的有效订单。
  2. invalid_orders.csv:异常订单,并说明异常原因。
  3. city_report.csv:按城市统计已支付订单数、总金额、平均金额。
  4. run_summary.json:本次任务行数、异常数、重复数、输出路径和处理状态。

为什么不能只写一个简单脚本

简单写法通常是这样:

python
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 就让任务失败,或者被错误忽略
不输出异常原因业务不知道该修哪一行
不按业务主键去重重复导入导致销售额翻倍
没有日志和摘要定时任务失败后不知道处理了多少数据
没有测试改规则后可能悄悄算错

端到端处理流程

mermaid
flowchart TD
    A["读取输入文件"] --> B["检查字段是否齐全"]
    B --> C["标准化文本、金额、日期"]
    C --> D["生成异常原因"]
    D --> E["拆分有效数据和异常数据"]
    E --> F["按 order_id 去重"]
    F --> G["按 paid 口径聚合城市报表"]
    G --> H["导出清洗结果、异常文件、报表"]
    H --> I["写入运行摘要和日志"]

这条链路的核心思想是:每一步都可解释、可复跑、可排查。

推荐文件结构

text
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清洗和聚合规则测试

依赖文件:

text
pandas>=2.0
pytest>=8.0

示例输入文件

保存为 data/orders.csv

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

python
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()

运行:

bash
python report_orders.py --input data/orders.csv --output output

运行后会生成:

text
output/
  clean_orders.csv
  invalid_orders.csv
  city_report.csv
  run_summary.json

核心过程逐步解释

第一步:读取时为什么指定 order_id 为字符串

python
pd.read_csv(path, dtype={"order_id": "string"})

订单号、资产编码、手机号、身份证号都不是数学数字。它们虽然由数字组成,但本质是标识符。

如果不指定字符串,可能出现:

问题示例
前导 0 丢失0001 变成 1
大数字精度丢失长订单号被当成浮点数
科学计数法显示1.23E+18
JOIN 对不上字符串订单号和数字订单号无法匹配

第二步:为什么先校验字段

python
missing = [column for column in REQUIRED_COLUMNS if column not in df.columns]

字段缺失是结构性错误,不能继续往下算。比如没有 updated_at,重复订单就不知道该保留哪一条;没有 status,就无法只统计已支付订单。

这类错误应该直接让任务失败,而不是生成一个看似成功但实际错误的报表。

第三步:为什么清洗要先 copy

python
result = df.copy()

copy 的作用是让函数不修改传入的原始 DataFrame。这样每个函数更像一个独立加工步骤,排查时可以比较原始数据、标准化数据、有效数据之间的差异。

如果函数偷偷修改入参,后续调试会很痛苦:你以为拿的是原始数据,其实已经被前面某一步改过了。

第四步:为什么异常原因要写清楚

异常文件只告诉业务“这一行错了”还不够,必须告诉他“为什么错”。

python
checked["invalid_reason"] = checked.apply(build_invalid_reason, axis=1)

异常原因的价值:

价值说明
业务可修复业务知道是金额错、日期错还是状态错
技术可排查能统计哪类错误最多
数据质量可治理反推上游系统字段质量
可审计后续追问为什么某行没统计,有证据

第五步:为什么去重在异常拆分之后

本例先拆异常,再对有效数据去重。原因是异常数据应该完整保留给业务修复。

如果先去重,可能会把一条异常重复数据丢掉,导致业务不知道上游到底传了多少问题数据。

不过真实项目要看业务规则:

规则适合场景
先拆异常再去重需要完整反馈源数据问题
先按主键取最新再校验上游多次更新同一订单,只关心最新状态
重复且内容冲突全部异常金额、状态冲突必须人工确认

第六步:为什么只统计 paid

商业报表一定要讲口径。销售额通常只统计支付成功订单,不统计退款、取消、无效订单。

python
paid = df[df["status"] == "paid"].copy()

如果口径不清,会出现这些争议:

问题结果
退款算不算负销售额财务报表和运营报表可能不同
取消订单是否计数订单量和成交量混淆
部分退款怎么算需要明细级退款金额
支付后关闭怎么算取决于业务状态机

测试代码

保存为 tests/test_report_orders.py

python
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

运行测试:

bash
pytest -q

测试的意义不是追求形式,而是保护业务口径。比如以后有人把 refund 也算进城市报表,这个测试就会失败。

大文件处理

如果文件很大,不能一次性读入内存,可以用分块:

python
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)

注意:分块读取不等于分块后直接聚合。因为同一个订单可能出现在不同块里,跨块去重必须在全量有效数据合并后做。

常见大文件策略:

问题解决方向
文件几百 MBchunksize 分块读取
文件数十 GB先入库,用 SQL 或 Spark 处理
需要本地高性能分析考虑 DuckDB 或 Polars
跨块去重很大用数据库唯一键或外部排序
内存持续上涨不要把所有明细放列表,尽早落盘或聚合

生产排查流程

mermaid
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 和订单号重复
某城市消失城市为空、城市名带空格、维表映射失败cityvalue_counts()
任务内存高文件过大、结果列表无限累积看文件大小、chunk 逻辑、中间对象
任务成功但报表为空paid 口径不匹配、状态大小写没标准化status 分布
Excel 打开乱码编码不适配CSV 导出用 utf-8-sig

商业扩展场景

场景改造点
订单日报按日期、城市、渠道聚合
库存对账按商品、仓库、批次对齐系统库存
医疗资产导入校验资产编码、科室、负责人、访问次数
异步入库前校验先清洗成有效文件,再发 MQ 入库
定时批处理接入 XXL-JOB、Airflow 或系统定时任务
数据质量平台把异常原因分类统计,形成质量看板

面试标准回答

数据处理脚本怎么设计才可靠?

可以这样答:

text
我不会把所有逻辑写在一个脚本里直接 groupby,而是先明确输入输出、字段口径和异常规则。实现上会拆成读取、字段校验、类型标准化、异常拆分、业务去重、聚合、导出和运行摘要几个步骤。核心字段转换失败不会静默丢弃,而是输出异常文件和原因。对金额、状态、重复订单、时间范围这类关键规则会写测试。生产运行时会记录原始行数、有效行数、异常行数、重复数和输出路径,方便定时任务失败或报表不一致时排查。

为什么异常数据不能直接删除?

异常数据代表上游数据质量问题。直接删除会让报表看似正常,但业务不知道哪些数据没有被统计,也无法修复源头。正确做法是拆分有效数据和异常数据,异常数据带上原因单独输出,必要时阻断任务或通知负责人。

大文件怎么处理?

大文件不能无脑一次性读入内存。可以先用 usecols 减少列,用 chunksize 分块读取,清洗后尽早聚合或落盘。但如果涉及跨块去重、全局排序、全局 TopN,就必须在最终阶段统一处理。数据量继续增大时,应考虑数据库、DuckDB、Polars 或 Spark。

报表金额对不上怎么排查?

先确认输入文件版本、行数和时间范围,再看运行摘要中的有效行、异常行、重复行;接着检查金额和日期转换失败、状态分布、重复订单、过滤条件和聚合口径;最后用明细抽样和聚合前后总金额与数据库或历史报表对账。

关联知识点

知识点作用
Python 数据处理总览建立清洗、校验、聚合、导出的主流程
Pandas 清洗分析深入理解 DataFrame、缺失、去重、groupby
NumPy 数组计算理解向量化、dtype、axis 和底层数组思维
异常与文件理解文件编码、异常处理和资源关闭
日志与调试让脚本在生产环境可排查
Python 测试给清洗规则和报表口径写测试
Python 面试题查看数据处理相关标准回答

本章小结

一个能在商业项目里长期运行的数据处理任务,关键不是会不会写 read_csv,而是能不能把数据变成可信结果。可信结果来自清晰的字段口径、严格的类型转换、可解释的异常文件、明确的去重规则、可验证的聚合逻辑、可追踪的日志摘要和自动化测试。