如何在Hadoop MapReduce中基于双Key实现按城市统计最高销量产品?
Got it, let's tackle this problem step by step. You already have a start with aggregating products per city, so we just need to extend that to find the top-selling product for each city. Here's a complete, working solution using Hadoop MapReduce, split into two jobs for clarity (this is more scalable for large datasets):
Job 1: Calculate Total Sales per Product per City
First, we'll count how many times each product appears in each city (if "sales" refers to transaction count; we'll cover amount-based sales later).
Mapper Class
This mapper reads each line, extracts the city and product, and outputs a composite key (city|product) paired with a count of 1.
import java.io.IOException; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; public class CityProductCountMapper extends Mapper<Object, Text, Text, IntWritable> { private final static IntWritable one = new IntWritable(1); private Text cityProductKey = new Text(); @Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { String[] fields = value.toString().split(";"); // Ensure we have at least 4 columns to avoid index errors if (fields.length >= 4) { String city = fields[2].trim(); String product = fields[3].trim(); // Use a pipe to separate city and product (avoid if your data has pipes!) cityProductKey.set(city + "|" + product); context.write(cityProductKey, one); } } }
Reducer Class
This reducer sums the counts for each city|product key, then outputs the city as the new key and a string like product:sales_count as the value.
import java.io.IOException; import org.apache.hadoop.io.IntWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; public class CityProductCountReducer extends Reducer<Text, IntWritable, Text, Text> { @Override protected void reduce(Text key, Iterable<IntWritable> values, Context context) throws IOException, InterruptedException { int totalSales = 0; for (IntWritable count : values) { totalSales += count.get(); } // Split the composite key back into city and product String[] keyParts = key.toString().split("\\|", 2); // Limit split to 2 parts to handle products with pipes if (keyParts.length == 2) { String city = keyParts[0]; String product = keyParts[1]; context.write(new Text(city), new Text(product + ":" + totalSales)); } } }
Job 2: Find Top-Selling Product per City
This job takes Job 1's output and finds the product with the highest sales count for each city.
Mapper Class
A simple mapper that passes through the city key and product:sales_count value directly.
import java.io.IOException; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Mapper; public class CityTopProductMapper extends Mapper<Object, Text, Text, Text> { @Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { // Split input by tab (Hadoop's default output separator) String[] parts = value.toString().split("\\t", 2); if (parts.length == 2) { String city = parts[0]; String productSales = parts[1]; context.write(new Text(city), new Text(productSales)); } } }
Reducer Class
This reducer iterates through all products for a city, tracks the one with the highest sales, and handles ties by listing all top products.
import java.io.IOException; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Reducer; public class CityTopProductReducer extends Reducer<Text, Text, Text, Text> { @Override protected void reduce(Text city, Iterable<Text> productSalesList, Context context) throws IOException, InterruptedException { String topProduct = ""; int maxSales = 0; for (Text productSales : productSalesList) { String[] parts = productSales.toString().split(":", 2); if (parts.length == 2) { String product = parts[0]; int sales = Integer.parseInt(parts[1]); if (sales > maxSales) { maxSales = sales; topProduct = product; } else if (sales == maxSales) { // Handle ties by appending products with the same max sales topProduct += ", " + product; } } } // Output the city and its top-selling product(s) with sales count context.write(city, new Text(topProduct + " (" + maxSales + " sales)")); } }
Driver Class to Run Both Jobs
This class submits both jobs, uses a temporary directory for Job 1's output, and cleans up the temp dir after Job 2 completes.
import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.IntWritable; 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; public class CityTopProductDriver { public static void main(String[] args) throws Exception { Configuration conf = new Configuration(); // Configure Job 1: Count sales per product per city Job job1 = Job.getInstance(conf, "City-Product-Sales-Count"); job1.setJarByClass(CityTopProductDriver.class); job1.setMapperClass(CityProductCountMapper.class); job1.setCombinerClass(CityProductCountReducer.class); // Optimize with combiner (safe for sum operations) job1.setReducerClass(CityProductCountReducer.class); job1.setMapOutputKeyClass(Text.class); job1.setMapOutputValueClass(IntWritable.class); job1.setOutputKeyClass(Text.class); job1.setOutputValueClass(Text.class); FileInputFormat.addInputPath(job1, new Path(args[0])); Path tempOutput = new Path(args[1] + "_temp"); FileOutputFormat.setOutputPath(job1, tempOutput); // Wait for Job 1 to finish before starting Job 2 if (!job1.waitForCompletion(true)) { System.exit(1); } // Configure Job 2: Find top product per city Job job2 = Job.getInstance(conf, "City-Top-Product"); job2.setJarByClass(CityTopProductDriver.class); job2.setMapperClass(CityTopProductMapper.class); job2.setReducerClass(CityTopProductReducer.class); job2.setOutputKeyClass(Text.class); job2.setOutputValueClass(Text.class); FileInputFormat.addInputPath(job2, tempOutput); FileOutputFormat.setOutputPath(job2, new Path(args[1])); // Run Job 2 and clean up temp directory on success boolean success = job2.waitForCompletion(true); if (success) { tempOutput.getFileSystem(conf).delete(tempOutput, true); } System.exit(success ? 0 : 1); } }
Notes for Customization
- If "sales" refers to monetary amount: Replace
IntWritablewithDoubleWritable, parse the 5th column (amount) instead of counting transactions, and sum doubles instead of integers. - Handle special characters: If your city/product names contain pipes or colons, use a different separator (like a rare character
\u0001) that won't appear in your data. - Single Job Alternative: For small datasets, you can skip Job 1 and track product counts directly in the reducer, but this is less scalable for large data.
内容的提问来源于stack exchange,提问作者Robert Almeida

