Spark中如何查找键值对值的交集并设为所有键值对的值?
Hey there! Let's break down how to solve your problem step by step. You're looking to compute the intersection of all value sets in your RDD and then assign that intersection as the value for every key. Here's how to do it properly in Spark:
Step 1: Compute the Global Intersection of All Value Sets
First, we need to extract all the sets from your RDD's values and compute their intersection. Since Spark RDDs have a values() method to isolate just the value part, we can use reduce() on that subset to iteratively calculate the intersection across all sets:
# Assuming your RDD is named 'original_rdd' global_intersection = original_rdd.values().reduce(lambda a, b: a.intersection(b))
This works because reduce() applies the lambda function (which takes two sets and returns their intersection) across every element in the values RDD. For your example input, this will result in Set(n2)—the common element present in all three sets.
Step 2: Reassign the Global Intersection to Every Key
Once we have the global intersection, we just need to map each key-value pair in the original RDD to keep the key and replace the value with our computed intersection:
result_rdd = original_rdd.map(lambda key_value: (key_value[0], global_intersection))
Running this will give you exactly the output you're expecting:(p1, Set(n2)), (p2, Set(n2)), (p3, Set(n2))
Why Your Previous reduce() Attempt Might Have Failed
If you tried using reduce() directly on the original key-value RDD, that's why it didn't work—reduce() on a key-value RDD tries to combine pairs using the provided function, which isn't what we need here. By first extracting the values with values(), we focus only on the sets we need to intersect.
And to address your note about no reduceByValue: that's correct, Spark doesn't have a reduceByValue method, but since we need a global operation (not per key), values().reduce() is the right approach here.
内容的提问来源于stack exchange,提问作者BuffaloSoldier

