pandas 百万级 CSV 聚合与异常检测:chunked 读取、Z-score 与业务实战

一、为什么 pandas.read_csv 是内存刺客

很多数据脚本从下面这一行开始:

1
2
3
import pandas as pd

df = pd.read_csv("access.log.csv")

小文件没有问题,但 CSV 并不是适合直接映射到内存的格式。它没有列类型信息、没有索引、没有压缩后的随机访问能力。读取时,程序需要完成文本解码、字段切分、类型推断和对象构造。一个 5GB 的 CSV,最终形成的 DataFrame 可能占用 15GB~30GB。

原因主要有三点。

第一,字符串列通常会产生大量 Python 对象。IP、User-Agent、URL 和 Referer 看起来只是文本,但每一行都可能拥有独立对象。第二,默认类型推断需要扫描数据,并且不稳定的列可能被整体推断为 object。第三,CSV 的文本表示比二进制列式格式冗余得多。

因此,32GB 机器读取 5GB 日志时出现 swap 并不罕见。进入 swap 后,CPU 使用率可能下降到个位数,磁盘 I/O 却持续升高,任务看起来像“卡死”。

处理大 CSV 的基本原则是:

  1. 明确列类型,避免自动推断。
  2. 只读取需要的列。
  3. 使用分块或惰性执行。
  4. 尽早聚合,避免保留明细。
  5. 结果写入 Parquet、数据库或聚合表。

二、三种分块读取方案

1. pandas chunked

read_csv(chunksize=...) 返回一个迭代器,每次只在内存中保留一个块。它最容易接入现有 pandas 代码,适合单机批处理。

1
2
3
4
5
6
7
for chunk in pd.read_csv(
"access.log.csv",
chunksize=100_000,
usecols=["timestamp", "ip", "status", "bytes"],
dtype={"ip": "string", "status": "int16", "bytes": "int64"},
):
process(chunk)

缺点是全局聚合需要自行合并中间结果,代码容易出现“每块省内存、最后合并爆内存”的问题。

2. Dask

Dask 将一个 DataFrame 拆成多个 partition,并构建任务图。它可以在多进程、多线程或多机环境中执行,适合数据规模明显超过单机内存的场景。

1
2
3
4
import dask.dataframe as dd

df = dd.read_csv("access-*.csv", blocksize="128MB")
result = df.groupby("status").bytes.sum().compute()

Dask 的成本是调试门槛更高。很多操作是惰性的,代码执行到 compute() 才真正运行;不熟悉任务图时,性能问题不容易定位。

3. Polars

Polars 使用 Rust 编写,支持多线程和惰性查询。scan_csv 不会立即读取数据,只有调用 collect() 时才执行查询计划。

1
2
3
4
5
6
7
8
9
10
import polars as pl

result = (
pl.scan_csv("access.log.csv")
.select(["timestamp", "status", "bytes"])
.with_columns(pl.col("timestamp").str.to_datetime())
.group_by("status")
.agg(pl.col("bytes").sum())
.collect()
)

在过滤、投影和聚合明确的任务中,Polars 往往比 pandas 快 5~10 倍。若团队已有大量 pandas 代码,chunked 是低风险迁移路径;若新项目以分析管道为主,Polars 值得优先评估。

三、实战 1:chunked 读取、按小时和状态码聚合

下面脚本处理包含 timestampipstatusbytes 四列的访问日志。它不会保存所有明细,只维护各块的聚合结果。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
from __future__ import annotations

from collections import defaultdict
from pathlib import Path
import pandas as pd

INPUT = Path("access.log.csv")
CHUNK_SIZE = 100_000

hourly = defaultdict(lambda: {"requests": 0, "bytes": 0})
by_ip = defaultdict(int)
by_status = defaultdict(int)

for chunk in pd.read_csv(
INPUT,
chunksize=CHUNK_SIZE,
usecols=["timestamp", "ip", "status", "bytes"],
dtype={"timestamp": "string", "ip": "string", "status": "int16", "bytes": "int64"},
):
chunk["timestamp"] = pd.to_datetime(
chunk["timestamp"], utc=True, errors="coerce"
)
chunk = chunk.dropna(subset=["timestamp", "ip"])

for hour, group in chunk.groupby(chunk["timestamp"].dt.floor("h")):
hourly[str(hour)]["requests"] += len(group)
hourly[str(hour)]["bytes"] += int(group["bytes"].sum())

for ip, count in chunk["ip"].value_counts().items():
by_ip[str(ip)] += int(count)

for status, count in chunk["status"].value_counts().items():
by_status[int(status)] += int(count)

pd.DataFrame.from_dict(hourly, orient="index").sort_index().to_csv(
"hourly_summary.csv", index_label="hour"
)
pd.DataFrame(
sorted(by_ip.items()), columns=["ip", "requests"]
).sort_values("requests", ascending=False).to_csv(
"ip_summary.csv", index=False
)
pd.DataFrame(
sorted(by_status.items()), columns=["status", "requests"]
).to_csv("status_summary.csv", index=False)

这里的关键不是 chunksize 本身,而是中间结果的规模。IP 的基数如果达到千万级,by_ip 仍然可能占用大量内存。这时应改成写入 SQLite、DuckDB 或每块输出 Parquet,再做二次聚合。

四、实战 2:Dask 处理十亿行日志

Dask 适合多个文件组成的日志目录。下面示例只读取指定列,并通过 split_out 增加 groupby 的并行输出分区。

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
import dask.dataframe as dd

