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

基于MapReduce的Python三方表关联问题:结果不完整排查

问题描述

现有三个独立CSV表:

  • customers(c_id, gender, address, dob)
  • meals(r_id, c_id, date)
  • restaurants(type, r_id)

需要通过Python MapReduce任务统计所有bistro类型餐厅的男性顾客用餐次数,对应SQL查询如下:

SELECT r.r_id, COUNT(*) AS count_meals 
FROM restaurants r 
INNER JOIN meals m ON r.r_id = m.r_id 
INNER JOIN customers c ON m.c_id = c.c_id 
WHERE c.gender = 'MALE' AND r.type = 'bistro' 
GROUP BY r.r_id

初始实现中,Mapper按字段长度区分表,但Reducer无法获取同key的匹配分组;改为多步骤关联后,输出结果仅为SQL查询的部分结果。相关代码片段如下:

初始Mapper代码

from mrjob.job import MRJob
from mrjob.step import MRStep
from mr3px.csvprotocol import CsvProtocol
import csv

class MRCustomers(MRJob):
    OUTPUT_PROTOCOL = CsvProtocol
    def mapper(self, _, line):
        if line.startswith('c_id') or line.startswith('r_id') or line.startswith('type'):
            return
        reader = csv.reader([line])
        columns = next(reader)
        if len(columns) == 4:
            if str(columns[1]) != 'MALE':
                return
            c_id = columns[0]
            yield c_id, "customer"
        elif len(columns) == 3:
            r_id = columns[0]
            c_id = columns[1]
            yield r_id, ("M", c_id)
        else:
            type = columns[0]
            if type == 'bistro':
                r_id = columns[1]
                yield r_id, "restaurant"

多步骤Reducer代码

def reducer1(self, key, values):
        joins = [x for x in values]
        if len(joins) > 1:
            if joins[0][0] == "restaurants":
                for tup in joins[1:]:
                    c_id = tup[1]
                    yield c_id, (tup[0], key)
        elif joins[0][0] == "customer":
            for customer in joins:
                yield key, ("customer", key)

    def mapper2(self, key, value):
        yield key, value

    def reducer2(self, key, values):
        joins = [x for x in values]
        if len(joins) > 1:
            if joins[0][0] == "customer":
                for tup in joins[1:]:
                    yield tup[1], 1

    def reducer3(self, key, values):
        yield None, (key, sum(values))

    def steps(self):
        first_step = MRStep(
            mapper=self.mapper1,
            reducer=self.reducer1,
        )
        second_step = MRStep(
            mapper=self.mapper2,
            reducer=self.reducer2,
        )
        third_step = MRStep(
            reducer=self.reducer3
        )

        return [first_step, second_step, third_step]
问题原因分析
  1. 数据标识不一致:Mapper输出的餐厅标识是"restaurant",但reducer1中判断的是"restaurants",导致匹配逻辑完全失效,大量有效数据被丢弃。
  2. 依赖不确定的数据顺序:reducer1直接取joins[0]判断数据类型,而MapReduce中values的顺序是随机的,无法保证第一个元素是餐厅数据,导致大量符合条件的分组被跳过。
  3. 无效数据未提前过滤:meals表数据未关联bistro餐厅就被输出,后续无法过滤这些无效记录,既导致统计结果错误,又浪费计算资源。
  4. Reducer逻辑不完整:reducer1仅处理餐厅数据在前的情况,未覆盖其他数据排列组合,导致部分关联数据无法传递到后续步骤;reducer3用None作为key,会把所有统计结果合并到同一组,不符合按r_id分组的逻辑,易引发数据混乱。
  5. 冗余步骤:mapper2仅原样输出键值对,无任何逻辑处理,属于冗余步骤。
正确实现方案

采用三步关联+统计的思路,每一步仅处理必要逻辑,提前过滤无效数据:

  1. 第一步:关联restaurants和meals,仅保留bistro类型餐厅的用餐记录,输出以c_id为key、r_id为value的键值对。
  2. 第二步:关联第一步结果和customers,仅保留男性顾客的用餐记录,输出以r_id为key、计数1为value的键值对。
  3. 第三步:按r_id分组统计用餐次数,输出最终结果。

完整代码如下:

from mrjob.job import MRJob
from mrjob.step import MRStep
from mr3px.csvprotocol import CsvProtocol
import csv

class MRBistroMaleMeals(MRJob):
    OUTPUT_PROTOCOL = CsvProtocol

    def mapper_join_init(self):
        # 缓存bistro类型的餐厅ID
        self.bistro_restaurants = set()

    def mapper_join(self, _, line):
        if line.startswith('c_id') or line.startswith('r_id') or line.startswith('type'):
            return
        
        reader = csv.reader([line])
        columns = next(reader)
        
        # 处理restaurants表,缓存bistro餐厅ID
        if len(columns) == 2:
            type_rest, r_id_rest = columns[0], columns[1]
            if type_rest == 'bistro':
                self.bistro_restaurants.add(r_id_rest)
        # 处理meals表,仅保留bistro餐厅的记录,输出(c_id, r_id)
        elif len(columns) == 3:
            r_id_meal, c_id_meal, _ = columns
            if r_id_meal in self.bistro_restaurants:
                yield c_id_meal, ('meal', r_id_meal)
        # 处理customers表,仅保留男性顾客,输出(c_id, 'male')
        elif len(columns) == 4:
            c_id_cust, gender_cust, _, _ = columns
            if gender_cust == 'MALE':
                yield c_id_cust, ('customer', 'male')

    def reducer_join(self, c_id, values):
        # 收集当前c_id对应的所有餐厅ID和顾客性别
        meal_r_ids = []
        is_male = False
        
        for val_type, data in values:
            if val_type == 'meal':
                meal_r_ids.append(data)
            elif val_type == 'customer':
                is_male = True
        
        # 仅当是男性顾客时,输出每个对应的餐厅ID和计数1
        if is_male and meal_r_ids:
            for r_id in meal_r_ids:
                yield r_id, 1

    def reducer_count(self, r_id, counts):
        # 统计每个餐厅的用餐次数
        yield r_id, sum(counts)

    def steps(self):
        return [
            MRStep(
                mapper_init=self.mapper_join_init,
                mapper=self.mapper_join,
                reducer=self.reducer_join
            ),
            MRStep(
                reducer=self.reducer_count
            )
        ]

if __name__ == '__main__':
    MRBistroMaleMeals.run()

代码说明

  • mapper_join_init:初始化缓存集合,存储所有bistro类型餐厅ID,避免重复判断。
  • mapper_join:分表处理数据,提前过滤非bistro餐厅和非男性顾客的数据,只输出需要关联的键值对。
  • reducer_join:按c_id分组,判断是否为男性顾客,若是则输出对应的餐厅ID和计数1。
  • reducer_count:按r_id分组,统计每个餐厅的用餐总次数,输出最终结果。

内容的提问来源于stack exchange,提问作者Dan G

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.22 19:17:01