PySpark groupBy后执行join报'GroupedData'无'join'属性如何解决
报错根本原因
你调用groupBy()方法后得到的是GroupedData类型对象,该类型仅支持count、sum、avg这类聚合操作,没有join方法,只有DataFrame才能执行join操作,因此触发AttributeError。
另外你的代码还存在两处隐性问题:
- 你通过Glue Catalog读取得到的是DynamicFrame,没有转成DataFrame就直接调用
select方法会报错 - 你
groupBy时用到了location_id字段,但前面select筛选字段时并没有保留该字段,运行时会提示字段不存在
解决方案
你的需求是分组后保留原表所有列,推荐使用窗口函数实现,无需额外join操作,示例代码如下:
from pyspark.sql.types import IntegerType from pyspark.sql.types import * from pyspark.sql.functions import * from pyspark.sql.window import Window # 读取DynamicFrame后先转成DataFrame retail_sales_transaction = glueContext.create_dynamic_frame.from_catalog( database="conform_main_mobconv", table_name="retail_sales_transaction" ).toDF() # 选择字段时要保留后续分组用到的location_id df_retail = retail_sales_transaction.select( "transaction_id","transaction_key","transaction_timestamp", "personnel_key","retail_site_id","personnel_role", "country_code","business_week","location_id" ) # 定义窗口,按你需要的分组字段分区 window_spec = Window.partitionBy("business_week", "location_id", "country_code", "personnel_key") # 在这里加你需要的聚合计算,比如统计每个分组的交易总数,原表所有字段都会保留 df_retail_with_agg = df_retail.withColumn("group_trans_count", count("transaction_id").over(window_spec)) # 如果需要去重取每个分组的唯一记录,可以配合dropDuplicates使用 # df_retail_distinct = df_retail.dropDuplicates(["business_week", "location_id", "country_code", "personnel_key"])
如果你一定要用groupBy+join的方式实现,参考代码如下:
# 先对分组结果执行聚合操作转成DataFrame group_df = df_retail.groupBy("business_week", "location_id", "country_code", "personnel_key")\ .agg(count("transaction_id").alias("group_trans_count")) # 再和原表join,注意join字段要和分组逻辑匹配,避免数据膨胀 result_df = df_retail.join(group_df, on=["business_week", "location_id", "country_code", "personnel_key"], how="left")
内容的提问来源于stack exchange,提问作者Sonam Garg
相关产品推荐
相关产品推荐

