PySpark中如何修改Pair RDD的键?创建指定键Pair RDD遇问题求助
Hey there! I see where the issues are in your code—let's get your Pair RDD working correctly step by step.
First, let's break down the problems in your original code:
- You have a typo:
lamdashould belambda(missing a 'b') - Variable case mismatch: You used uppercase
XinX.split(",")but your parameter is lowercasex—this would cause an error or incorrect data extraction - The value part doesn't match your requirement: You're passing the entire row
xas the value, but you need(name, age)instead
Here's the corrected code:
First, let's assume you've already loaded your CSV into an RDD:
rdd = sc.textFile("your_csv_file.csv")
If your CSV has a header row (NAME,AGE,NATIONALITY), we should skip it first to avoid treating the header as data:
# Grab the header row header = rdd.first() # Filter out the header to keep only actual data data_rdd = rdd.filter(lambda row: row != header)
Now create the Pair RDD with NATIONALITY as the key and (name, age) as the value. We can optimize by splitting each row cleanly to avoid redundant operations:
# Split each row once, then map to key-value pairs t1 = data_rdd.map(lambda x: (x.split(",")[2], (x.split(",")[0], int(x.split(",")[1]))) )
Note: I converted the age to an integer with int()—this is optional but useful if you plan to do numerical operations on age later. If your age field has non-numeric values, skip this conversion.
Test it out:
You can now check the keys and values properly:
# Retrieve all nationalities (keys) print(t1.keys().collect()) # Retrieve all (name, age) pairs (values) print(t1.values().collect())
This should give you exactly the Pair RDD you're looking for, matching the functionality you'd get in Scala!
内容的提问来源于stack exchange,提问作者Vee JayBee

