【问题标题】:how to use mix two RDD with python如何将两个RDD与python混合使用
【发布时间】:2018-07-15 23:22:46
【问题描述】:

我有以下两个RDD,第一个是:

training2 = training.map(lambda x:(x[0],(x[1],x[2])))
training2.collect()

#[(u'1', (u'4298118681424644510', u'7686695')),
# (u'1', (u'4860571499428580850', u'21560664')), 
# (u'1', (u'9704320783495875564', u'21748480')),
# (u'1', (u'13677630321509009335', u'3517124')),

第二个是:

user_id2 = user_id.map(lambda x:(x[0],(x[1],x[2])))
user_id2.collect()

#[(u'1', (u'1', u'5')),
# (u'2', (u'2', u'3')),
# (u'3', (u'1', u'5')),
# (u'4', (u'1', u'3')),
# (u'5', (u'2', u'1')),

在两个RDD中,参数u'1',u'2'...表示用户ID,所以我需要按键混合两个RDD,输出必须为每个键组合,如下所示:

u'1', (u'1', u'5', u'4298118681424644510', u'7686695')

【问题讨论】:

  • 您可以使用rdd.toDF()将两个rdds转换为spark DataFrame,然后使用join()将它们组合起来。

标签: apache-spark join lambda pyspark rdd


【解决方案1】:

如何添加两个rdd并使用aggregateByKey(self, zeroValue, seqFunc, combFunc, numPartitions=None)

您也可以使用reduceByKeygroupByKey

例如

zero_value=set()
def seq_op(x, y) :
    x.add(y)
    return x

def comb_op(x, y) :
    return x.union(y)

numbers = sc.parallelize([0,0,1,2,5,4,5,5,5]).map(lambda x : ["Even" if (x % 2 == 0) else "Odd", x])
numbers.collect()

numbers.aggregateByKey(zero_value, seq_op, comb_op).collect()
# results looks like [("Even", {0, 2, 4,}), ....]

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-11-30
    • 1970-01-01
    • 2015-12-23
    • 2019-04-24
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多