一、为什么 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 的基本原则是:
明确列类型,避免自动推断。
只读取需要的列。
使用分块或惰性执行。
尽早聚合,避免保留明细。
结果写入 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 读取、按小时和状态码聚合 下面脚本处理包含 timestamp、ip、status、bytes 四列的访问日志。它不会保存所有明细,只维护各块的聚合结果。
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 annotationsfrom collections import defaultdictfrom pathlib import Pathimport 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 pdfrom 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。