超大规模数据库(2.5亿行)极速过滤方案咨询
超大规模数据集高效过滤方案咨询
问题背景
- 数据集规模:2.5亿行,15列(计划扩展至20列),当前存储为CSV格式
- 需求:对数据集执行多条件过滤并检索指定列数据,目标耗时控制在1秒以内
- 当前实现:使用Dask处理,耗时超60秒,过滤逻辑如下:
from dask import dataframe as dd df = dd.read_csv(path, sep=",") df = df[(df["prob"] > 0.99) & (df["actions_preflop"] == "re") & (df["actions_flop"] == "re") & (df["actions_turn"] == "-") & (df["actions_river"] == "-") & (df["Table"] == "14s13s10h") ] df = df.compute() raise_std = df["raise"].values call_std = df["call"].values
可行解决方案
1. 转换为列式存储格式(Parquet/ORC)
CSV是行式存储,扫描效率极低,转换为Parquet或ORC这类列式存储格式,可实现谓词下推(只扫描需要的列和符合条件的行),同时支持压缩减少IO开销。建议对高选择性的过滤列(如Table)进行分区,进一步减少扫描范围。
示例代码(用Dask转换并读取Parquet):
# 首次转换CSV为Parquet,按Table分区 df = dd.read_csv(path, sep=",") df.to_parquet("data.parquet", partition_on="Table", compression="snappy") # 后续读取与过滤 df = dd.read_parquet("data.parquet", filters=[("Table", "==", "14s13s10h")]) df = df[(df["prob"] > 0.99) & (df["actions_preflop"] == "re") & (df["actions_flop"] == "re") & (df["actions_turn"] == "-") & (df["actions_river"] == "-") ] # 仅加载需要的列,减少数据传输 result = df[["raise", "call"]].compute() raise_std = result["raise"].values call_std = result["call"].values
2. 使用OLAP引擎(DuckDB/ClickHouse)
这类引擎专为大规模数据分析优化,支持向量查询、列存储和高效索引,单机器即可轻松实现秒级过滤。
DuckDB(嵌入式,无需部署)
import duckdb # 直接查询CSV(或Parquet,速度更快) con = duckdb.connect() result = con.execute(""" SELECT raise, call FROM read_csv(?) WHERE prob > 0.99 AND actions_preflop = 're' AND actions_flop = 're' AND actions_turn = '-' AND actions_river = '-' AND Table = '14s13s10h' """, [path]).fetchdf() raise_std = result["raise"].values call_std = result["call"].values
ClickHouse(单机/集群部署)
如果数据量持续增长,可部署ClickHouse,通过Python客户端连接查询:
from clickhouse_driver import Client client = Client(host='localhost') result = client.execute(""" SELECT raise, call FROM your_table WHERE prob > 0.99 AND actions_preflop = 're' AND actions_flop = 're' AND actions_turn = '-' AND actions_river = '-' AND Table = '14s13s10h' """) # 转换为数组 raise_std = [row[0] for row in result] call_std = [row[1] for row in result]
3. 内存映射工具(Vaex)
Vaex利用内存映射技术,无需将全量数据加载到内存,即可实现快速过滤,适合单机器内存不足但需要快速查询的场景。
示例代码:
import vaex df = vaex.read_csv(path) filtered_df = df[(df.prob > 0.99) & (df.actions_preflop == 're') & (df.actions_flop == 're') & (df.actions_turn == '-') & (df.actions_river == '-') & (df.Table == '14s13s10h')] raise_std = filtered_df.raise.values call_std = filtered_df.call.values
4. 过滤逻辑与读取优化
如果坚持使用Dask,可通过以下优化缩短耗时:
- 指定读取列:只加载过滤和输出需要的列,减少IO和内存占用
cols = ["prob", "actions_preflop", "actions_flop", "actions_turn", "actions_river", "Table", "raise", "call"] df = dd.read_csv(path, sep=",", usecols=cols) - 调整过滤顺序:将选择性最高的条件(如
Table等值过滤)放在最前面,让引擎优先过滤掉大部分数据 - 增加分区数:Dask默认分区可能过大,可通过
blocksize参数调整分区大小(如blocksize="64MB"),提升并行处理效率
内容的提问来源于stack exchange,提问作者Raphaël Ambit
相关产品推荐
相关产品推荐

