Scala中使用reduceByKey时遭遇类型不匹配问题求助
reduceByKey Let's work through this type mismatch issue step by step to get your code running correctly.
First, let's clarify the type of your RDDs to understand the root problem:
p7is an RDD of type(String, (String, String, Int))— each value is a 3-element tuple:(firstString, secondString, count)
When using reduceByKey, the function you pass has a strict requirement: it must take two values of the same type as the RDD's value, and return a value of that exact same type.
Your current reduce function:
(a,b) => (a._1,(a._2, a._3+b._3))
Returns a (String, (String, Int)) — a tuple where the second element is a 2-element tuple. This doesn't match the original value type (String, String, Int), which is why you're seeing the type mismatch error.
The Corrected Code
Assuming your goal is to sum the third integer value while retaining the first two string values (e.g., keeping the strings from the first occurrence of each key), adjust the reduce function to return a 3-element tuple that matches the original value type:
val p8 = p7.reduceByKey( (a,b) => (a._1, a._2, a._3 + b._3) )
Customizing the Logic (If Needed)
If you want to handle the first two strings differently (like concatenating them instead of keeping the first occurrence), you can modify the function while maintaining the correct return type:
// Example: Concatenate the first two strings and sum the count val p8 = p7.reduceByKey( (a,b) => (a._1 + b._1, a._2 + b._2, a._3 + b._3) )
Here's how the corrected code runs in the Scala Shell:
scala> val p6 = sc.parallelize(List( ("a","b"),("b","c"))) p6: org.apache.spark.rdd.RDD[(String, String)] = ParallelCollectionRDD[10] at parallelize at <console>:24 scala> val p7 = p6.map(a => ((a._1+a._2), (a._1, a._2, 1))) p7: org.apache.spark.rdd.RDD[(String, (String, String, Int))] = MapPartitionsRDD[11] at map at <console>:26 scala> val p8 = p7.reduceByKey( (a,b) => (a._1, a._2, a._3 + b._3) ) p8: org.apache.spark.rdd.RDD[(String, (String, String, Int))] = ShuffledRDD[12] at reduceByKey at <console>:28 scala> p8.collect() res0: Array[(String, (String, String, Int))] = Array((ab,(a,b,1)), (bc,(b,c,1)))
内容的提问来源于stack exchange,提问作者Litchy

