使用Python的Hadoop MapReduce处理Pandas DataFrame触发KeyError
Let's break down why this is happening and how to fix it:
Root Cause
When you run mapper.py locally, sys.stdin is the full CSV file (including the header row with column names like 'Time'). Pandas reads this correctly and recognizes the columns.
But in Hadoop Streaming, each mapper task gets a chunk/slice of the input CSV—Hadoop doesn't automatically pass the header row to every mapper. Some mapper processes will start reading from a data row instead of the header, so Pandas ends up creating a DataFrame with default column names (0,1,2,...) instead of your actual CSV columns. That's why you get KeyError: 'Time'.
Solutions
1. Manually Specify Column Names in Pandas
If you know the exact column names and order of your CSV, explicitly define them when reading the input. This bypasses the need for the header row in each mapper's input.
First, get your CSV's full column list (e.g., from the original file's header):
# Replace this with your actual CSV column names in order csv_columns = ["Column1", "Time", "Column3", ..., "Summary", ...]
Then modify your mapper's pd.read_csv line:
df = pd.read_csv(sys.stdin, header=None, names=csv_columns)
Important: You need to remove the header row from your input Reviews.csv before uploading it to Hadoop. Otherwise, the first mapper will treat the header as a data row.
2. Use a Separate Header File (More Robust)
If you don't want to modify the input CSV, upload the header as a separate file and load it in your mapper:
Step 1: Save the CSV header to a file
Extract the first line of Reviews.csv into headers.txt locally.
Step 2: Update your Hadoop command to include the header file
bin/hadoop jar share/hadoop/tools/lib/hadoop-streaming-3.1.0.jar \ -file /mypath/mapper.py -mapper 'python /mypath/mapper.py' \ -file /mypath/reducer.py -reducer 'python /mypath/reducer.py' \ -file /mypath/headers.txt \ -input /user/andreone/input/Reviews.csv -output /user/andreone/output/out_1
Step 3: Modify the mapper to load the header and filter out header rows
import sys import string import pandas as pd # Load the header from the uploaded file with open('headers.txt', 'r') as f: csv_columns = f.readline().strip().split(',') # Read input with explicit column names df = pd.read_csv(sys.stdin, header=None, names=csv_columns) # Filter out any rows that match the header (in case a mapper gets the header row) df = df[~df.apply(lambda x: x.tolist() == csv_columns, axis=1)] # Rest of your processing code remains the same df['Time'] = pd.to_datetime(df['Time'], unit='s').apply(lambda x : x.year) df['Summary'] = df['Summary'].str.lower() df['Summary'] = df['Summary'].str.replace('[{}]'.format(string.punctuation), '') for index, row in df.iterrows(): key = '' key += str(row.iloc[7]) key += '-' for word in str(row.iloc[8]).split(): key += word print('{}\t{}'.format(key, 1)) key = key.replace(word, '')
3. Ditch Pandas for Streamed Row-by-Row Processing (Best for Hadoop)
Pandas is great for batch processing, but Hadoop Streaming is designed for streaming, row-by-row operations. Using Pandas to load entire chunks can cause memory issues with large datasets, and it's prone to this header problem.
Here's a revised mapper that processes rows one at a time (using Python's built-in csv module to handle complex CSV formatting like quoted fields):
import sys import string import csv from datetime import datetime def main(): reader = csv.reader(sys.stdin) # Skip the header row header = next(reader) # Get indices of the columns we need time_col_idx = header.index('Time') summary_col_idx = header.index('Summary') # row.iloc[7] refers to the 8th column (0-indexed) base_key_col_idx = 7 for row in reader: if not row: continue # Process Time column to get year time_seconds = int(row[time_col_idx]) year = datetime.fromtimestamp(time_seconds).year # Clean Summary column summary = row[summary_col_idx].lower() # Remove punctuation summary = summary.translate(str.maketrans('', '', string.punctuation)) # Build keys and emit output base_key = f"{row[base_key_col_idx]}-" for word in summary.split(): full_key = base_key + word print(f"{full_key}\t1") if __name__ == "__main__": main()
This approach is more efficient for Hadoop's distributed model, avoids header issues entirely, and handles edge cases like commas inside quoted CSV fields.
Final Checks
- Verify that your input
Reviews.csvdoesn't have duplicate or malformed rows that might confuse the column parsing. - If you're using Python 2 vs 3 in Hadoop vs your local machine, make sure there are no version incompatibilities (e.g.,
str.translatesyntax changes).
内容的提问来源于stack exchange,提问作者andy

