【发布时间】:2016-08-03 05:43:27
【问题描述】:
我想在这两个 PySpark DataFrame 之间执行连接:
from pyspark import SparkContext
from pyspark.sql.functions import col
sc = SparkContext()
df1 = sc.parallelize([
['owner1', 'obj1', 0.5],
['owner1', 'obj1', 0.2],
['owner2', 'obj2', 0.1]
]).toDF(('owner', 'object', 'score'))
df2 = sc.parallelize(
[Row(owner=u'owner1',
objects=[Row(name=u'obj1', value=Row(fav=True, ratio=0.3))])]).toDF()
必须对对象的名称执行连接,即 objects 中的字段 name 用于 df2 和 object 用于 df1。
我可以在嵌套字段上执行 SELECT,如
df2.where(df2.owner == 'owner1').select(col("objects.value.ratio")).show()
但我无法运行此连接:
df2.alias('u').join(df1.alias('s'), col('u.objects.name') == col('s.object'))
返回错误
pyspark.sql.utils.AnalysisException: u"无法解析 '(objects.name = cast(object as double))' 由于数据类型 不匹配:'(objects.name = cast(object as double))' (数组和双精度).;"
有什么办法解决这个问题吗?
【问题讨论】:
标签: apache-spark dataframe join pyspark apache-spark-sql