【发布时间】:2017-02-07 17:54:14
【问题描述】:
我有以下 Spark DataFrame:
df = sql.createDataFrame([
(1, [
{'name': 'john', 'score': '0.8'},
{'name': 'johnson', 'score': '0.9'},
]),
(2, [
{'name': 'jane', 'score': '0.9'},
{'name': 'janine', 'score': '0.4'},
]),
(3, [
{'name': 'sarah', 'score': '0.2'},
{'name': 'sara', 'score': '0.9'},
]),
], schema=['id', 'names'])
Spark 正确推断架构:
root
|-- id: long (nullable = true)
|-- names: array (nullable = true)
| |-- element: map (containsNull = true)
| | |-- key: string
| | |-- value: string (valueContainsNull = true)
对于每一行,我想选择得分最高的名称。我可以使用 Python UDF 执行此操作,如下所示:
import pyspark.sql.types as T
import pyspark.sql.functions as F
def top_name(names):
return sorted(names, key=lambda d: d['score'], reverse=True)[0]['name']
top_name_udf = F.udf(top_name, T.StringType())
df.withColumn('top_name', top_name_udf('names')) \
.select('id', 'top_name') \
.show(truncate=False)
如您所愿,您会得到:
+---+--------+
|id |top_name|
+---+--------+
|1 |johnson |
|2 |jane |
|3 |sara |
+---+--------+
如何使用 Spark SQL 做到这一点?是否可以在没有 Python UDF 的情况下做到这一点,这样数据就不会在Python 和Java 之间序列化?1
1很遗憾,我运行的是 Spark 1.5,无法在 Spark 2.1 中使用 registerJavaFunction。
【问题讨论】:
标签: apache-spark pyspark apache-spark-sql pyspark-sql