基于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]
问题原因分析
- 数据标识不一致:Mapper输出的餐厅标识是
"restaurant",但reducer1中判断的是"restaurants",导致匹配逻辑完全失效,大量有效数据被丢弃。 - 依赖不确定的数据顺序:reducer1直接取
joins[0]判断数据类型,而MapReduce中values的顺序是随机的,无法保证第一个元素是餐厅数据,导致大量符合条件的分组被跳过。 - 无效数据未提前过滤:meals表数据未关联bistro餐厅就被输出,后续无法过滤这些无效记录,既导致统计结果错误,又浪费计算资源。
- Reducer逻辑不完整:reducer1仅处理餐厅数据在前的情况,未覆盖其他数据排列组合,导致部分关联数据无法传递到后续步骤;reducer3用
None作为key,会把所有统计结果合并到同一组,不符合按r_id分组的逻辑,易引发数据混乱。 - 冗余步骤:mapper2仅原样输出键值对,无任何逻辑处理,属于冗余步骤。
正确实现方案
采用三步关联+统计的思路,每一步仅处理必要逻辑,提前过滤无效数据:
- 第一步:关联
restaurants和meals,仅保留bistro类型餐厅的用餐记录,输出以c_id为key、r_id为value的键值对。 - 第二步:关联第一步结果和
customers,仅保留男性顾客的用餐记录,输出以r_id为key、计数1为value的键值对。 - 第三步:按
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
相关产品推荐
相关产品推荐

