【问题标题】:PySpark isin functionPySpark isin 函数
【发布时间】:2017-06-08 20:04:40
【问题描述】:

我正在使用 PySpark 将我的旧 Python 代码转换为 Spark。

我想获得一个 PySpark 等价物:

usersofinterest = actdataall[actdataall['ORDValue'].isin(orddata['ORDER_ID'].unique())]['User ID']

actdataallorddata 都是 Spark 数据帧。

鉴于与之相关的缺点,我不想使用toPandas() 函数。

【问题讨论】:

    标签: apache-spark pyspark


    【解决方案1】:

    您的代码最直接的翻译是:

    from pyspark.sql import functions as F
    
    # collect all the unique ORDER_IDs to the driver
    order_ids = [x.ORDER_ID for x in orddata.select('ORDER_ID').distinct().collect()]
    
    # filter ORDValue column by list of order_ids, then select only User ID column
    usersofinterest = actdataall.filter(F.col('ORDValue').isin(order_ids)).select('User ID')
    

    但是,只有在“ORDER_ID”的数量肯定很小(可能

    如果“ORDER_ID”的数量很大,您应该使用广播变量将 order_ids 列表发送到每个执行程序,以便它可以与本地 order_ids 进行比较以加快处理速度。请注意,即使 'ORDER_ID' 很小,这也会起作用。

    order_ids = [x.ORDER_ID for x in orddata.select('ORDER_ID').distinct().collect()]
    order_ids_broadcast = sc.broadcast(order_ids)  # send to broadcast variable
    usersofinterest = actdataall.filter(F.col('ORDValue').isin(order_ids_broadcast.value)).select('User ID')
    

    有关广播变量的更多信息,请查看:https://jaceklaskowski.gitbooks.io/mastering-apache-spark/spark-broadcast.html

    【讨论】:

      【解决方案2】:
      • 如果两个数据框都很大,您应该考虑使用内连接作为过滤器:

        首先让我们创建一个包含我们想要保留的订单 ID 的数据框:

        orderid_df = orddata.select(orddata.ORDER_ID.alias("ORDValue")).distinct()
        

        现在让我们将它与我们的 actdataall 数据框连接起来:

        usersofinterest = actdataall.join(orderid_df, "ORDValue", "inner").select('User ID').distinct()
        
      • 如果您的订单 ID 的目标列表很小,那么您可以使用 furianpandit 帖子中提到的 pyspark.sql isin 函数,不要忘记在使用之前广播您的变量(spark 会将对象复制到每个节点,使其任务更快):

        orderid_list = orddata.select('ORDER_ID').distinct().rdd.flatMap(lambda x:x).collect()[0]
        sc.broadcast(orderid_list)
        

      【讨论】:

        【解决方案3】:

        所以,您有两个 spark 数据框。一个是actdataall,一个是orddata,然后使用下面的命令得到你想要的结果。

        usersofinterest  = actdataall.where(actdataall['ORDValue'].isin(orddata.select('ORDER_ID').distinct().rdd.flatMap(lambda x:x).collect()[0])).select('User ID')
        

        【讨论】:

          猜你喜欢
          • 1970-01-01
          • 1970-01-01
          • 2020-10-29
          • 2022-11-16
          • 2017-06-06
          • 2020-11-13
          • 2020-07-13
          • 1970-01-01
          • 1970-01-01
          相关资源
          最近更新 更多