You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

8亿行CSV按用户ID分组、时间戳排序去重Col1的Dask实现咨询

问题背景

输入数据

多个结构相同的CSV文件,总计8亿行,列包含:Time Stamp(时间戳)、User ID(用户ID)、Col1、Col2、Col3

可用资源

60GB内存、24核CPU(服务器还有其他负载,无法全量使用内存)

需求目标

按User ID分组,每组内按Time Stamp排序,对Col1进行去重并保留基于时间戳的出现顺序

已尝试方案

  1. 使用joblib并行加载CSV,通过pandas排序时出错
  2. 尝试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/")

疑问

  1. 假设Dask通过磁盘排序输出多份文件,后续用read_csv读取这些文件时,是否能保留顺序?
  2. 如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 09:35:21