如何编写PySpark pivot操作对应的纯SQL等价查询语句
PySpark pivot操作等价SQL实现方案
你需要补全的PIVOT部分完整代码如下,运行后得到的结果和PySpark API的groupBy('id').pivot('name').max()完全一致:
ds.createOrReplaceTempView('ds') dp = spark.sql(""" SELECT * FROM ds PIVOT (MAX(count) FOR name IN ('down', 'left', 'right', 'up') ) """).toPandas()
语法说明
FOR关键字后跟随要做行转列的源列名,这里就是你要转置的name列IN后列出name列的所有唯一枚举值,每个值用单引号包裹,列的顺序和PySpark默认pivot生成的列顺序对齐,无匹配值的位置会自动返回NULL,转pandas后会转为NaN,和原结果完全匹配
如果不想手动写死枚举值,可以先动态获取name列的所有唯一值再拼接SQL:
# 动态获取name列所有唯一值,按字典序排序和PySpark默认pivot的列顺序一致 name_list = [r['name'] for r in ds.select('name').distinct().orderBy('name').collect()] in_part = ', '.join(f"'{name}'" for name in name_list) dp = spark.sql(f""" SELECT * FROM ds PIVOT (MAX(count) FOR name IN ({in_part}) ) """).toPandas()
内容的提问来源于stack exchange,提问作者Russell Burdt
相关产品推荐
相关产品推荐

