PySpark:行数据转键值对实现及RDD读取数据输出异常排查
First, let's address the b... output: this happens because your RDD contains byte strings instead of regular Unicode strings. We'll decode each element to UTF-8 first. Then, we'll parse the structured text into proper key-value pairs where the key is a tuple of items and the value is the count.
Step 1: Read and Decode the File
First, load the file and convert each byte string to a regular string (adding a check ensures compatibility across different PySpark versions):
import re # Load the file and decode bytes to UTF-8 strings if needed data = sc.textFile("hdfs://h1:9000/data.txt").map(lambda line: line.decode("utf-8") if isinstance(line, bytes) else line)
Step 2: Parse Each Line into Key-Value Pairs
Next, we'll use regex to extract each item group and its count, then convert them into hashable tuple keys and integer values:
def parse_line(line): # Clean up any trailing invalid characters (like the extra ] in your example) cleaned_line = line.rstrip("] ") # Extract all (item group, count) pairs using regex matches = re.findall(r'\(\[(.*?)\],(\d+)\)', cleaned_line) # Convert each match to a key-value pair: (tuple of items, integer count) return [(tuple(items.split(", ")), int(count)) for items, count in matches] # Use flatMap to split multi-entry lines into individual key-value pairs key_value_rdd = data.flatMap(parse_line)
Step 3: Verify the Result
Now, when you collect the RDD, you'll see proper key-value pairs instead of byte strings:
print(key_value_rdd.collect())
Sample output:
[ (('Fire Grilled Steak Bowl', 'Veri Veri Teriyaki Marinade & Sauce'), 1), (('Fat Free Strawberry Yogurt', 'Strawberry Banana Juice'), 1), (('Organic Unsweetened & Salt Free Sunflower Seed Butter', 'Organic Broccoli Crowns'), 1), (('Roma Tomato', 'Vanilla Soymilk'), 3), (('15 Calories Light Lemonade', 'Lemon-Lime 12 Pack Soda'), 1) ]
Explanation of the Parsing Logic
- Regex:
r'\(\[(.*?)\],(\d+)\)'reliably matches each entry like([Item1, Item2],N), capturing the comma-separated items and numeric count. - FlatMap: We use
flatMapinstead ofmapbecause your example has multiple entries in one line—this ensures each entry becomes a separate element in the RDD. - Tuple Keys: Converting the item list to a tuple makes it hashable, which is required for operations like
groupByKeyorreduceByKeylater.
If your file has one entry per line instead of multiple entries, simplify the parsing function:
def parse_single_entry(line): cleaned_line = line.strip("()[] ") items_part, count_part = cleaned_line.split("], ") items = tuple(items_part.split(", ")) count = int(count_part) return (items, count) key_value_rdd = data.map(parse_single_entry)
This resolves both the byte string issue and converts your data into usable key-value pairs for further Spark operations.
内容的提问来源于stack exchange,提问作者learning_dev

