【问题标题】:How do you perform basic joins of two RDD tables in Spark using Python?如何使用 Python 在 Spark 中执行两个 RDD 表的基本连接?
【发布时间】:2015-09-24 06:18:21
【问题描述】:

您将如何使用 python 在 Spark 中执行基本连接?在 R 中,您可以使用 merg() 来执行此操作。在 spark 上使用 python 的语法是什么:

  1. 内连接
  2. 左外连接
  3. 交叉连接

有两个表 (RDD),每个表都有一个列,每个列都有一个公共键。

RDD(1):(key,U)
RDD(2):(key,V)

我认为内部连接是这样的:

rdd1.join(rdd2).map(case (key, u, v) => (key, ls ++ rs));

对吗?我搜索了互联网,找不到一个很好的连接示例。提前致谢。

【问题讨论】:

    标签: python join apache-spark pyspark rdd


    【解决方案1】:

    可以使用PairRDDFunctions 或 Spark 数据帧来完成。由于数据帧操作受益于Catalyst Optimizer,第二个选项值得考虑。

    假设您的数据如下所示:

    rdd1 =  sc.parallelize([("foo", 1), ("bar", 2), ("baz", 3)])
    rdd2 =  sc.parallelize([("foo", 4), ("bar", 5), ("bar", 6)])
    

    使用 PairRDD:

    内连接:

    rdd1.join(rdd2)
    

    左外连接:

    rdd1.leftOuterJoin(rdd2)
    

    笛卡尔积(不需要RDD[(T, U)]):

    rdd1.cartesian(rdd2)
    

    广播加入(不需要RDD[(T, U)]):

    最后是 cogroup,它没有直接的 SQL 等效项,但在某些情况下很有用:

    cogrouped = rdd1.cogroup(rdd2)
    
    cogrouped.mapValues(lambda x: (list(x[0]), list(x[1]))).collect()
    ## [('foo', ([1], [4])), ('bar', ([2], [5, 6])), ('baz', ([3], []))]
    

    使用 Spark 数据帧

    您可以使用 SQL DSL 或使用 sqlContext.sql 执行原始 SQL。

    df1 = spark.createDataFrame(rdd1, ('k', 'v1'))
    df2 = spark.createDataFrame(rdd2, ('k', 'v2'))
    
    # Register temporary tables to be able to use `sparkSession.sql`
    df1.createOrReplaceTempView('df1')
    df2.createOrReplaceTempView('df2')
    

    内连接:

    # inner is a default value so it could be omitted
    df1.join(df2, df1.k == df2.k, how='inner') 
    spark.sql('SELECT * FROM df1 JOIN df2 ON df1.k = df2.k')
    

    左外连接:

    df1.join(df2, df1.k == df2.k, how='left_outer')
    spark.sql('SELECT * FROM df1 LEFT OUTER JOIN df2 ON df1.k = df2.k')
    

    交叉连接(Spark.2.0 - spark.sql.crossJoin.enabled for Spark 2.x 需要显式交叉连接或配置更改):

    df1.crossJoin(df2)
    spark.sql('SELECT * FROM df1 CROSS JOIN df2')
    

    df1.join(df2)
    sqlContext.sql('SELECT * FROM df JOIN df2')
    

    从 1.6(Scala 中的 1.5)开始,这些中的每一个都可以与 broadcast 函数结合使用:

    from pyspark.sql.functions import broadcast
    
    df1.join(broadcast(df2), df1.k == df2.k)
    

    执行广播加入。另见Why my BroadcastHashJoin is slower than ShuffledHashJoin in Spark

    【讨论】:

    • 请注意:笛卡尔实际上在 RDD(不是 PairRDD)上可用
    • df1.join(df2, df1.k == df2.k, joinType='left_outer') 你如何将多个逻辑输入到参数中? df1.k == df2.k | df1.k2 == df2.k2 ?
    • @paradox (df1.k == df2.k) | (df1.k2 == df2.k2) 但将其设为union 或融化并转换为等值连接会更有意义。
    • @zero323:很好的答案。我建议的唯一更改是从 2.0 版开始,他们将“joinType”更改为“how”。
    • 是否可以在连接条件中添加一个函数。假设我有一个检查两个字符串相似度并返回相似度百分比的函数。例如:df1.join(df2, stringFunction(df1.k ,df2.k) > 80, how='left_outer')
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-10-19
    • 2018-04-08
    • 2015-01-31
    • 2015-05-30
    • 1970-01-01
    • 2015-12-27
    相关资源
    最近更新 更多