Spark DataFrame实现count_row_asc列:统计Remarks为Yes的员工并生成递增序号
实现分组递增序号
count_row_asc的方案 要生成需求中的count_row_asc列,我们可以利用PySpark的**窗口函数(Window)**实现分组内的序号递增。结合你已经完成的count_row_yes统计逻辑,完整实现步骤如下:
完整代码实现
from pyspark.sql import Window from pyspark.sql.functions import col, count, when, row_number # 1. 复用你已实现的:统计每个员工的Yes记录数 yes_count_df = df1.groupBy("Employeed_ID").agg( count(when(col("Remarks") == "Yes", True)).alias("count_row_yes") ) # 2. 关联原始表与统计结果,筛选出Remarks为Yes的记录 filtered_df = df1.join(yes_count_df, on="Employeed_ID", how="inner") \ .filter(col("Remarks") == "Yes") # 3. 定义窗口规则:按员工ID分组,生成递增序号 # 若需要按特定顺序排序(比如原始数据行序),可替换orderBy的字段 window_spec = Window.partitionBy("Employeed_ID").orderBy(col("Employeed_ID")) # 4. 添加递增序号列 final_df = filtered_df.withColumn("count_row_asc", row_number().over(window_spec)) # 调整列顺序,匹配期望数据集结构 final_df = final_df.select( "Employeed_ID", "Remarks", "count_row", "count_row_yes", "count_row_asc" ) # 查看最终结果 final_df.show()
关键说明
row_number().over(window_spec)会在每个Employeed_ID分组内,为每条符合条件的记录生成从1开始的连续递增序号。- 如果需要按原始数据的行顺序生成序号,可在
orderBy中指定原始数据的行标识字段(若存在),否则按Employeed_ID排序即可满足需求。
内容的提问来源于stack exchange,提问作者kimchigirl
相关产品推荐
相关产品推荐

