如何在MapReduce中使用多个CSV文件?Mapper代码异常求助
Hey there! Let's work through this MapReduce CSV merging problem together. It sounds like you're aiming to do a side join—using one CSV as the main job input, and loading the other directly in your Mapper to combine data via the shared BibNum key. This is a super common use case, but there are a few key pitfalls that might be breaking your code. Let's break down the fixes and best practices:
The biggest mistake people make here is loading the secondary CSV inside the map() method—this runs every time a record is processed, which is inefficient and can cause file access errors. Instead, load it once per Mapper task using the setup() method:
import org.apache.hadoop.io.LongWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; import java.io.BufferedReader; import java.io.FileReader; import java.io.IOException; import java.util.Arrays; import java.util.HashMap; import java.util.Map; public class CSVJoinMapper extends Mapper<LongWritable, Text, Text, Text> { private Map<String, String> secondaryCSVData = new HashMap<>(); @Override protected void setup(Context context) throws IOException, InterruptedException { // Fetch the path to your second CSV from the job configuration String secondaryCSVPath = context.getConfiguration().get("secondary.csv.path"); // Read the CSV and store relevant data in a map (key = BibNum, value = your desired columns) try (BufferedReader br = new BufferedReader(new FileReader(secondaryCSVPath))) { String line; // Skip header row if your CSV has one br.readLine(); while ((line = br.readLine()) != null) { // Split while preserving empty columns (the -1 flag is critical!) String[] fields = line.split(",", -1); if (fields.length >= 2) { // Adjust based on how many columns you need from the second CSV String bibNum = fields[0].trim(); // Combine the columns you want into a single string (e.g., col2,col3,col4) String combinedValues = String.join(",", Arrays.copyOfRange(fields, 1, fields.length)); secondaryCSVData.put(bibNum, combinedValues); } } } } @Override protected void map(LongWritable key, Text value, Context context) throws IOException, InterruptedException { // Parse the main input CSV (first file) String[] mainFields = value.toString().split(",", -1); if (mainFields.length >= 3) { // Match your main CSV's columns: BibNum, checkoutdatetime, itemtype String bibNum = mainFields[0].trim(); String checkoutDateTime = mainFields[1].trim(); String itemType = mainFields[2].trim(); // Look up the matching data from the secondary CSV String secondaryValues = secondaryCSVData.get(bibNum); if (secondaryValues != null) { // Output the joined data—use a unique separator (like |) to avoid conflicts with commas String outputValue = checkoutDateTime + "|" + itemType + "|" + secondaryValues; context.write(new Text(bibNum), new Text(outputValue)); } } } }
When submitting your job, you need to pass the path to the secondary CSV and ensure it's accessible to all Mapper tasks on the cluster:
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.output.FileOutputFormat; import java.io.IOException; public class CSVJoinDriver { public static void main(String[] args) throws IOException, InterruptedException, ClassNotFoundException { Configuration conf = new Configuration(); Job job = Job.getInstance(conf, "CSV_Merge_Job"); // Pass the secondary CSV path to the Mapper's configuration job.getConfiguration().set("secondary.csv.path", args[1]); // Assume args[1] is the second CSV path // For cluster mode: Add the file to Distributed Cache so all nodes can access it job.addCacheFile(new Path(args[1]).toUri()); job.setJarByClass(CSVJoinDriver.class); job.setMapperClass(CSVJoinMapper.class); job.setOutputKeyClass(Text.class); job.setOutputValueClass(Text.class); // Set main input (first CSV) and output paths FileInputFormat.addInputPath(job, new Path(args[0])); FileOutputFormat.setOutputPath(job, new Path(args[2])); System.exit(job.waitForCompletion(true) ? 0 : 1); } }
- CSV Parsing Issues: If your CSV has commas inside quoted values (e.g.,
"Doe, John"), splitting withsplit(",")will break fields. Use a dedicated CSV parser like OpenCSV instead:import com.opencsv.CSVReader; // ... CSVReader reader = new CSVReader(new FileReader(secondaryCSVPath)); String[] fields; while ((fields = reader.readNext()) != null) { // Process fields here } - Empty Columns: Always use
split(",", -1)to preserve trailing empty columns—without the-1, split will drop them, leading to incorrect field indices. - Cluster File Access: If running on a Hadoop cluster, the secondary CSV must be stored in HDFS, not just your local machine. Use
hdfs dfs -putto upload it first.
If you don't need to aggregate records with the same BibNum, you can skip the Reducer entirely by setting job.setNumReduceTasks(0). If you do need to combine multiple records for the same key, here's a simple Reducer:
import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; import java.io.IOException; public class CSVJoinReducer extends Reducer<Text, Text, Text, Text> { @Override protected void reduce(Text key, Iterable<Text> values, Context context) throws IOException, InterruptedException { // Output all joined records (modify if you need to aggregate) for (Text value : values) { context.write(key, value); } } }
Don't forget to set the Reducer class in your driver with job.setReducerClass(CSVJoinReducer.class) if you use it.
内容的提问来源于stack exchange,提问作者Gideok Seong

