【问题标题】:Spark - speed up cartesian productSpark - 加速笛卡尔积
【发布时间】:2019-06-06 02:30:00
【问题描述】:

是否有适当的方法来加速两个数据集之间的笛卡尔积?我正在使用 350k 个元素的数据集,我想获得它的组合(n 取 2)。

我使用经典策略在 Spark 中查找 2-组合:

words_comb = dataset.cartesian(dataset).filter(lambda x: x[0] < x[1])

我正在使用 Databricks 框架,它需要超过 45 分钟才能完成(火花驱动程序在 45 分钟时停止在 Databricks 中......)。我们都同意这样一个事实,即这个特定问题的瓶颈是数据集的笛卡尔积,其时间复杂度为 O(n^2)。

有没有办法改善这种情况?有没有更好的方法来解决这个问题?

(谢谢)

【问题讨论】:

  • databricks 论坛上有一个article 可以帮助您。
  • @Jeremy 谢谢。
  • 问题将是 O(n^2) 无论你如何孤立地这样做,这个问题有点毫无意义。生成所有组合建议使用蛮力方法,因此我将专注于改进整体算法。
  • 我同意@user8371915
  • 你想用笛卡尔积做什么?降低复杂性的一般方法是确定(根据您的用例)哪些值不需要生成,并在调用笛卡尔方法之前将它们过滤掉。

标签: apache-spark pyspark combinations


【解决方案1】:

我需要构建一个图 G=(N,F),其中 F(n) 是一个函数,它以具有 edit_distance(n,s) = 1 的单词 s 的子集 S 作为图像。要做到这一点,我hae 从所有单词组合开始,然后我依次过滤了所有不满足 edit_distance = 1 约束的单词对。

你的方法效率低得离谱。 n 个词的平均长度字符串 s 或多或少 O(n2 s2) (n2 edit_distance 电话)。同时您的数据很小(根据the comment 为4.1MB),分布式处理及其开销并不是很有用。你应该重新考虑你的方法。

我的建议是使用高效的查找结构(例如Trie 或BWT),它可以促进高效的不匹配搜索。使用整个数据集构建索引,如果需要,使用线程来并行化搜索。

【讨论】:

    【解决方案2】:
    It's possible to get rid of Cartesian product, without using a special data structure.
    
    I will demonstrate the method by an example with pyspark
    
    Suppose you have a data set of (user, query).
    
    input:
    
    users_queries_rdd = sc.parallelize([
         ('u1', 'q1'), ('u1', 'q2'), ('u1', 'q3'), ('u1', 'q4'),
         ('u2', 'q2'), ('u2', 'q4'), ('u2', 'q5'),
         ('u3', 'q1'), ('u3', 'q2'), ('u3', 'q4')
     ])
    
    You would like to count the occurrences of 2 queries for different users.
    
    expected output:
    [(('q4', 'q2'), 3),
     (('q5', 'q2'), 1),
     (('q3', 'q1'), 1),
     (('q5', 'q4'), 1),
     (('q4', 'q1'), 2),
     (('q3', 'q2'), 1),
     (('q2', 'q1'), 2),
     (('q4', 'q3'), 1)]
    

    方法 1 - 使用笛卡尔积:

    pair_queries_count_rdd = users_queries_rdd\
    .cartesian(users_queries_rdd)\
    .filter(lambda line: line[0] > line[1])\
    .filter(lambda line: line[0][0] == line[1][0])\
    .map(lambda line: (line[0][1], line[1][1]))\
    .map(lambda line: (line, 1))\
    .reduceByKey(add)
    

    方法 2 - 摆脱笛卡尔积:

    pair_queries_count_rdd_no_cartesian = users_queries_rdd\
                            .map(lambda line: (line[0], [line[1]]))\
                            .reduceByKey(add)\
                            .map(lambda line: tuple(combinations(line[1], 2)))\
                            .flatMap(lambda line: [(x, 1) for x in line])\
                            .reduceByKey(add)
    

    解释:

    方法一:

    1.1 .cartesian(users_queries_rdd)
    在 rdd 与自身之间创建笛卡尔积。 它会生成 n^2 个组合。

    1.2 .filter(lambda line: line[0] > line[1])
    如果包含 (q1,q2),则 (q2,q1) 将被过滤掉。

    1.3 .filter(lambda line: line[0][0] == line[1][0])
    按用户分组查询(聚合)。

    1.4 .map(lambda line: (line[0][1], line[1][1]))
    省略用户列。只保留查询对。

    1.5 .map(lambda line: (line, 1))
    每行将 (q[i], q[j]) 映射到 ((q[i], q[j]), 1)

    1.6 .reduceByKey(add)
    计算每个查询对的出现次数。

    方法二:

    2.1 .map(lambda line: (line[0], [line[1]]))
    将每一行的 (u[i], q[j]) 映射到 (u[i], [q[j]]) (查询被封装到一个列表中)

    2.2 .reduceByKey(add)
    为每个用户创建一个包含所有查询的列表。

    2.3 .map(lambda line: tuple(combinations(line[1], 2)))
    对于每个用户,省略用户列,并创建他们所有查询的组合

    2.4 .flatMap(lambda line: [(x, 1) for x in line])
    平面映射将平面所有键: (q[i],q[j]) 映射到 ((q[i],q[j]), 1)

    2.5 .reduceByKey(add)
    计算每个查询对的出现次数。

    我们假设每个用户的查询数量相对较少(一个常数)。 所以方法2的效率是O(n)

    这个假设是必不可少的,因为每个用户都使用了 iter.combinations 函数。

    在大多数实际情况下,方法 2 更有效。

    【讨论】:

      猜你喜欢
      • 2015-07-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多