开始使用免费开始使用

将窗口函数从点号表示法转换为 SQL

我们要在列车时刻表中新增一列,使每一行都包含列车到达下一站所需的分钟数。

  • 我们有一个 dataframe df,其中 df.columns == ['train_id', 'station', 'time']
  • df 已注册为名为 schedule 的 SQL 表。
  • 下面的窗口函数查询使用点号表示法。它会生成一个新的 dataframe dot_df
window = Window.partitionBy('train_id').orderBy('time')
dot_df = df.withColumn('diff_min', 
                    (unix_timestamp(lead('time', 1).over(window),'H:m') 
                     - unix_timestamp('time', 'H:m'))/60)

请注意 unix_timestamp 函数的用法,它等价于 SQL 中的 UNIX_TIMESTAMP 函数。

请留意示例代码中的脚手架结构。按照脚手架格式书写答案,可以避免因格式问题被误判为错误提交。

本练习是课程的一部分

Python 中的 Spark SQL 入门

查看课程

练习说明

  • 创建一个 SQL 查询,使其结果与 dot_df 完全一致。请按照脚手架(即占位下划线 _____)的格式来书写查询。

交互式实操练习

通过完成这段示例代码来试试这个练习。

# Create a SQL query to obtain an identical result to dot_df
query = """
SELECT *, 
(____(____(time, 1) ____ (____ BY train_id ____ BY time),'H:m') 
 - ____(time, 'H:m'))/60 AS diff_min 
FROM schedule 
"""
sql_df = spark.sql(query)
sql_df.show()
编辑并运行代码