代码之家  ›  专栏  ›  技术社区  ›  Sai

如何使用pyspark中其他Rdd元素的所有可能组合创建新的Rdd?

  •  2
  • Sai  · 技术社区  · 7 年前

    rdd1=sc.parallelize(['P','T','K'])
    rdd1.collect()
    ['P', 'T', 'K']
    

    现在我想创建一个新的 RDD2 使用所有可能的组合,如下文所示,使用新的RDD,即,除了相同的元素组合,如(p,p),(k,k),(t,t)。

    当我正在做时,我的预期输出

    RDD2.collect()
    
    [
        ('P'),('T'),('K'),
        ('P','T'),('P','K'),('T','K'),('T','P'),('K','P'),('K','T'),
        ('P','T','K'),('P','K','T'),('T','P','K'),('T','K','P'),('K','P','T'),('K','T','P')
    ]
    
    2 回复  |  直到 7 年前
        1
  •  3
  •   pault Tanjin    7 年前

    似乎您希望在您的应用程序中生成元素的所有排列 rdd 其中每行包含唯一的值。

    一种方法是首先创建一个helper函数来生成所需的长度组合 n :

    from functools import reduce
    from itertools import chain
    
    def combinations_of_length_n(rdd, n):
        # for n > 0
        return reduce(
            lambda a, b: a.cartesian(b).map(lambda x: tuple(chain.from_iterable(x))),
            [rdd]*n
        ).filter(lambda x: len(set(x))==n)
    

    从本质上讲,这个函数就可以了 你的笛卡尔积 rdd 只保留所有值都不同的行。

    我们可以测试一下 n = [2, 3] :

    print(combinations_of_length_n(rdd1, n=2).collect())
    #[('P', 'T'), ('P', 'K'), ('T', 'P'), ('K', 'P'), ('T', 'K'), ('K', 'T')]
    
    print(combinations_of_length_n(rdd1, n=3).collect())
    #[('P', 'T', 'K'),
    # ('P', 'K', 'T'),
    # ('T', 'P', 'K'),
    # ('K', 'P', 'T'),
    # ('T', 'K', 'P'),
    # ('K', 'T', 'P')]
    

    您想要的最终输出只是 union 这些中间结果与原始结果一致 rdd tuple s) 。

    rdd1.map(lambda x: tuple((x,)))\
        .union(combinations_of_length_n(rdd1, 2))\
        .union(combinations_of_length_n(rdd1, 3)).collect()
    #[('P',),
    # ('T',),
    # ('K',),
    # ('P', 'T'),
    # ('P', 'K'),
    # ('T', 'P'),
    # ('K', 'P'),
    # ('T', 'K'),
    # ('K', 'T'),
    # ('P', 'T', 'K'),
    # ('P', 'K', 'T'),
    # ('T', 'P', 'K'),
    # ('K', 'P', 'T'),
    # ('T', 'K', 'P'),
    # ('K', 'T', 'P')]
    

    要概括任何最大重复次数,请执行以下操作:

    num_reps = 3
    reduce(
        lambda a, b: a.union(b),
        [
            combinations_of_length_n(rdd1.map(lambda x: tuple((x,))), i+1) 
            for i in range(num_reps)
        ]
    ).collect()
    #Same as above
    

        2
  •  0
  •   Ali Yesilli    7 年前

    有几种方法。您可以运行一个循环,获取排列并将其存储在列表中,然后将列表转换为rdd

    >>> rdd1.collect()
    ['P', 'T', 'K']
    >>> 
    >>> l = []
    >>> for i in range(2,rdd1.count()+1):
    ...     x = list(itertools.permutations(rdd1.toLocalIterator(),i))
    ...     l = l+x
    ... 
    >>> rdd2 = sc.parallelize(l)
    >>> 
    >>> rdd2.collect()
    [('P', 'T'), ('P', 'K'), ('T', 'P'), ('T', 'K'), ('K', 'P'), ('K', 'T'), ('P', 'T', 'K'), ('P', 'K', 'T'), ('T', 'P', 'K'), ('T', 'K', 'P'), ('K', 'P', 'T'), ('K', 'T', 'P')]