You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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 IntWritable with DoubleWritable, 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.26 10:22:26