【问题标题】:How to join three RDDs using the Python Core API (Apache Spark)?如何使用 Python Core API (Apache Spark) 加入三个 RDD?
【发布时间】:2021-06-05 12:00:27
【问题描述】:

我正在尝试通过 Apache Spark 使用 Python Core API 将你们的 RDD 连接在一起;但是,我没有运气尝试完成此操作。

目前,我有这三个具有共同属性的 RDD:

  • users_rdd: user_id
  • reviews_rdd:review_id、company_id 和 user_id
  • companies_rdd: company_id

现在,当将两个 RDD 连接在一起时,它可以正常工作,没有任何问题:

user_rev_rdd = (users_rdd
  .keyBy(lambda user: user['user_id'])
  .join(
      reviews_rdd.keyBy(lambda rev: rev['user_id'])
  )
)

虽然,为了将这三个结合在一起,我已经尝试过这个,但由于某种原因它根本不适合我:

user_rev_com_rdd = (users_rdd
  .keyBy(lambda user: user['user_id'])
  .join(
      reviews_rdd.keyBy(lambda rev: rev['user_id'])
  )
 .join(
      companies_rdd.keyBy(lambda com: com['company_id'])
  )
)

任何关于如何将我的所有三个 RDD 连接在一起的帮助都会非常有帮助,因为我不确定如何正确地做这样的事情。谢谢。

【问题讨论】:

    标签: python apache-spark join pyspark rdd


    【解决方案1】:

    第一次加入后,key为user_id,但你加入companies_rdd,key为company_id,所以加入key不正确。您需要将密钥更改为company_id,例如

    user_rev_com_rdd = (users_rdd
        .keyBy(lambda user: user['user_id'])
        .join(
            reviews_rdd.keyBy(lambda rev: rev['user_id'])
        )
        .map(lambda r: (r[1][1]['company_id'], r[1]))
        .join(
            companies_rdd.keyBy(lambda com: com['company_id'])
        )
    )
    

    要合并三个RDD中的元素并在加入后删除加入键,可以在末尾添加map

    user_rev_com_rdd = (users_rdd
        .keyBy(lambda user: user['user_id'])
        .join(
            reviews_rdd.keyBy(lambda rev: rev['user_id'])
        )
        .map(lambda r: (r[1][1]['company_id'], r[1]))
        .join(
            companies_rdd.keyBy(lambda com: com['company_id'])
        )
        .map(lambda r: (*r[1][0], r[1][1]))
    )
    

    【讨论】:

      猜你喜欢
      • 2019-03-25
      • 1970-01-01
      • 2014-05-13
      • 2015-06-15
      • 2017-06-27
      • 1970-01-01
      • 2015-07-05
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多