求助:将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
相关产品推荐
相关产品推荐

