PySpark:使用Pandas UDF实现月度重置的序列计数
使用PySpark + Pandas UDF实现月度重置的序列计数
我来帮你搞定这个需求!针对同一月份内递增、跨月重置的序列计数,咱们可以通过按月份分组 + Pandas UDF组内生成序列的方式来实现,下面是具体的步骤和代码示例:
步骤说明
- 从日期列中提取年月作为分组标识,把同一个月的数据归为一组;
- 对每个月内的数据按日期(或原数据顺序)排序,确保序列按时间递增;
- 用Pandas UDF为每个分组生成从1开始的连续递增序列,跨月后自动重置计数。
完整代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import col, year, month, concat_ws from pyspark.sql.types import IntegerType import pandas as pd # 初始化SparkSession spark = SparkSession.builder.appName("MonthlyResetSequence").getOrCreate() # 模拟你提供的数据集结构 data = [ ("2016-12-14", 0, 0, 0, 14, 0), ("2016-12-14", 0, 0, 0, 14, 0), ("2016-12-14", 0, 0, 0, 18, 0), ("2016-12-14", 0, 0, 0, 19, 0), ("2016-12-14", 0, 0, 0, 20, 0), ("2016-12-14", 0, 0, 0, 26, 0), ("2016-12-14", 0, 0, 0, 60, 0), ("2017-01-02", 0, 0, 0, 5, 0), ("2017-01-02", 0, 0, 0, 10, 0) ] columns = ["Date", "column", "column_2", "co_3", "column_4", "column_5"] df = spark.createDataFrame(data, columns) # 1. 提取年月作为分组键 df_with_month = df.withColumn( "year_month", concat_ws("-", year(col("Date")), month(col("Date"))) ) # 2. 按年月分组,用applyInPandas生成月度序列 result_df = df_with_month.groupBy("year_month")\ .applyInPandas( lambda group_df: group_df.assign(monthly_sequence=range(1, len(group_df)+1)), schema=df_with_month.schema.add("monthly_sequence", IntegerType()) ) # 展示关键列结果 result_df.select("Date", "year_month", "column_4", "monthly_sequence").show()
代码细节解释
- 提取年月分组:用
year()和month()拆分日期,拼接成year_month作为分组依据,确保同月份数据被归为一组; - Pandas UDF灵活处理:
applyInPandas允许直接对每个分组的Pandas DataFrame操作,调用assign生成从1开始的序列,跨月后因为分组不同,计数会自动重置; - 排序注意事项:如果你的原始数据不是按日期排序的,建议在分组前加上
.orderBy("Date"),保证序列严格按时间递增。
输出示例
+----------+----------+--------+----------------+ | Date|year_month|column_4|monthly_sequence| +----------+----------+--------+----------------+ |2016-12-14| 2016-12| 14| 1| |2016-12-14| 2016-12| 14| 2| |2016-12-14| 2016-12| 18| 3| |2016-12-14| 2016-12| 19| 4| |2016-12-14| 2016-12| 20| 5| |2016-12-14| 2016-12| 26| 6| |2016-12-14| 2016-12| 60| 7| |2017-01-02| 2017-1| 5| 1| |2017-01-02| 2017-1| 10| 2| +----------+----------+--------+----------------+
内容的提问来源于stack exchange,提问作者user1157751
相关产品推荐
相关产品推荐

