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

求助:将Pandas的DataFrame分组聚合代码转换为PySpark实现

等价的PySpark实现方案

直接上与你提供的Pandas代码功能完全一致的PySpark代码:

from pyspark.sql import functions as F

df = df.groupBy('col').agg(
    F.first('col').alias('col'),
    F.concat_ws(', ', F.collect_list('country')).alias('country')
)

对应逻辑说明:

  • 原Pandas代码里的groupby('col', as_index=False),在PySpark中直接用groupBy('col')即可——PySpark分组后不会将分组列转为索引,默认就保持列结构。
  • 针对agg({'col': 'first'}),PySpark用F.first('col').alias('col')实现:first函数提取分组内col列的首个值,alias保证输出列名还是col。
  • 针对{'country': ', '.join},PySpark需要两步组合:先用collect_list把分组内的所有country值收集成列表,再用concat_ws(指定分隔符为, )把列表元素拼接成字符串,最后用alias保留列名country。

额外注意:

如果你的业务需要严格保证first取到的是原数据中的首个出现值,由于PySpark是分布式计算,默认不保证行顺序,建议先添加自增ID排序后再分组:

from pyspark.sql import functions as F

df = df.withColumn("row_id", F.monotonically_increasing_id()) \
       .orderBy("row_id") \
       .groupBy('col') \
       .agg(
           F.first('col').alias('col'),
           F.concat_ws(', ', F.collect_list('country')).alias('country')
       ) \
       .drop("row_id")

内容的提问来源于stack exchange,提问作者Rao Lokeswar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 16:10:30