【问题标题】:Pyspark: How to get a unique value pairs with reduceByKey()? (without using Distinct() method)Pyspark:如何使用 reduceByKey() 获得唯一值对? (不使用 Distinct() 方法)
【发布时间】:2018-02-07 20:10:05
【问题描述】:

我正在尝试从 [1,2,3] 的多个列中获取唯一值对。 数据量很大,有多个文件(总大小约1TB)。

我只想过滤带有“client”字符串的行,并 grep 每个文件中唯一的 [1,2,3] 列。 我首先使用了 tuple 和 Distinct() 函数,但是进程因 Java 内存错误而停止。

if __name__ == "__main__":
    sc=SparkContext(appName="someapp")
    cmd = 'hdfs dfs -ls /user/path'.split()
    files = subprocess.check_output(cmd).strip().split('\n')
    rdds=[]
    for ff in files[1:]:
        rdd=sc.textFile(ff.split()[-1])
        rdd2=rdd.filter(lambda x: "client" in x.lower())
        rdd3=rdd2.map(lambda x: tuple(x.split("\t")[y] for y in [1,2,3]))
        rdd4=rdd3.distinct()
        rdds.append(rdd4)

     rdd0=sc.union(rdds)
     rdd0.collect()
     rdd0.saveAsTextFile('/somedir')

所以我尝试了另一个使用reduceByKey() 方法的脚本,效果很好。

if __name__ == "__main__":
        sc=SparkContext(appName="someapp")
        cmd = "hdfs dfs -ls airties/eventU".split()
        files = subprocess.check_output(cmd).strip().split('\n')
        rdds=[]
        for ff in files[1:]:
                rdd=sc.textFile(ff.split()[-1])
                rdd2=rdd.filter(lambda x: "client" in x.lower())
                rdd3=rdd2.map(lambda x: ','.join([x.split("\t")[y] for y in [1,2,3]]))
                rdds.append(rdd3)
        rdd0=sc.union(rdds)
        rddA=rdd0.map(lambda x: (x,1)).reduceByKey(lambda a,b: a+b)
        rddA.collect()
        rddA.saveAsTextFile('/somedir')

但我试图理解为什么Distinct() 不能很好地工作,但reduceByKey() 方法有效。 distinct() 不是找到唯一值的正确方法吗?

还试图找到一种更好的方法来优化多个文件的处理,在每个文件中找到唯一值对并聚合它们。当每个文件都包含专有内容时,我只需要对每个文件应用唯一性并在最后一步中完全聚合。但似乎我当前的代码给系统带来了太多开销。

数据是这样的:大量冗余

+-----+---+------+----------+
|1    |2  |3     |4         |
+-----+---+------+----------+
|    1|  1|     A|2017-01-01|
|    2|  6|client|2017-01-02|
|    2|  3|     B|2017-01-02|
|    3|  5|     A|2017-01-03|
|    3|  5|client|2017-01-03|
|    2|  2|client|2017-01-02|
|    3|  5|     A|2017-01-03|
|    1|  3|     B|2017-01-02|
|    3|  5|client|2017-01-03|
|    3|  5|client|2017-01-04|
+-----+---+------+----------+

数据是这样的:大量冗余

+-----+---+------+
|1    |2  |3     |
+-----+---+------+
|    2|  6|client|
|    3|  5|client|
|    2|  2|client|
|    3|  5|client|
|    3|  5|client|
+-----+---+------+

第 3 列是多余的,但仅以此作为示例场景。

【问题讨论】:

  • 您能否展示一小部分数据样本和所需的输出?更多信息:stackoverflow.com/questions/48427185/…
  • 是样本输出还是输入?字符串'client' 在哪里?另外,您是否考虑过使用数据框?
  • 我已经更新了这个例子。我相信使用数据框可能会导致过多的开销。

标签: pyspark unique distinct aggregation reduce


【解决方案1】:

distinct 使用reduceByKey 实现。实际上,您只是重新实现了distinct(几乎)与currently implemented 完全相同。

但是,您的 2 个代码 sn-ps 不相等。在

首先sn-p它会

  • 处理 RDD
  • 保存 RDD 的不同元素
  • 将 RDD 附加到列表中,以便稍后创建聚合 RDD

第二个sn-p它会

  • 处理 RDD
  • 将 RDD 附加到列表中,以便稍后创建聚合 RDD
  • 在聚合 RDD 中保存不同的元素

如果不同文件中存在重复行,它们将在第一个 sn-p 中重复,而不是第二个。这可能是您内存不足的原因。请注意,这里不需要通过 collect 调用将所有记录提供给驱动程序,这会严重影响性能

【讨论】:

    猜你喜欢
    • 2021-10-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-03-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多