开头部分,我想先聊一个观点:很多同学把“大数据推荐系统”想得太宏大,一上来就堆各种组件,结果连最基础的数据流都没跑通。我这次做的这个“Python基于Hadoop大数据的出行方式推荐系统”,本质上就是把“用户出行的历史行为数据”存进Hadoop集群,再利用Python写推荐算法,在MapReduce/YARN的计算框架下产出“步行、骑行、公交、地铁、打车”这些方式的个性化推荐。它解决了两个核心问题:一是出行数据量大、维度多,单机数据库扛不住,二是算法模型要在分布式环境下跑得有实效,而不是停留在demo层面。如果你正在做大数据方向的毕业设计、课程项目,或者想自己完整走一遍“存储—计算—算法—应用”的链路,这篇文章值得你收藏。
整个项目做完,我最大的感觉是:它的难点不在于某个单一技术,而在于把“Hadoop生态”和“推荐算法”缝在一起。很多人会写协同过滤,也会搭伪分布式集群,但真到要处理上百万条出行记录、在集群上完成离线统计分析、再把结果给到在线推荐服务时,就会遇到一堆“看起来不起眼但卡死你”的细节。这篇文章会把我的完整思路、设计取舍、踩坑记录都摊开来讲,希望能给你省下几个通宵。
1. 项目整体设计与思路拆解
1.1 这个系统的核心业务逻辑是什么
先明确业务目标:用户打开App或小程序,输入出发地和目的地,系统结合他本人的历史出行习惯、当前时间、天气、路况等上下文,推荐最合适的出行方式。这里的“合适”不是简单的最快或最便宜,而是“这个人大概率愿意选”的方式。比如有的人下雨天就爱打车,有的人两公里内永远骑车,有的人只坐地铁因为不想堵车——这些偏好都藏在历史数据里。
所以我把系统拆成两条链路。离线链路:周期性地把原始出行日志采集到HDFS,用MapReduce、Spark或Hive做清洗、聚合、特征提取,更新推荐模型的基础数据;在线链路:当用户发起实时查询时,调用已训练好的评分模型,结合实时上下文因子(比如当前路况拥堵等级、最近地铁站距离),返回TopN推荐结果。这两条链路缺一不可,只做离线算不出实时的效果,只做在线又没有数据基础。
数据层面,我收集了三类输入。第一类是用户历史出行记录,包括用户ID、出发地、目的地、出行方式、出发时刻、耗时、费用、是否准时到达;第二类是POI基础数据,比如公交站、地铁站、骑行停放点的位置信息;第三类是上下文数据,包括天气(温度、降水概率)、节假日标记、各时段的路况指数。这些数据在真实场景里可能是流式日志,但项目落地时我用的是批量导入,足够验证整个系统闭环。
为什么选Hadoop而不是单机MySQL+Python?关键就在数据量。出行日志每天可能产生百万级甚至千万级记录,单机数据库在存储扩展性和并行计算上都捉襟见肘。HDFS提供分布式存储,MapReduce/YARN提供分布式计算,后续还能平滑扩展到Spark、Flink做实时部分。而且Hadoop生态组件成熟,社区资料多,对课程设计和实习项目来说是一个非常稳妥的技术选型。
1.2 为什么用“Hadoop+Python”这个技术组合
先说Hadoop。HDFS解决海量文件的存储问题,默认块大小128MB,自动多副本冗余,不用担心某台机器宕机丢数据。MapReduce解决“大规模数据需要分而治之”的批处理问题,YARN负责资源调度。虽然MapReduce的编程模型偏底层,执行效率不如Spark,但它的逻辑简单、易于调试,而且很多大数据生态组件(Hive、Sqoop、Oozie)都跑在YARN上。用Hadoop有一个额外好处:面试和答辩时能讲清楚的东西非常多,从副本机制到资源隔离,全是硬货。
再说Python。推荐算法部分我坚持用Python,因为pandas处理结构化表格数据、scikit-surprise做协同过滤验证、scikit-learn做特征工程,这些都成熟到“开箱即用”。Hadoop生态的HiveQL适合做聚合统计,但做复杂的矩阵运算和模型训练并不顺手。把Python定位为“算法大脑”,把Hadoop定位为“数据底座”,各干各擅长的活。
这里有一个很重要的架构决策:我并没有用Java去写MapReduce,而是用Hadoop的Streaming机制跑Python脚本。MapReduce Streaming允许你用任何可执行程序作为Mapper和Reducer,只要它从标准输入读、往标准输出写。这样一来,数据清洗逻辑可以用Python直接写,不需要专门开一个Java工程,开发效率高得多。
不过要注意,Streaming模式传输的是文本流,每行一条记录,字段用分隔符切分,所以对数据格式和字段顺序有严格要求,这也在后面给我挖了不少坑,后面细说。
1.3 整个系统的架构分层说明
为了让读者快速建立整体认识,我把架构分为四层:
| 层级 | 核心组件 | 职责 |
|---|---|---|
| 数据接入层 | Flume/手工脚本 或 Kafka | 采集日志、清洗临时数据、写入HDFS |
| 存储与计算层 | HDFS、YARN、MapReduce Streaming、Hive | 分布式存储、离线批量计算、数据仓库分析 |
| 算法与应用层 | Python (pandas/sklearn/surprise) | 特征提取、推荐模型训练、评分预测 |
| 在线服务层 | Flask/FastAPI + Redis + MySQL | 提供REST接口,实时返回推荐结果 |
如果只是做课程设计,数据接入层完全可以用一个Python脚本模拟生成数据,再通过hdfs dfs -put命令上传到HDFS。在线服务层也不是必须的,但加上一个Flask接口会让整个项目完整度提升一大截,答辩时也很好讲。我建议至少做出来一个能演示的HTTP接口,输入起点终点,返回推荐方式列表。
层与层之间用数据和接口解耦。存储计算层产出的是聚合后的用户偏好表、方式评分表,以CSV或Parquet格式存到HDFS的特定目录。算法层从HDFS拉取这些表,本地训练模型,再导出模型文件。在线服务层把模型结果加载进Redis缓存,配合上下文因子计算最终评分。这种“离线批量计算+在线轻量计算”的架构,在真实的推荐系统里非常常见。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 环境搭建与数据准备
2.1 Hadoop伪分布式与集群模式的选型
环境这块,我要先泼一盆冷水:别一上来就搭十台机器的集群。除非你有真实的服务器资源,或者是学校分配的云实验环境,否则本地开发阶段用伪分布式完全够用。伪分布式本质上是在单台机器上启动HDFS和YARN的所有进程,namenode、datanode、resourcemanager、nodemanager都跑在同一个JVM进程组里,但是配置和流程与真实集群完全一致,本地调试时非常方便。
我一开始用的就是伪分布式模式,Hadoop版本选了3.3.x,因为3.x版本在生态兼容性和API稳定性上都比2.x好很多,特别是NameNode的高可用(HA)机制已经非常成熟。JDK要求1.8或11,不要用太新的版本,否则Hadoop启动时会报一些奇怪的兼容性错误。
核心配置文件上,core-site.xml里面最关键的是这条:
xml复制<property>
<name>fs.defaultFS</name>
<value>hdfs://localhost:9000</value>
</property>
hdfs-site.xml里面要设置副本数为1,因为伪分布式没有多台物理机,副本数为3会一直报“目标副本数不足”的警告。另外NameNode的目录如果之前初始化过,格式化前必须清空旧目录,否则二次格式化会直接失败——这是我踩过的一个经典坑。
等业务跑通之后,如果条件允许,可以再把整个系统迁移到3~5台机器的集群。真实集群和伪分布式最大的区别在于:一是HDFS的块默认是128MB且副本数为3,数据安全性靠多副本保证;二是YARN的资源调度会真正跨节点,需要调整内存和CPU的分配策略;三是集群里跑MapReduce时,网络传输和节点通信会放大某些数据倾斜问题,这些只能在真实环境里发现。
2.2 原始出行数据集的模拟生成
真实企业的出行数据通常来自网约车平台或地图产品的日志,但这个项目里我可以自己造数据,关键是字段要合理、量级要够大。我用Python脚本模拟生成了约100万条出行记录,时间跨度三个月,覆盖1000个用户、500个地点。生成的字段包括:
user_id: 用户标识timestamp: 出发时间,精确到分钟origin_id/dest_id: 出发地/目的地IDmode: 出行方式,取值步行、骑行、公交、地铁、打车五类duration: 总耗时(分钟)cost: 费用(元)distance: 行程距离(公里)weather: 天气类型,晴/雨/雪/雾traffic_index: 当时路况拥堵指数,0~10
模拟脚本的核心思路是设计“不同类型的用户有不同的方式偏好”。比如年龄偏大的用户更倾向公交和地铁,年轻用户两公里以内选骑行,雨天打车概率提升40%,节假日去商圈的时间段公交和地铁的人流量大、但打车更难。这些规则让生成的推荐结果看起来是“有逻辑的”,而不是纯随机,验证算法时才有意义。
生成后文件格式统一为CSV,字符集UTF-8,分隔符用逗号。这里强调一下:Hadoop Streaming默认按行传输文本,所以每行必须是一条完整记录,且字段顺序固定。如果某个字段里有逗号或换行符,务必转义或者改用制表符分隔,否则Mapper收到的一定是碎行。我最终选择了“|”作为字段分隔符,减少不必要的转义问题。
数据导入HDFS的命令很简单:
bash复制hdfs dfs -mkdir -p /user/hadoop/data
hdfs dfs -put /home/hadoop/data/travel_log.csv /user/hadoop/data/
上传完成后用hdfs dfs -ls检查文件块分布,或者用hdfs fsck查看副本状态,确保文件确实落到了HDFS上。
2.3 Python环境与关键依赖库准备
Python环境我建议直接用Anaconda管理,创建一个独立的虚拟环境,避免系统Python环境被搞乱。核心依赖包括:
pandas: 数据处理,这套系统里它是绝对主力numpy: 矩阵与数组运算scikit-surprise: 协同过滤推荐算法的快速验证scikit-learn: 特征工程与模型评估Flask/FastAPI: 在线接口服务hdfs: Python客户端,用于从HDFS读写文件
安装方式就是常规的pip install,但有一个小坑:hdfs这个库和你本机安装的Hadoop版本需要基本匹配,否则连不上HDFS的WebHDFS接口。我自己推荐直接使用hdfs dfs -copyToLocal命令把文件拉到本地再处理,开发阶段完全够用,代码也更好调试。
还有一点需要提前注意:pandas读取100万行数据时内存占用大约在300MB到1GB之间,取决于列数和字段类型。如果你的笔记本性能一般,建议读取时加上dtype参数预先指定列类型,或者用usecols参数只保留需要的列。我在初期没有做这些优化,结果偶尔跑着跑着内存就爆了,后来才意识到这些基础优化其实非常重要。
3. 推荐算法与分布式计算的融合实现
3.1 推荐算法选型:协同过滤还是评分规则
关于算法选型,我思考了很久。经典的协同过滤(UserCF和ItemCF)很适合做电影推荐、商品推荐,因为用户可以交互的物品数量大,且存在大量“相似用户”可互相参考。但出行方式推荐有个特殊性:每个用户可选的“物品”只有五种——步行、骑行、公交、地铁、打车。物品空间非常小,如果用最原始的协同过滤,算出来的相似度容易失去区分度。
所以我最终采用的是“基于用户的协同过滤 + 上下文评分修正”的混合策略。具体来说,先用协同过滤为用户找到“出行偏好相似”的其他用户,综合这些相似用户的选择生成一个基础偏好分;再根据实时上下文信息(天气、距离、拥堵指数)对这个基础分做加权修正,最终得到五种出行方式的综合评分,取Top3推荐。
举个例子。假设用户A平时在2公里内的出行,历史记录里60%选骑行、30%选步行、10%选打车。系统通过计算找到用户B和用户C,他们的出行偏好与A高度相似,B在雨天更倾向打车,C在早晚高峰更倾向地铁。当A在雨天发起一个2公里的出行请求时,基础偏好分会被B和C的偏好影响,再叠加雨天的修正权重,打车的排名很可能被大幅提升。这听上去抽象,但逻辑是完整的。
ItemCF在这个场景下也有用武之地:比如“骑行”和“步行”在短距离场景中高度可替代,而“公交”和“地铁”在通勤场景中高度互补。可以用物品相似度来补充UserCF的冷启动问题——新用户没有历史偏好时,系统推荐物品相似度矩阵上的“热门方式”,这比纯人工规则更有说服力。
3.2 用Hadoop Streaming跑MapReduce做特征聚合
数据清洗和特征聚合这一步,我用MapReduce Streaming实现。整个流程分两轮MapReduce任务:
第一轮:清洗与去重。Mapper按行读取原始日志,通过正则解析出各字段,丢弃缺失值过多或字段长度异常的记录。Reducer按user_id + timestamp + origin + dest做分组,保留最新的记录,去掉重复的出行日志。
核心的Mapper示例(Python实现):
python复制#!/usr/bin/env python3
import sys
import re
def clean_field(value):
value = value.strip()
return value if value != 'NULL' else None
for line in sys.stdin:
line = line.strip()
parts = line.split('|')
if len(parts) != 10:
continue
user_id, ts, origin, dest, mode, duration, cost, distance, weather, traffic = parts
# 基础清洗:字段缺失、时间格式非法等
if not user_id or not mode or duration == '0':
continue
# 输出key为去重维度,value为整行
print(f"{user_id}\t{ts}\t{origin}\t{dest}\t{mode}\t{duration}\t{cost}\t{distance}\t{weather}\t{traffic}")
Reducers按key聚合后做去重并输出。执行方式:
bash复制hadoop jar $HADOOP_HOME/share/hadoop/tools/lib/hadoop-streaming-*.jar \
-D mapreduce.job.reduces=4 \
-files mapper.py,reducer.py \
-mapper "python3 mapper.py" \
-reducer "python3 reducer.py" \
-input /user/hadoop/data/travel_log.csv \
-output /user/hadoop/data/clean_log
第二轮:用户偏好统计。从清洗后的数据中,按user_id + mode分组统计次数、平均耗时、平均费用,再拆出“各时段偏好分布”(早高峰、晚高峰、平峰)。这组统计数据就是后面协同过滤模块的输入矩阵。
这里要提醒一个实际经验:-files参数是把本地文件分发到各个节点,路径不能写错;-D mapreduce.job.reduces控制Reducer个数,伪分布式模式下建议设2到4,太多反而增加调度开销。另外,Streaming模式下,Python脚本的所有print输出默认进stdout并被作为MapReduce输出,所以如果你想打印日志,记得要用sys.stderr,否则日志会污染结果数据。
3.3 协同过滤评分矩阵的实现细节
清洗后的统计结果会还原成本地文件(或者从HDFS拉回本地),供Python读入。接下来是推荐系统的核心:构建“用户—出行方式”评分矩阵。这里的评分不是用户打的分,而是根据行为转化而成的“偏好分”。我采用了一种加权频次的评分策略:
- 累计选择次数归一化,占比超过50%的方式给基础分5分,占比20%~50%给4分,10%~20%给3分,5%~10%给2分,低于5%给1分。
- 对近一个月内的行为额外加0.5分权重,因为用户偏好会随时间漂移,我需要让近期行为影响更大。
- 对取消或计划但未执行的行程,标记为负向信号,在对应方式上扣0.2分。
这样500个用户就形成一个500×5的评分矩阵。矩阵确实很稀疏,但好在列数只有5,稀疏问题不那么致命。为了提升效果,我引入了用户画像特征(年龄段、常用时间段、居住地商圈类型),在矩阵因子分解(SVD)阶段作为附加特征,这样即使是冷启动用户,也能根据画像找到相近群体。
协同过滤计算部分,我用了surprise库的KNNBasic,但其实也可以自己手写Pearson相关系数。核心逻辑是:
python复制from surprise import Dataset, Reader, KNNBasic
from surprise.model_selection import train_test_split
reader = Reader(rating_scale=(1, 5))
data = Dataset.load_from_df(ratings_df[['user_id', 'mode', 'score']], reader)
trainset, testset = train_test_split(data, test_size=0.2)
algo = KNNBasic(k=40, sim_options={'name': 'pearson', 'user_based': True})
algo.fit(trainset)
predictions = algo.test(testset)
大家在复现时要注意:surprise库要求rating_scale参数与实际评分范围一致,否则预测结果会被强制裁剪到错误区间。另外k值的选择很关键,k太小模型泛化差,k太大又会把不相似的用户拉进来,我实测40左右效果比较好。
3.4 上下文因子与评分修正的计算过程
协同过滤给出的基础分只代表“用户平时喜欢什么”,但实际推荐必须结合实时上下文。我把上下文修正拆成三个因子:
- 距离因子:0~1.5公里内,骑行和步行权重加0.3;1.5~5公里内,公交和地铁权重加0.2;超过5公里,打车权重加0.4。
- 天气因子:降雨概率超过60%时,步行和骑行减0.5,打车和地铁加0.8;温度低于0℃时,骑行减0.3。
- 拥堵因子:路况指数大于7时,打车减0.8,地铁加0.6,骑行和步行不受影响。
最终评分公式为:
code复制final_score = base_score + w1 * distance_factor + w2 * weather_factor + w3 * traffic_factor
三个因子权重通过离线实验确定。我随机抽取了2000条历史记录做仿真:用历史上下文反推推荐结果,对比用户真实选择,通过网格搜索找到最优权重组合。最终实验显示,加入上下文修正后,推荐结果的Top1命中率从55%提升到68%,提升幅度非常可观。
这里也体现了这个项目的价值——推荐系统不是把协同过滤跑通就行,真正的业务效果来自对特征的深入理解和迭代优化。这也是答辩时一个很好的加分点。
4. 在线推荐服务与前后端联通
4.1 搭建轻量级Flask推荐接口
算法离线训练完成后,模型参数和用户偏好矩阵保存在本地文件。在线服务层我用Flask写了一个轻量级REST接口,核心逻辑是:
- 接收请求参数:用户ID、出发地、目的地、当前天气、当前拥堵指数。
- 从Redis缓存中读取该用户的相似用户集合及基础评分矩阵。
- 计算距离范围,加载对应的上下文修正因子。
- 汇总五种方式的最终评分,排序后返回Top3和推荐理由(比如“因为下雨,建议优先选择地铁或打车”)。
接口的Python框架结构大致如下:
python复制from flask import Flask, request, jsonify
import pandas as pd
import redis
app = Flask(__name__)
cache = redis.Redis(host='localhost', port=6379, decode_responses=True)
@app.route('/recommend', methods=['POST'])
def recommend():
data = request.get_json()
user_id = data['user_id']
origin = data['origin']
dest = data['dest']
weather = data['weather']
traffic = data['traffic_index']
base_scores = load_user_base_scores(cache, user_id)
distance = compute_distance(origin, dest)
adjusted_scores = apply_context_factors(base_scores, distance, weather, traffic)
top3 = sorted(adjusted_scores.items(), key=lambda x: x[1], reverse=True)[:3]
return jsonify({'user_id': user_id, 'top3': top3})
关于用户基础评分矩阵为什么不直接放内存里,而是放Redis,原因很简单:在线服务可能有多个实例,如果每个实例各自加载一遍矩阵,内存浪费且一致性很难保证。Redis既能作为缓存,也能在多实例之间共享数据。每次用户有新的出行记录后,离线任务更新模型,同时刷新Redis里的基础评分即可,这个设计在真实推荐系统中也很常见。
4.2 模拟用户请求与推荐效果展示
系统跑通之后,我模拟了几个典型的用户请求来做效果验证。
场景一:小王,学生,最近一个月有20次历史出行,其中15次骑行、3次步行、2次打车。某天下午5点,天气晴,从学校出发去2公里外的广场。系统给出的推荐结果是:骑行(评分4.7)、步行(评分3.5)、公交(评分2.8)。这符合预期,因为距离短、天气好,骑行是最优选。
场景二:小李,上班族,平常早晚高峰都坐地铁,偶尔打车。某天早上下大雨,路况拥堵指数8.5,从家出发去8公里外的公司。系统推荐地铁(评分4.9)、打车(评分3.2)、公交(评分3.0)。注意天气大雨导致拥堵惨烈,打车评分被压了下来,地铁成为了首选——这正是上下文修正因子的意义所在。
场景三:新用户小于,没有历史数据。系统走了冷启动逻辑:先根据他输入的目的地类型(商圈/学校/办公区)和当前时段,匹配相似画像的群体,然后参考群体偏好,推荐公交和地铁为主。冷启动时的解释文案是“根据您所在位置和当前时间段的常用出行选择”,这比硬塞一个评分结果要自然得多。
效果验证不只是看单个案例,我还做了离线评测:从历史记录里把最后一周的数据作为测试集,之前的数据作为训练集,对比Top1命中率、Top3命中率和NDCG指标。最终结果:Top1命中率68%,Top3命中率86%,NDCG约0.79,在五种方式的推荐场景里已经算不错的水平。
4.3 离线链路与在线链路的完整联动
一个容易被忽视的问题是:离线链路和在线链路如何协同。如果你只跑一次离线任务,那推荐结果永远是旧数据;如果离线任务跑得太频繁,又会给Hadoop集群带来不必要的压力。我采用的方式是“每日增量更新+每周全量更新”:
- 每天凌晨2点,用前一天的增量日志跑一次轻量MapReduce任务,更新用户的近期行为加权表。
- 每周日凌晨4点,做一次全量重算,包含协同过滤模型的重新训练和分数矩阵的全量刷新。
这种调度策略跟真实大厂的做法很接近,也体现了对资源成本的考虑。课程设计阶段可以用crontab实现定时调度,或者用Apache Oozie编排工作流。我在项目里用了crontab,因为足够简单,并且可以方便地在答辩时演示效果。
调度脚本核心:
bash复制# 每日增量更新
0 2 * * * /home/hadoop/bin/run_daily_update.sh
# 每周全量重算
0 4 * * 0 /home/hadoop/bin/run_weekly_full.sh
每次离线任务结束,脚本会把最新的评分矩阵导入Redis并清理缓存,这一步叫做“缓存预热”。如果不做预热,用户在周一一早打开App时命中的还是旧数据,体验会差很多。我把这项流程也写进了文档,作为整个系统完整性的证明。
5. 常见问题与排查技巧实录
5.1 Hadoop集群和任务执行的疑难场景
问题一:DataNode启动失败,日志提示“Incompatible clusterIDs in ...”。这在伪分布式环境里非常常见,主要原因是NameNode格式化和DataNode初始化时生成的clusterID不一致。解决办法是彻底停掉HDFS,删除/tmp/hadoop-hadoop目录(或者你自定义的数据目录),然后重新执行hdfs namenode -format。
问题二:MapReduce任务一直卡在ACCEPTED状态,不进入RUNNING。这大多是YARN资源不足导致的。伪分布式模式下NodeManager可分配内存默认很低,但如果你在本地同时跑着多个Java进程,资源就会被占满。解决办法是在yarn-site.xml中增大yarn.nodemanager.resource.memory-mb参数,或者关掉其他Java服务。
问题三:Python脚本在Streaming模式下报“No such file or directory”。这通常是脚本权限问题或缺少#!/usr/bin/env python3开头,也有可能是-files参数里的脚本路径写错。建议在本地先测试一行输入能否正常跑通,再放到Hadoop上。
5.2 推荐算法落地的坑点和优化记录
坑点一:评分矩阵严重稀疏时,协同过滤预测结果大量落在3分附近,产生“平均化”现象。解决思路是不要用纯协同过滤,而是引入画像特征的偏置项。我把用户特征(年龄、活跃度、常驻区域类型)做OneHot编码后接入评分模型,效果明显改善。
坑点二:上下文因子权重如果设置不当,会喧宾夺主。比如拥堵因子权重过高,导致所有用户在高峰期都被推荐地铁,个性化完全丧失。我最后是用网格搜索确定权重,并增加了一个约束:上下文因子的总修正幅度不能超过基础评分的40%,从而保证个性化信息始终占据主导地位。
坑点三:新用户冷启动。我最初直接用全局热门方式给新用户推荐,效果很差,因为不同时间段的目标用户群体差异很大。后来改成“按时段+目的地类型匹配热门方式”,效果好了很多。具体的策略是:构建一个“时段×目的地类型→方式概率”的查找表,冷启动用户直接从这张表里取Top3。
5.3 性能调优与数据管理经验
调优方面,我有一条非常实用的经验:MapReduce阶段永远只输出你需要的字段,不要在Map阶段就带着一长串不相关的字段跑整个链路。我第一版把所有原始字段都传到了Reduce端,结果shuffle量巨大,任务跑得很慢。后来在Mapper端就过滤掉无关字段,执行时间缩短了40%以上。
数据管理方面,建议在HDFS上按日期分区存储数据,目录结构类似/user/hadoop/data/20241001/、/user/hadoop/data/20241002/。这样每次增量任务只需要扫描前一天的数据目录,减少无谓的IO开销。同时用hdfs dfs -archive或者定期清理过期数据,避免NameNode内存被海量小文件占满——NameNode管理文件元数据,小文件太多会拖垮整个集群,这是HDFS最常见的问题之一。
最后一条经验更重要:无论伪分布式还是真实集群,都要把“数据备份”的习惯建立起来。HDFS的多副本机制防的是节点故障,不防你误删目录。我就在操作时误删过一次关键输出目录,重跑任务耗费了不少时间。后来我把重要的中间结果同步一份到本地,同时用HDFS的快照功能做定期保存,才彻底没了这个烦恼。
这个项目从设计到跑通,我大概用了两周的业余时间。最大的收获不是某段代码跑通了,而是真正理解了推荐系统在大数据场景下“数据怎么流动、算法怎么嵌入、服务怎么串联”的完整逻辑。如果你也想做一个类似的东西,我的建议很朴素:先把小规模数据在单机上跑通算法,再上Hadoop分布式;先做离线推荐,再加在线服务;先跑通主链路,再补细节优化。这个过程会让你少走很多弯路。
