这是解决这个问题的方法:
from datetime import datetime
from datetime import timedelta
from pyspark.sql.types import *
df = spark.createDataFrame([(2018, 'April', '01 W'),
(2018, 'April', '02 W'),
(2018, 'April', '03 W'),
(2018, 'April', '04 W'),
(2018, 'May', '01 W'),
(2018, 'May', '02 W'),
(2018, 'May', '03 W'),
(2018, 'May', '04 W'),
(2018, 'June', '01 W')
],
["Year", "Month", "Weeks"])
df = df.withColumn('week_number', F.regexp_extract(df['Weeks'], r'(\d+) ',1).cast(IntegerType()))
md = {'April':'04', 'May':'05', 'June':'06'}
df = df.withColumn('month_number', F.udf(lambda r: md[r])(df['Month']))
df = df.withColumn('yyyymm', F.concat_ws('-', df['Year'], df['month_number']))
df = df.withColumn('first_date', F.to_date(df['yyyymm'], 'yyyy-MM'))
df = df.withColumn('first_date', F.date_sub(df['first_date'], 1))
df = df.withColumn('first_date', F.next_day(df['first_date'], 'Fri'))
df = df.withColumn('date', F.lit(''))
df.show()
@pandas_udf(df.schema, PandasUDFType.GROUPED_MAP)
def _calc_fri(pdf):
s = pd.to_datetime(pdf['first_date'], format = '%Y-%m-%d')
days = s + pd.to_timedelta((pdf['week_number']-1)*7, unit='day')
pdf['date'] = days.dt.strftime("%Y-%m-%d")
return pdf
df = df.groupby(['Year', 'Month']).apply(_calc_fri).orderBy(['Year', 'month_number', 'week_number'])
df.show()
输出:
+----+-----+-----+-----------+------------+-------+----------+----------+
|Year|Month|Weeks|week_number|month_number| yyyymm|first_date| date|
+----+-----+-----+-----------+------------+-------+----------+----------+
|2018|April| 01 W| 1| 04|2018-04|2018-04-06|2018-04-06|
|2018|April| 02 W| 2| 04|2018-04|2018-04-06|2018-04-13|
|2018|April| 03 W| 3| 04|2018-04|2018-04-06|2018-04-20|
|2018|April| 04 W| 4| 04|2018-04|2018-04-06|2018-04-27|
|2018| May| 01 W| 1| 05|2018-05|2018-05-04|2018-05-04|
|2018| May| 02 W| 2| 05|2018-05|2018-05-04|2018-05-11|
|2018| May| 03 W| 3| 05|2018-05|2018-05-04|2018-05-18|
|2018| May| 04 W| 4| 05|2018-05|2018-05-04|2018-05-25|
|2018| June| 01 W| 1| 06|2018-06|2018-06-01|2018-06-01|
+----+-----+-----+-----------+------------+-------+----------+----------+
猜猜你也可以把所有的工作都放到pandas_udf,或者使用udf,我个人会尽量少做任何udf的工作。