【发布时间】:2022-10-05 23:20:59
【问题描述】:
我正在尝试使用 Snowpark 将 Snowflake 中的一些季度数据重新采样为每日数据,我有一些代码可以在 PySpark 中完成此操作;但是,似乎函数 \"explode()\" 在 Snowpark 中不支持。
# define function to create date range
def date_range(t1, t2, step=60*60*24):
\"\"\"Return a list of equally spaced points between t1 and t2 with stepsize step.\"\"\"
return [t1 + step*x for x in range(int((t2-t1)/step)+1)]
def resample(df, date_column=\'REPORTING_DATE\', groupby=\'ID\'):
# define udf
date_range_udf = udf(date_range)
# obtain min and max of time period for each group
df_base = df.groupBy(groupby)\\
.agg(F.min(date_column).cast(\'integer\').alias(\'epoch_min\')).select(\'epoch_min\', F.current_timestamp().cast(\'integer\').alias(\'epoch_max\'))
# generate timegrid and explode
df_base = df_base.withColumn(date_column, F.explode(date_range_udf(\"epoch_min\", \"epoch_max\")))\\
.drop(\'epoch_min\', \'epoch_max\')
# convert epoch to timestamp
df_base = df_base.withColumn(date_column, F.date_format(df_base[date_column].cast(dataType=T.TimestampType()), \'yyyy-MM-dd\')).orderBy(date_column, ascending=True)
# outer left join on reporting_date to resample data
df = df_base.join(df, [date_column], \'leftouter\')
# window for forward fill
window = Window.orderBy(date_column).partitionBy(groupby).rowsBetween(Window.unboundedPreceding, Window.currentRow)
# apply forward fill to all columns
for column in df.columns:
df = df.withColumn(column, F.last(column, ignorenulls=True).over(window))
return df
有人可以建议一个替代方案/提供一个示例代码来帮助我。谢谢 :)
标签: python dataframe resample snowpark