logs = dd.read_csv(
"logs/access-*.csv",
blocksize="256MB",
usecols=["timestamp", "status", "bytes"],
dtype={"status": "int16", "bytes": "int64"},
assume_missing=True,
)

logs["timestamp"] = dd.to_datetime(logs["timestamp"], utc=True)
logs["hour"] = logs["timestamp"].dt.floor("h")

result = (
logs.groupby(["hour", "status"], observed=True)
.agg(requests=("status", "size"), total_bytes=("bytes", "sum"))
.reset_index()
.compute()
)

result.to_parquet("summary.parquet", index=False)

对于十亿行数据,不能只看 Python 代码。还要观察 partition 大小、磁盘带宽和 shuffle。partition 太小会产生大量调度开销,太大则会让单个 worker 内存峰值过高。通常可以从 128MB~512MB 进行测试。

五、实战 3:Polars 高性能聚合

Polars 的惰性 API 会尽可能提前过滤和裁剪列:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
import polars as pl

query = (
pl.scan_csv(
"access.log.csv",
schema_overrides={
"ip": pl.String,
"status": pl.Int16,
"bytes": pl.Int64,
},
)
.with_columns(
pl.col("timestamp").str.to_datetime(strict=False).dt.replace_time_zone("UTC")
)
.filter(pl.col("status") >= 400)
.group_by(
pl.col("timestamp").dt.truncate("1h").alias("hour"),
"status",
)
.agg(
pl.len().alias("requests"),
pl.col("bytes").sum().alias("total_bytes"),
)
.sort("hour")
)

query.collect().write_parquet("error_summary.parquet")

Polars 与 pandas 的差异需要特别注意:字符串类型、缺失值、日期时间和表达式语法都不同。不要机械地把 pandas 代码逐行翻译;应尽量使用表达式,让优化器理解整个查询。

六、三种异常检测算法

Z-score

Z-score 为:

[
z = \frac{x-\mu}{\sigma}
]

通常 abs(z) > 3 被视为异常。它适合近似正态分布的数据,例如延迟在稳定系统中的某个固定区间内波动。

1
2
3
4
5
6
7
8
9
import pandas as pd

df = pd.read_csv("hourly_metrics.csv")
mean = df["requests"].mean()
std = df["requests"].std(ddof=0)

df["z_score"] = (df["requests"] - mean) / std
outliers = df[df["z_score"].abs() > 3]
outliers.to_csv("zscore_outliers.csv", index=False)

如果标准差为零,需要提前处理;如果数据存在趋势,应先按小时、星期或业务分组,否则白天和夜间的正常差异也可能被判定为异常。

IQR

IQR 为第三四分位数减第一四分位数:

[
IQR=Q_3-Q_1
]

低于 Q1 - 1.5*IQR 或高于 Q3 + 1.5*IQR 的数据可标记为异常。IQR 对偏态分布和少量极端值更稳健。

1
2
3
4
5
6
7
8
9
10
11
import pandas as pd

df = pd.read_csv("latency.csv")
q1 = df["latency_ms"].quantile(0.25)
q3 = df["latency_ms"].quantile(0.75)
iqr = q3 - q1

low = q1 - 1.5 * iqr
high = q3 + 1.5 * iqr
df["is_outlier"] = ~df["latency_ms"].between(low, high)
df[df["is_outlier"]].to_csv("iqr_outliers.csv", index=False)

IsolationForest

IsolationForest 通过随机切分特征空间隔离样本。异常点通常更容易被孤立,适合请求数、延迟、错误率、响应大小等多维特征。

1
2
3
4
5
6
7
8
9
10
11
12
13
import pandas as pd
from sklearn.ensemble import IsolationForest

df = pd.read_csv("metrics.csv")
features = ["requests", "latency_ms", "error_rate", "bytes"]

model = IsolationForest(
n_estimators=200,
contamination=0.01,
random_state=42,
)
df["anomaly"] = model.fit_predict(df[features].fillna(0))
df[df["anomaly"] == -1].to_csv("iforest_outliers.csv", index=False)

contamination 不是越小越好。它应结合历史告警比例、人工复核成本和业务容忍度校准。无监督算法只能发现“与历史模式不同”的点,不能自动证明这些点就是故障。

七、踩坑记录

第一,分块后直接 concat 如果每个块都被保留下来,最后的 concat 仍会产生全量内存峰值。应只保存聚合结果,或直接写 Parquet。

第二,字符串列占用过大。 可以使用 category,但只有低基数列适合。千万级不同 URL 使用 category 可能造成类别表本身膨胀。

第三,时区不一致。 pd.to_datetime(..., utc=True) 可以将输入统一为 UTC。不要把带时区和不带时区的时间直接比较。

第四,异常检测污染。 计算均值和标准差之前,数据中若已经混有大面积故障,Z-score 的基准就会被污染。生产系统应使用滚动窗口或历史健康区间。

八、小结

场景 推荐方案 主要优势
单机、改造旧 pandas 脚本 chunked 成本最低
多文件、数据远超内存 Dask 分区和分布式执行
新建分析管道 Polars 多线程、惰性优化
近似正态分布 Z-score 解释简单
偏态或极端值较多 IQR 稳健
多维、无标签数据 IsolationForest 不需要人工阈值

最重要的优化通常不是换库,而是只读必要列、尽早过滤、避免明细级全量缓存,并把中间结果保存为 Parquet。


pandas 百万级 CSV 聚合与异常检测:chunked 读取、Z-score 与业务实战
https://blog.calcguide.tech/2026-08-10-pandas百万级CSV聚合与异常检测/
作者
王争气
发布于
2026年8月10日
许可协议