【问题标题】:Joining PySpark DataFrames on nested field在嵌套字段上加入 PySpark DataFrames
【发布时间】: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


    【解决方案1】:

    由于您要匹配和提取特定元素,最简单的方法是explode该行:

    matches = df2.withColumn("object", explode(col("objects"))).alias("u").join(
      df1.alias("s"),
      col("s.object") == col("u.object.name")
    )
    
    matches.show()
    ## +-------------------+------+-----------------+------+------+-----+
    ## |            objects| owner|           object| owner|object|score|
    ## +-------------------+------+-----------------+------+------+-----+
    ## |[[obj1,[true,0.3]]]|owner1|[obj1,[true,0.3]]|owner1|  obj1|  0.5|
    ## |[[obj1,[true,0.3]]]|owner1|[obj1,[true,0.3]]|owner1|  obj1|  0.2|
    ## +-------------------+------+-----------------+------+------+-----+
    

    另一种但效率非常低的方法是使用array_contains

    matches_contains = df1.alias("s").join(
      df2.alias("u"), expr("array_contains(objects.name, object)"))
    

    这是无效的,因为它将扩展到笛卡尔积:

    matches_contains.explain()
    ## == Physical Plan ==
    ## Filter array_contains(objects#6.name,object#4)
    ## +- CartesianProduct
    ##    :- Scan ExistingRDD[owner#3,object#4,score#5] 
    ##    +- Scan ExistingRDD[objects#6,owner#7]
    

    如果数组的大小相对较小,则可以生成array_contains 的优化版本,正如我在这里展示的那样:Filter by whether column value equals a list in spark

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2014-11-11
      • 1970-01-01
      • 2021-09-29
      • 1970-01-01
      • 2015-07-18
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多