8亿行CSV按用户ID分组、时间戳排序去重Col1的Dask实现咨询
输入数据
多个结构相同的CSV文件,总计8亿行,列包含:Time Stamp(时间戳)、User ID(用户ID)、Col1、Col2、Col3
可用资源
60GB内存、24核CPU(服务器还有其他负载,无法全量使用内存)
需求目标
按User ID分组,每组内按Time Stamp排序,对Col1进行去重并保留基于时间戳的出现顺序
已尝试方案
- 使用
joblib并行加载CSV,通过pandas排序时出错 - 尝试Dask(新手阶段),代码如下:
from dask.distributed import LocalCluster from dask.dataframe import read_csv port = 8787 # 示例端口 cluster = LocalCluster(dashboard_address=f':{port}', n_workers=4, threads_per_worker=4, memory_limit='7GB') # 服务器有其他负载,无法用满60G内存 ddf = read_csv("/path/*.csv") ddf = ddf.set_index("Time Stamp") ddf.to_csv("/outdir/")
疑问
- 假设Dask通过磁盘排序输出多份文件,后续用
read_csv读取这些文件时,是否能保留顺序? - 如何在Dask中实现需求的核心逻辑?在pandas中可以用以下函数实现:
def getUnique(user_group): # 假设每个用户的行已按时间戳排序 res = list() for val in user_group["Col1"]: if val not in res: res.append(val) return res
另外,有没有比Dask更合适的替代方案?
疑问1解答
Dask输出的多份CSV文件不会自动保留全局排序顺序。因为Dask的set_index是分块执行的,输出的每个文件对应一个索引区间的数据块,块间是有序的,但单个文件内是该区间内的有序数据。如果后续直接用read_csv读取所有输出文件,只会按文件读取顺序拼接,不会自动合并成全局有序的结果。
如果需要后续读取后仍保留顺序,有两种可行方式:
- 输出时指定
write_index=True,后续读取后重新执行set_index完成全局排序 - 若内存允许,用
ddf.compute()合并为单个pandas DataFrame后保存(但8亿行数据大概率内存不足,不推荐)
疑问2:Dask实现核心逻辑的方案
核心需求拆解为按用户分组→组内按时间戳排序→Col1按时间顺序去重,Dask可通过以下步骤实现:
步骤1:数据分区与组内排序
先按User ID分区,确保同一用户的数据落在同一分区,减少跨分区计算开销;再对每个用户组内的数据按时间戳排序:
from dask.distributed import LocalCluster, Client from dask.dataframe import read_csv # 初始化集群:根据服务器负载调整参数,示例配置总内存约48GB(6*8),不超60GB上限 cluster = LocalCluster(dashboard_address=':8787', n_workers=6, threads_per_worker=4, memory_limit='8GB') client = Client(cluster) # 读取CSV,指定dtype减少内存占用,解析时间戳为datetime类型 ddf = read_csv( "/path/*.csv", dtype={ "User ID": "int64", "Col1": "object", # 根据实际数据类型调整,比如int64 "Col2": "float64", "Col3": "float64" }, parse_dates=["Time Stamp"] ) # 按User ID分区,组内按Time Stamp排序 ddf = ddf.set_partition("User ID") ddf_sorted = ddf.groupby("User ID").apply( lambda df: df.sort_values("Time Stamp"), meta=ddf._meta # 指定输出元数据结构,确保类型一致 )
步骤2:组内Col1去重(保留时间顺序)
推荐用向量化方法替代循环,性能更优;也可以沿用你熟悉的pandas逻辑:
方法1:向量化去重(推荐)
利用shift对比前后行的Col1值,标记首次出现的行:
def dedup_col1_vectorized(group): # 确保组内已按时间戳排序 group = group.sort_values("Time Stamp") # 标记Col1首次出现的行(第一行shift后为NaN,用fillna(True)保留) group["is_first"] = group["Col1"].ne(group["Col1"].shift()).fillna(True) # 保留首次出现的行,删除辅助列 return group[group["is_first"]].drop(columns=["is_first"]) # 应用去重函数,指定元数据 meta = ddf_sorted._meta ddf_deduped = ddf_sorted.groupby("User ID").apply(dedup_col1_vectorized, meta=meta)
方法2:沿用自定义循环逻辑
如果习惯原pandas函数,直接在Dask的apply中使用即可:
def getUnique(user_group): res = list() # 先确保组内按时间戳排序 user_group = user_group.sort_values("Time Stamp") # 过滤出未出现过的Col1值 filtered = user_group[~user_group["Col1"].isin(res)] for val in filtered["Col1"]: res.append(val) return user_group[user_group["Col1"].isin(res)] ddf_deduped = ddf_sorted.groupby("User ID").apply(getUnique, meta=meta)
步骤3:输出结果
按User ID分区保存,方便后续读取和处理:
ddf_deduped.to_csv("/outdir/deduped_*.csv", index=False)
替代方案推荐
PySpark
适合超大规模分布式数据处理,API成熟稳定,处理8亿行数据的容错性更好。核心逻辑示例:
from pyspark.sql import SparkSession from pyspark.sql.window import Window from pyspark.sql.functions import col, lag spark = SparkSession.builder.appName("UserCol1Dedup").getOrCreate() df = spark.read.csv("/path/*.csv", header=True, inferSchema=True) # 定义窗口:按User ID分组,按Time Stamp排序 window = Window.partitionBy("User ID").orderBy("Time Stamp") # 标记Col1首次出现的行 df = df.withColumn("prev_col1", lag("Col1", 1).over(window)) df_deduped = df.filter((col("prev_col1").isNull()) | (col("Col1") != col("prev_col1"))).drop("prev_col1") df_deduped.write.csv("/outdir/spark_deduped", header=True)
Vaex
基于内存映射的单机器大数据处理库,API与pandas高度兼容,适合60GB内存的场景,无需分布式集群:
import vaex # 内存映射读取所有CSV df = vaex.open("/path/*.csv") # 按User ID和Time Stamp排序,组内保留Col1首次出现的行 df_deduped = df.sort(["User ID", "Time Stamp"]).groupby("User ID").agg({ "Time Stamp": "first", "Col1": "first", "Col2": "first", # 根据需求选择聚合方式,比如last/mean "Col3": "first" }) # 导出结果 df_deduped.export_csv("/outdir/vaex_deduped.csv")
内容的提问来源于stack exchange,提问作者willfigureitout

