如何优化Pandas分组日期填充逻辑并在PySpark中实现?
Pandas性能优化与PySpark实现方案
需求说明
现有包含id、place、date、value列的DataFrame,value取值为core和not core;另有月度最后日期列表dates。需完成:
- 为每个
id-place分组填充dates中的缺失日期 - 生成
status列:当**上月value为core且当前value不为core**时,status为0,否则为1
示例输入DataFrame
id place value date 1 A core 2023-08-31 1 A not core 2023-09-30 1 A core 2023-11-30 2 A core 2023-10-30 2 A core 2023-11-30 2 B not core 2023-07-31 2 B core 2023-10-31
日期列表
['2023-07-31', '2023-08-31', '2023-09-30', '2023-10-31', '2023-11-30']
期望输出DataFrame
id place value date prev_month_value status 1 A NONE 2023-07-31 NONE 1 1 A core 2023-08-31 NONE 1 1 A not core 2023-09-30 core 0 1 A NONE 2023-10-31 not core 1 1 A core 2023-11-30 NONE 1 2 A NONE 2023-07-31 NONE 1 2 A NONE 2023-08-31 NONE 1 2 A NONE 2023-09-30 NONE 1 2 A core 2023-10-31 NONE 1 2 A core 2023-11-30 core 1 2 B not core 2023-07-31 NONE 1 2 B NONE 2023-08-31 not core 1 2 B NONE 2023-09-30 NONE 1 2 B core 2023-10-31 NONE 1 2 B NONE 2023-11-30 core 0
用户原Pandas方案因循环分组+逐行apply导致效率低下,以下是针对性优化方案:
1. Pandas性能与内存优化方案
原代码核心问题是循环遍历分组和逐行apply,这是Pandas中效率最低的操作模式。优化思路是用矢量化操作和全量笛卡尔积关联替代循环:
优化步骤
- 生成
id-place唯一组合与日期列表的笛卡尔积,构建全量时间序列框架 - 右关联原数据,缺失
value填充为'NONE' - 按
id-place分组后用shift(1)获取上月值,避免逐行处理 - 用矢量化布尔判断生成
status列
优化后代码
import pandas as pd # 1. 构建全量id-place-date框架 dates_df = pd.DataFrame({'date': dates}) id_place_unique = df[['id', 'place']].drop_duplicates() full_frame = id_place_unique.merge(dates_df, how='cross') # 2. 关联原数据并填充缺失值 result_df = full_frame.merge(df, on=['id', 'place', 'date'], how='left') result_df['value'] = result_df['value'].fillna('NONE') # 3. 获取分组内上月值 result_df['prev_month_value'] = result_df.groupby(['id', 'place'])['value'].shift(1) result_df['prev_month_value'] = result_df['prev_month_value'].fillna('NONE') # 4. 矢量化生成status列 result_df['status'] = 1 mask = (result_df['prev_month_value'] == 'core') & (result_df['value'] != 'core') result_df.loc[mask, 'status'] = 0 # 排序并重置索引 result_df = result_df.sort_values(['id', 'place', 'date']).reset_index(drop=True)
优化点说明
- 移除
groupby循环,用merge cross生成全量框架,减少内存碎片与循环开销 - 用
fillna和矢量化布尔判断替代apply,速度提升10~100倍(取决于数据规模) - 所有操作均为Pandas底层优化的矢量化运算,内存使用更高效
2. PySpark实现方案
PySpark不支持分组迭代,需用窗口函数和笛卡尔积关联实现逻辑:
实现步骤
- 将日期列表转为Spark DataFrame
- 获取
id-place唯一组合,与日期做笛卡尔积生成全量框架 - 左关联原数据,缺失
value填充为'NONE' - 用窗口函数
lag(1)获取分组内上月值 - 用
when/otherwise条件判断生成status列
PySpark代码
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.window import Window spark = SparkSession.builder.appName("DateFillStatus").getOrCreate() # 1. 转换日期列表为Spark DataFrame dates_df = spark.createDataFrame([(d,) for d in dates], ['date']) dates_df = dates_df.withColumn('date', F.to_date('date')) # 2. 生成全量id-place-date框架 id_place_unique = df.select('id', 'place').dropDuplicates() full_frame = id_place_unique.crossJoin(dates_df) # 3. 关联原数据并填充缺失值 result_df = full_frame.join(df, on=['id', 'place', 'date'], how='left') result_df = result_df.fillna({'value': 'NONE'}) # 4. 窗口函数获取上月值 window_spec = Window.partitionBy('id', 'place').orderBy('date') result_df = result_df.withColumn( 'prev_month_value', F.lag('value', 1).over(window_spec) ).fillna({'prev_month_value': 'NONE'}) # 5. 生成status列 result_df = result_df.withColumn( 'status', F.when( (F.col('prev_month_value') == 'core') & (F.col('value') != 'core'), 0 ).otherwise(1) ) # 排序输出 result_df = result_df.orderBy('id', 'place', 'date') result_df.show()
关键说明
- 用
crossJoin实现笛卡尔积,适配分布式场景下的全量时间序列生成 lag窗口函数替代shift,实现分组内的上月值获取fillna和when/otherwise均为Spark矢量化操作,适合大数据量处理
内容的提问来源于stack exchange,提问作者Nourless
相关产品推荐
相关产品推荐

