欢迎光临
我们一直在努力

基于Spark的旅游路线推荐系统代码讲解文档

本文档旨在帮助学生理解系统代码结构和实现原理,方便答辩时进行代码讲解。



系统概述

在这里插入图片描述

系统简介

本系统是一个基于Spark的旅游路线推荐系统,采用Django作为后端框架,结合Hadoop/Spark大数据技术进行数据分析,实现了旅游路线的智能推荐功能。系统主要面向旅游爱好者,提供旅游路线浏览、景点信息查询、路线推荐、评论互动等功能。

系统功能

功能模块功能描述
用户管理 用户注册、登录、个人信息管理
旅游路线 路线浏览、添加、编辑、删除、分类管理
景点信息 景点浏览、收藏、评分查看
智能推荐 协同过滤推荐、点击热度推荐
评论功能 路线评论、点赞踩功能
旅游资讯 新闻发布、浏览、点赞收藏
数据分析 Hadoop MapReduce统计分析
收藏功能 路线、景点、资讯收藏

项目架构

目录结构

djangoy1pr16u7/
├── dj2/ # Django项目配置目录
│ ├── settings.py # 项目配置文件
│ ├── urls.py # 主路由配置
│ ├── views.py # 主视图文件
│ └── wsgi.py # WSGI部署配置
├── main/ # 主应用目录
│ ├── models.py # 数据模型定义
│ ├── model.py # 基础模型类
│ ├── urls.py # 应用路由配置
│ ├── views.py # 应用视图
│ ├── hadoop_v.py # Hadoop分析视图
│ ├── Lvyouluxian_v.py # 旅游路线视图控制器
│ ├── Yonghu_v.py # 用户视图控制器
│ ├── Sightinfo_v.py # 景点信息视图控制器
│ ├── group_mapper.py # MapReduce分组Mapper
│ ├── group_reducer.py # MapReduce分组Reducer
│ ├── value_mapper.py # MapReduce数值Mapper
│ ├── value_reducer.py # MapReduce数值Reducer
├── util/ # 工具类目录
│ ├── auth.py # 用户认证工具
│ ├── common.py # 公共工具类
│ ├── codes.py # 状态码定义
│ ├── spark_func.py # Spark分析函数
│ ├── hive_func.py # Hive数据库函数
│ ├── mapreduce_func.py # MapReduce函数
│ ├── hdfs_func.py # HDFS文件操作
│ ├── configread.py # 配置文件读取
├── xmiddleware/ # 中间件目录
│ ├── xauth.py # 认证中间件
│ ├── xparam.py # 参数处理中间件
├── templates/ # 前端模板目录
│ ├── front/ # 前端页面
│ └── upload/ # 上传文件目录
├── media/ # 媒体文件目录
├── db/ # 数据库SQL文件
├── config.ini # 数据库配置文件
├── manage.py # Django管理脚本
└── requirements.txt # 依赖包列表

架构分层

系统采用MVT(Model-View-Template)架构模式:

  • Model层:数据模型定义,负责与数据库交互
  • View层:业务逻辑处理,接收请求返回响应
  • Template层:前端页面展示,用户交互界面
  • Middleware层:中间件,处理认证和参数预处理
  • Util层:工具类,提供公共功能支持

  • 技术栈介绍

    后端技术

    技术版本说明
    Python 3.x 后端开发语言
    Django 2.0 Web框架
    MySQL 5.5+ 关系型数据库
    Hadoop 3.3.0 分布式计算框架
    Spark 大数据分析框架
    Hive 数据仓库工具
    HDFS 分布式文件系统

    前端技术

    技术说明
    HTML5 页面结构
    CSS3 样式设计
    JavaScript 页面交互
    Vue.js 前端框架(部分页面)
    jQuery DOM操作库
    Layui UI组件库

    关键依赖包

    Django==2.0
    pymysql # MySQL数据库连接
    pycrypto # 加密工具
    mrjob # MapReduce作业
    hdfs # HDFS客户端
    impyla # Hive连接
    pyspark # Spark分析
    pandas # 数据处理
    requests # HTTP请求
    xlrd # Excel读取


    核心模块详解

    1. 配置模块 (dj2/settings.py)

    代码讲解

    # 第110行 – 读取数据库配置
    dbtype, host, port, user, passwd, dbName, charset, hasHadoop = config_read("config.ini")

    讲解要点:

    • config_read函数从config.ini文件读取数据库配置参数
    • 支持MySQL和MSSQL两种数据库类型
    • hasHadoop参数标识是否启用Hadoop大数据功能

    # 第114-136行 – MySQL数据库配置
    if dbtype == 'mysql':
    DATABASES = {
    'default': {
    'ENGINE': 'django.db.backends.mysql',
    'NAME': dbName,
    'USER': user,
    'PASSWORD': passwd,
    'HOST': host,
    'PORT': port,
    'charset': charset,
    },
    }

    讲解要点:

    • Django使用ORM方式连接MySQL
    • ENGINE指定数据库引擎为MySQL
    • 配置项与config.ini中的参数对应

    # 第44-57行 – 中间件配置
    MIDDLEWARE = [
    'django.middleware.security.SecurityMiddleware',
    'django.contrib.sessions.middleware.SessionMiddleware',
    'django.middleware.common.CommonMiddleware',
    'django.contrib.auth.middleware.AuthenticationMiddleware',
    'django.middleware.clickjacking.XFrameOptionsMiddleware',
    'threadlocals.middleware.ThreadLocalMiddleware',
    "xmiddleware.xparam.Xparam", # 参数预处理中间件
    "xmiddleware.xauth.Xauth", # 认证中间件
    'corsheaders.middleware.CorsMiddleware', # 跨域中间件
    ]

    讲解要点:

    • 中间件按顺序执行,请求依次通过每个中间件
    • xparam中间件处理请求参数
    • xauth中间件验证用户身份
    • corsheaders中间件处理跨域请求

    2. 数据模型模块 (main/models.py)

    基础模型类 BaseModel (main/model.py)

    # 第15-17行 – 基础模型类定义
    class BaseModel(models.Model):
    class Meta:
    abstract = True # 抽象基类,不会创建实际表

    讲解要点:

    • abstract = True表示这是一个抽象基类
    • 所有数据模型继承此基类,获得公共方法
    • 不会在数据库中创建对应的表

    # 第197-287行 – createbyreq方法(创建记录)
    def __CreateByReq(self, model, params):
    # 自动生成ID(毫秒级时间戳)
    if model.__tablename__ != 'users':
    params['id'] = int(float(time.time()) * 1000)

    # 处理不同类型字段
    for col in model._meta.fields:
    if str(col.get_internal_type()).lower() == "integerfield":
    column_list.append(col.name)

    # 设置创建时间
    paramss["addtime"] = datetime.datetime.now()

    # 创建并保存记录
    m = model(**paramss)
    ret = m.save()
    return m.id

    讲解要点:

    • 使用毫秒级时间戳生成唯一ID
    • 根据字段类型自动处理数据转换
    • 自动设置addtime创建时间字段

    # 第35-174行 – page方法(分页查询)
    def __Page(self, model, params, request, q):
    # 获取分页参数
    page = params.get('page') if params.get('page') != None else 1
    limit = params.get('limit') if params.get('limit') != None else 99999

    # 处理排序
    sort = params.get('sort') if params.get('sort') else 'id'
    order = params.get('order')

    # 模糊查询处理
    for k, v in params.items():
    if "%" in str(v):
    fuzzy_val = v.replace("%", "")
    contain_str += '.filter({}__icontains="{}")'.format(k, fuzzy_val)

    # 使用Django Paginator分页
    datas = model.objects.filter(**condition).filter(q).all()
    p = Paginator(datas, int(limit))
    p2 = p.page(int(page))

    return newDataa, page, pages, total, limit

    讲解要点:

    • 支持分页查询,默认每页显示数据量可配置
    • sort和order参数控制排序
    • %符号表示模糊查询,使用icontains实现
    用户模型 yonghu

    # models.py 第10-45行
    class yonghu(BaseModel):
    __tablename__ = 'yonghu'

    # 用户表特殊属性
    __loginUser__ = 'yonghuzhanghao' # 登录用户字段
    __authPeople__ = '是' # 标识为用户表
    __loginUserColumn__ = 'yonghuzhanghao' # 登录字段名

    # 字段定义
    yonghuzhanghao = models.CharField(max_length=255, unique=True, verbose_name='用户账号')
    mima = models.CharField(max_length=255, verbose_name='密码')
    yonghuxingming = models.CharField(max_length=255, verbose_name='用户姓名')
    xingbie = models.CharField(max_length=255, verbose_name='性别')
    touxiang = models.TextField(verbose_name='头像')
    lianxifangshi = models.CharField(max_length=255, verbose_name='联系方式')

    class Meta:
    db_table = 'yonghu'
    verbose_name = '用户'

    讲解要点:

    • __loginUserColumn__指定登录时使用的账号字段
    • unique=True保证账号唯一性
    • verbose_name设置字段在后台显示的中文名称
    旅游路线模型 lvyouluxian

    # models.py 第70-121行
    class lvyouluxian(BaseModel):
    __tablename__ = 'lvyouluxian'

    # 授权表映射
    __authTables__ = {'yonghuzhanghao': 'yonghu'}

    # 智能推荐属性(使用协同过滤)
    __intelRecom__ = '用协'

    # 前台列表可查看
    __foreEndList__ = '是'

    # 字段定义
    luxianmingcheng = models.CharField(max_length=255, verbose_name='路线名称')
    luxianfenlei = models.CharField(max_length=255, verbose_name='路线分类')
    luxiantupian = models.TextField(verbose_name='路线图片')
    xingchenganpai = models.CharField(max_length=255, verbose_name='行程安排')
    qidian = models.CharField(max_length=255, verbose_name='起点')
    zhongdian = models.CharField(max_length=255, verbose_name='终点')
    feiyongyusuan = models.FloatField(verbose_name='费用预算')
    tujingjingdian = models.CharField(max_length=255, verbose_name='途经景点')
    jiaotongjianyi = models.CharField(max_length=255, verbose_name='交通建议')
    luxianxiangqing = models.TextField(verbose_name='路线详情')
    yonghuzhanghao = models.CharField(max_length=255, verbose_name='用户账号')

    # 推荐相关字段
    clicktime = models.DateTimeField(auto_now=True, verbose_name='最近点击时间')
    discussnum = models.IntegerField(default='0', verbose_name='评论数')
    storeupnum = models.IntegerField(default='0', verbose_name='收藏数')

    讲解要点:

    • __intelRecom__ = '用协'表示使用协同过滤算法推荐
    • __authTables__映射用户账号到用户表,实现权限关联
    • clicktime记录最近点击时间,用于热度推荐
    • discussnum和storeupnum统计评论和收藏数量
    景点信息模型 sightinfo

    # models.py 第122-171行
    class sightinfo(BaseModel):
    __tablename__ = 'sightinfo'

    poiname = models.CharField(max_length=255, verbose_name='景点名称')
    imgurl = models.TextField(verbose_name='图片')
    commentscore = models.FloatField(verbose_name='评分')
    commentcount = models.IntegerField(verbose_name='评论数')
    heatscore = models.FloatField(verbose_name='热度')
    jiage = models.FloatField(verbose_name='价格')
    districtname = models.CharField(max_length=255, verbose_name='城市')
    zonename = models.CharField(max_length=255, verbose_name='地区')
    sightcategoryinfo = models.CharField(max_length=255, verbose_name='口碑榜')
    sightlevel = models.CharField(max_length=255, verbose_name='等级')
    features = models.CharField(max_length=255, verbose_name='特征')
    tagname = models.CharField(max_length=255, verbose_name='标签')

    讲解要点:

    • 景点信息表用于Hadoop数据分析
    • commentscore评分字段用于计算平均评分
    • heatscore热度字段用于热门景点统计
    • districtname城市字段用于地区分布统计

    3. 视图控制器模块

    旅游路线控制器 (main/Lvyouluxian_v.py)
    分页查询接口

    # 第141-157行
    def lvyouluxian_page(request):
    if request.method in ["POST", "GET"]:
    msg = {"code": normal_code, "msg": mes.normal_code}
    req_dict = request.session.get("req_dict")

    # 当前登录用户信息
    tablename = request.session.get("tablename")
    if tablename == 'yonghu':
    # 用户只能查看自己发布的路线
    req_dict['yonghuzhanghao'] = request.session.get("params").get(yonghu.__loginUserColumn__)

    # 调用分页方法
    msg['data']['list'], msg['data']['currPage'], msg['data']['totalPage'],
    msg['data']['total'], msg['data']['pageSize'] = lvyouluxian.page(lvyouluxian, lvyouluxian, req_dict, request)

    return JsonResponse(msg, encoder=CustomJsonEncoder)

    讲解要点:

    • 从request.session获取请求参数字典
    • 如果是普通用户登录,自动过滤只显示该用户的路线
    • 使用JsonResponse返回JSON格式数据
    详情查询接口

    # 第434-468行
    def lvyouluxian_detail(request, id_):
    if request.method in ["POST", "GET"]:
    msg = {"code": normal_code, "msg": mes.normal_code}

    # 获取数据
    data = lvyouluxian.getbyid(lvyouluxian, lvyouluxian, int(id_))
    if len(data) > 0:
    msg['data'] = data[0]

    # 浏览点击次数更新
    try:
    __browseClick__ = lvyouluxian.__browseClick__
    except:
    __browseClick__ = None

    if __browseClick__ == "是" and "clicknum" in lvyouluxian.getallcolumn():
    clicknum = int(data[0].get("clicknum", 0)) + 1
    click_dict = {"id": int(id_), "clicknum": clicknum, "clicktime": datetime.datetime.now()}
    lvyouluxian.updatebyparams(lvyouluxian, lvyouluxian, click_dict)

    return JsonResponse(msg)

    讲解要点:

    • 根据ID查询路线详情
    • __browseClick__属性控制是否记录点击次数
    • 每次访问详情页自动增加点击次数,更新点击时间
    协同过滤推荐接口

    # 第571-632行
    import math

    # 余弦相似度计算
    def cosine_similarity(a, b):
    numerator = sum([a[key] * b[key] for key in a if key in b])
    denominator = math.sqrt(sum([a[key]**2 for key in a])) * math.sqrt(sum([b[key]**2 for key in b]))
    return numerator / denominator

    # 协同过滤推荐
    def lvyouluxian_autoSort2(request):
    if request.method in ["POST", "GET"]:
    req_dict = request.session.get("req_dict")
    cursor = connection.cursor()

    # 查询收藏数据
    cursor.execute("select * from storeup where type = 1 and tablename = 'lvyouluxian'")
    data_dict = [dict(zip([col[0] for col in desc], row)) for row in cursor.fetchall()]

    # 构建用户-路线矩阵
    user_ratings = {}
    for item in data_dict:
    if user_ratings.__contains__(item["userid"]):
    ratings_dict = user_ratings[item["userid"]]
    if ratings_dict.__contains__(item["refid"]):
    ratings_dict[str(item["refid"])] += 1
    else:
    ratings_dict[str(item["refid"])] = 1
    else:
    user_ratings[item["userid"]] = {str(item["refid"]): 1}

    # 计算用户相似度
    current_user_id = request.session.get("params").get("id")
    similarities = {other_user: cosine_similarity(user_ratings[current_user_id], user_ratings[other_user])
    for other_user in user_ratings if other_user != current_user_id}

    # 找最相似用户
    most_similar_user = sorted(similarities, key=similarities.get, reverse=True)[0]

    # 推荐该用户收藏但当前用户未收藏的路线
    recommended_goods = {goods: rating for goods, rating in user_ratings[most_similar_user].items()
    if goods not in user_ratings[current_user_id]}

    sorted_recommended_goods = sorted(recommended_goods, key=recommended_goods.get, reverse=True)

    return JsonResponse({"code": 0, "data": {"list": L[0:int(req_dict["limit"])]}})

    讲解要点:

    • 协同过滤算法原理:基于用户收藏行为相似度推荐
    • 用户-路线矩阵:记录每个用户收藏的路线及权重
    • 余弦相似度:计算两个用户收藏行为的相似程度
    • 推荐逻辑:找到最相似用户,推荐他收藏但当前用户未收藏的路线

    余弦相似度公式:
    similarity=∑iAi×Bi∑iAi2×∑iBi2similarity = \\frac{\\sum_{i} A_i \\times B_i}{\\sqrt{\\sum_i A_i^2} \\times \\sqrt{\\sum_i B_i^2}}similarity=iAi2×iBi2iAi×Bi


    4. 路由配置模块 (main/urls.py)

    动态路由生成

    # 第43-102行
    # 自动扫描视图文件生成路由
    for i in os.listdir(mainDir):
    if i not in excludeList and i[5:] == "_v.py":
    tableName = i[:5] # 去掉_v.py后缀

    urlpatterns.extend([
    path(r'{}/register'.format(tableName.lower()),
    eval("{}_v.{}_register".format(tableName.capitalize(), tableName.lower()))),
    path(r'{}/login'.format(tableName.lower()),
    eval("{}_v.{}_login".format(tableName.capitalize(), tableName.lower()))),
    path(r'{}/page'.format(tableName.lower()),
    eval("{}_v.{}_page".format(tableName.capitalize(), tableName.lower()))),
    path(r'{}/list'.format(tableName.lower()),
    eval("{}_v.{}_list".format(tableName.capitalize(), tableName.lower()))),
    path(r'{}/add'.format(tableName.lower()),
    eval("{}_v.{}_add".format(tableName.capitalize(), tableName.lower()))),
    path(r'{}/info/<id_>'.format(tableName.lower()),
    eval("{}_v.{}_info".format(tableName.capitalize(), tableName.lower()))),
    path(r'{}/detail/<id_>'.format(tableName.lower()),
    eval("{}_v.{}_detail".format(tableName.capitalize(), tableName.lower()))),
    path(r'{}/update'.format(tableName.lower()),
    eval("{}_v.{}_update".format(tableName.capitalize(), tableName.lower()))),
    path(r'{}/delete'.format(tableName.lower()),
    eval("{}_v.{}_delete".format(tableName.capitalize(), tableName.lower()))),
    ])

    讲解要点:

    • 动态路由机制:自动扫描_v.py结尾的视图文件
    • 路由命名规则:小写表名作为URL路径
    • eval动态执行:根据文件名动态生成视图函数调用
    • 自动生成的标准接口包括:register、login、page、list、add、info、detail、update、delete

    5. 认证中间件 (xmiddleware/xauth.py)

    # 第12-94行
    class Xauth(MiddlewareMixin):
    def process_request(self, request):
    fullPath = request.get_full_path()

    # WebSocket请求不拦截
    if request.META.get('HTTP_UPGRADE') == 'websocket':
    return

    # GET请求白名单过滤
    if request.method == 'GET':
    filterList = [
    "/index",
    "/login",
    "/register",
    "/detail",
    ".js",
    ".css",
    ".jpg",
    ".png",
    # …静态资源
    ]

    # 前台列表接口不拦截
    for m in allModels:
    foreEndList = m.__foreEndList__
    if foreEndList is None or foreEndList != "前要登":
    filterList.append("/{}/list".format(m.__tablename__))
    filterList.append("/{}/detail".format(m.__tablename__))

    # 检查是否需要认证
    auth = True
    for i in filterList:
    if i in fullPath:
    auth = False

    if auth == True:
    result = Auth.identify(Auth, request)
    if result.get('code') != normal_code:
    return JsonResponse(result)

    讲解要点:

    • 中间件作用:拦截请求验证用户身份
    • 白名单机制:登录、注册、静态资源不需要认证
    • __foreEndList__属性决定前台列表是否需要登录
    • 未认证请求返回401错误

    6. 用户认证工具 (util/auth.py)

    # 第11-27行
    class Auth(object):
    def authenticate(self, model, req_dict):
    """
    用户登录认证
    """

    msg = {'code': normal_code, 'msg': mes.normal_code}
    tablename = model.__tablename__

    # 生成Token
    encode_dict = {"tablename": tablename, "params": req_dict}
    encode_str = base64.b64encode(str(encode_dict).encode("utf-8"))
    msg['token'] = encode_str.decode('utf-8')

    return JsonResponse(msg)

    def identify(self, request):
    """
    用户身份验证
    """

    token = request.META.get('HTTP_TOKEN')

    if token and token != "null":
    # 解码Token
    decode_str = base64.b64decode(token).decode("utf8")
    decode_dict = eval(decode_str)

    tablename2 = decode_dict.get("tablename")
    params2 = decode_dict.get("params")

    # 验证用户是否存在
    datas = model.getbyparams(model, model, params2)
    if not datas:
    return {'code': 401, 'msg': '找不到该用户信息'}
    else:
    request.session['tablename'] = tablename2
    request.session['params'] = params2
    return {'code': normal_code, 'msg': '身份验证通过'}

    讲解要点:

    • Token机制:使用Base64编码用户信息生成Token
    • Token包含用户表名和用户参数
    • 身份验证时解码Token,查询用户是否存在
    • 验证成功后将用户信息存入session

    数据库设计

    数据库配置 (config.ini)

    [sql]
    type = mysql
    host = 127.0.0.1
    port = 3306
    user = root
    passwd = 123456
    db = djangoy1pr16u7
    charset = utf8
    hasHadoop = 是

    核心数据表

    表名说明主要字段
    yonghu 用户表 yonghuzhanghao(账号), mima(密码), yonghuxingming(姓名)
    lvyouluxian 旅游路线表 luxianmingcheng(名称), luxianfenlei(分类), feiyongyusuan(预算)
    luxianfenlei 路线分类表 luxianfenlei(分类名)
    sightinfo 景点信息表 poiname(名称), commentscore(评分), heatscore(热度)
    commentinfo 评论信息表 usernick(用户), plcontent(内容), score(评分)
    news 旅游资讯表 title(标题), typename(分类), content(内容)
    storeup 收藏表 userid(用户ID), refid(关联ID), tablename(表名)
    discusslvyouluxian 路线评论表 refid(路线ID), content(内容), userid(用户ID)

    ER关系图说明

    用户表(yonghu)

    ├── 发布 ──→ 旅游路线表(lvyouluxian)
    │ │
    │ ├── 收藏 ──→ 收藏表(storeup)
    │ │
    │ └── 评论 ──→ 路线评论表(discusslvyouluxian)

    └── 收藏 ──→ 景点信息表(sightinfo)


    关键功能实现

    1. 用户登录流程

    前端发起请求 ──→ xauth中间件(放行) ──→ users_login视图


    查询用户数据(users.getbyparams)


    生成Token(Auth.authenticate)


    返回Token给前端保存

    代码流程:

    # users_v.py 第11-25行
    def users_login(request):
    req_dict = request.session.get("req_dict")

    # 查询用户
    datas = users.getbyparams(users, users, req_dict)
    if not datas:
    return JsonResponse({'code': password_error_code, 'msg': '密码错误'})

    # 生成Token
    req_dict['id'] = datas[0].get('id')
    return Auth.authenticate(Auth, users, req_dict)

    2. 数据新增流程

    前端提交数据 ──→ xparam中间件(参数处理) ──→ xxx_save视图


    验证字段类型(model._meta.fields)


    自动填充字段(addtime, userid)


    保存到数据库(model.save)


    返回新增记录ID

    3. 分页查询流程

    前端请求(page, limit) ──→ xxx_page视图 ──→ 参数解析


    构建查询条件(filter)


    应用排序(order_by)


    Django Paginator分页


    转换为字典列表(to_list)


    返回分页数据JSON


    Hadoop大数据分析模块

    模块概述

    本系统使用Hadoop MapReduce对景点信息和评论数据进行统计分析,主要实现:

  • 景点评分统计:按景点名称统计平均评分
  • 景点评论数统计:按景点名称统计评论总数
  • 景点热度统计:按景点名称统计热度总和
  • 景点价格统计:按景点名称统计平均价格
  • 城市分布统计:按城市统计景点数量
  • IP属地统计:按IP属地统计评论数量
  • 旅游类型统计:按旅游类型统计评论数量
  • hadoop_v.py 核心代码解析

    # 第35-87行 – 数据上传到HDFS
    def upload_csv_mapreduce_hadoop():
    # 查询景点数据
    query = "SELECT * FROM sightinfo"
    df = pd.read_sql(query, connection)

    # 导出CSV文件
    local_csv_path = os.path.join(parent_directory, "sightinfo.csv")
    df.to_csv(local_csv_path, index=False)

    # 上传到HDFS
    hdfs_csv_path = '/input/sightinfo.csv'
    hadoop_client.upload(hdfs_csv_path, local_csv_path)

    # 上传MapReduce代码
    hadoop_client.upload('/input/group_mapper.py', group_mapper_local_path)
    hadoop_client.upload('/input/group_reducer.py', group_reducer_local_path)

    讲解要点:

    • 使用pandas从MySQL读取数据导出为CSV
    • HDFS客户端上传数据文件到分布式文件系统
    • 同时上传MapReduce的Mapper和Reducer代码

    # 第89-159行 – 执行MapReduce作业
    def send_cmd():
    job_commands = [
    # 景点评分统计
    [
    f"{hadoop_path}/bin/hadoop.cmd", "jar", f"{hadoop_path}/share/hadoop/tools/lib/hadoop-streaming-3.3.0.jar",
    "-files", "\\"hdfs://localhost:9000/input/value_mapper.py,hdfs://localhost:9000/input/value_reducer.py\\"",
    "-mapper", f"\\"python value_mapper.py {csv_index('sightinfo.csv','poiname')} {csv_index('sightinfo.csv','commentscore')} 无\\"",
    "-reducer", "\\"python value_reducer.py poiname\\"",
    "-input", "hdfs://localhost:9000/input/sightinfo.csv",
    "-output", "hdfs://localhost:9000/output/sightinfo/valuepoinamecommentscore"
    ],
    # …其他分析任务
    ]

    # 多进程并行执行
    processes = []
    for job_command in job_commands:
    p = multiprocessing.Process(target=run_mapreduce_job_on_remote, args=(job_command, table_name, fileName))
    p.start()
    processes.append(p)

    for p in processes:
    p.join()

    讲解要点:

    • 使用Hadoop Streaming执行Python MapReduce
    • -mapper参数指定Mapper脚本和字段索引
    • -reducer参数指定Reducer脚本
    • 多进程并行执行多个分析任务

    MapReduce代码解析

    group_mapper.py (分组统计Mapper)

    # 第6-18行
    index = int(sys.argv[1]) # 获取字段索引

    for line in sys.stdin:
    try:
    line = line.strip()
    lists = line.split(',')
    if lists[0] == "id": # 跳过表头
    continue
    print(f"{lists[index]}\\t") # 输出分组字段
    except Exception as e:
    continue

    讲解要点:

    • 从命令行参数获取要统计的字段索引
    • 读取CSV文件的每一行
    • 跳过第一行表头(id字段)
    • 输出指定字段值作为Mapper结果
    group_reducer.py (分组统计Reducer)

    # 第6-23行
    c_name = sys.argv[1] # 获取字段名

    for line in sys.stdin:
    line = line.strip()
    role = line.split('\\t')[0] # 获取分组值

    # 统计每个值出现次数
    d = next((item for item in json_list if item.get(c_name) == role), None)
    if d:
    d["total"] += 1
    else:
    json_list.append({c_name: role, "total": 1})

    print(json.dumps(json_list, ensure_ascii=False))

    讲解要点:

    • 接收Mapper输出结果
    • 对每个分组值进行计数统计
    • 输出JSON格式的统计结果
    value_mapper.py (数值统计Mapper)

    # 第6-27行
    x_index = int(sys.argv[1]) # X轴字段索引(分组字段)
    y_index = sys.argv[2] # Y轴字段索引(数值字段)

    for line in sys.stdin:
    lists = line.strip().split(',')
    if lists[0] == "id":
    continue

    x_value = lists[x_index] # 分组字段值
    y_value = lists[int(y_index)] # 数值字段值

    print(f"{x_value}\\t{y_value}")

    讲解要点:

    • 提取分组字段和数值字段
    • 输出两个字段供Reducer聚合
    value_reducer.py (数值统计Reducer)

    # 第7-38行
    x_name = sys.argv[1]

    for line in sys.stdin:
    xname = line.split('\\t')[0] # 分组字段值
    yname = line.split('\\t')[1] # 数值

    d = next((item for item in json_list if item.get(x_name) == xname), None)
    if d:
    d["total"] = str(round(Decimal(d["total"]) + Decimal(yname), 2)) # 累加
    else:
    json_list.append({x_name: xname, "total": str(round(Decimal(yname), 2))})

    print(json.dumps(json_list, ensure_ascii=False))

    讲解要点:

    • 对相同分组值的数值进行累加
    • 使用Decimal保证精度
    • 输出每个分组的数值总和

    MapReduce数据流图

    MySQL数据


    导出CSV文件


    上传到HDFS (/input/)


    MapReduce任务

    ├── Mapper ──→ 读取数据,提取字段

    ├── Shuffle ──→ 按key分组传输

    ├── Reducer ──→ 聚合统计


    输出到HDFS (/output/)


    下载结果JSON文件


    协同过滤推荐算法

    算法原理

    协同过滤(Collaborative Filtering)是一种基于用户行为的推荐算法。本系统采用基于用户的协同过滤(User-based CF):

  • 收集用户收藏行为数据
  • 计算用户之间的相似度
  • 找到与当前用户最相似的用户
  • 推荐相似用户收藏但当前用户未收藏的内容
  • 相似度计算 – 余弦相似度

    # Lvyouluxian_v.py 第571-574行
    import math

    def cosine_similarity(a, b):
    # 计算向量内积
    numerator = sum([a[key] * b[key] for key in a if key in b])

    # 计算向量模长乘积
    denominator = math.sqrt(sum([a[key]**2 for key in a])) * math.sqrt(sum([b[key]**2 for key in b]))

    return numerator / denominator

    数学公式:
    cos(A,B)=A⋅B∣A∣×∣B∣=∑iAiBi∑iAi2∑iBi2cos(A, B) = \\frac{A \\cdot B}{|A| \\times |B|} = \\frac{\\sum_i A_i B_i}{\\sqrt{\\sum_i A_i^2} \\sqrt{\\sum_i B_i^2}}cos(A,B)=A×BAB=iAi2iBi2iAiBi

    讲解要点:

    • 余弦相似度衡量两个向量方向的相似程度
    • 值范围在[-1, 1],越接近1表示越相似
    • 用户向量:用户收藏的所有路线及其收藏次数

    推荐实现流程

    # Lvyouluxian_v.py 第577-632行
    def lvyouluxian_autoSort2(request):
    # 1. 查询所有收藏数据
    cursor.execute("select * from storeup where type = 1 and tablename = 'lvyouluxian'")

    # 2. 构建用户-路线评分矩阵
    user_ratings = {}
    for item in data_dict:
    if user_ratings.__contains__(item["userid"]):
    ratings_dict = user_ratings[item["userid"]]
    ratings_dict[str(item["refid"])] = 1
    else:
    user_ratings[item["userid"]] = {str(item["refid"]): 1}

    # 3. 计算当前用户与其他用户相似度
    current_user_id = request.session.get("params").get("id")
    similarities = {other_user: cosine_similarity(user_ratings[current_user_id], user_ratings[other_user])
    for other_user in user_ratings if other_user != current_user_id}

    # 4. 找最相似用户
    most_similar_user = sorted(similarities, key=similarities.get, reverse=True)[0]

    # 5. 推荐路线
    recommended_goods = {goods: rating for goods, rating in user_ratings[most_similar_user].items()
    if goods not in user_ratings[current_user_id]}

    # 6. 按评分排序返回
    sorted_recommended_goods = sorted(recommended_goods, key=recommended_goods.get, reverse=True)

    推荐流程图

    用户收藏数据


    构建用户-路线矩阵


    计算用户相似度


    找出最相似用户


    筛选推荐路线


    排序返回结果

    算法示例

    假设有以下收藏数据:

    用户收藏路线
    用户A 路线1, 路线2, 路线3
    用户B 路线1, 路线2, 路线4
    用户C 路线2, 路线5

    用户-路线矩阵:

    用户A: {路线1: 1, 路线2: 1, 路线3: 1}
    用户B: {路线1: 1, 路线2: 1, 路线4: 1}
    用户C: {路线2: 1, 路线5: 1}

    相似度计算:

    • 用户A与用户B:similarity = 2/√9 * √9 = 2/3 ≈ 0.67
    • 用户A与用户C:similarity = 1/√9 * √4 = 1/6 ≈ 0.17

    推荐结果:

    • 用户A的最相似用户是用户B
    • 用户B收藏了路线4,用户A未收藏
    • 推荐路线4给用户A

    API接口设计

    接口规范

    所有接口遵循RESTful风格,返回JSON格式数据:

    {
    "code": 0, // 状态码:0成功,其他失败
    "msg": "操作成功", // 消息
    "data": { // 数据
    "list": [], // 列表数据
    "total": 100, // 总数
    "currPage": 1, // 当前页
    "totalPage": 10, // 总页数
    "pageSize": 10 // 每页条数
    }
    }

    主要接口列表

    接口路径方法说明
    /{table}/login POST 用户登录
    /{table}/register POST 用户注册
    /{table}/page GET 后台分页查询
    /{table}/list GET 前台列表查询
    /{table}/detail/{id} GET 详情查询
    /{table}/add POST 新增数据
    /{table}/update POST 更新数据
    /{table}/delete POST 删除数据
    /{table}/autoSort GET 智能推荐(热度)
    /{table}/autoSort2 GET 协同过滤推荐
    /hadoop/analyze GET 执行数据分析

    接口调用示例

    用户登录

    POST /djangoy1pr16u7/yonghu/login
    Content-Type: application/json

    {
    "yonghuzhanghao": "user001",
    "mima": "123456"
    }

    响应:

    {
    "code": 0,
    "msg": "成功",
    "token": "eyJ0YWJsZW5hbWUiOi…"
    }

    路线列表查询

    GET /djangoy1pr16u7/lvyouluxian/list?page=1&limit=10&luxianfenlei=自然风光
    Token: eyJ0YWJsZW5hbWUiOi…

    响应:

    {
    "code": 0,
    "data": {
    "list": [
    {
    "id": 164800001,
    "luxianmingcheng": "西湖一日游",
    "luxianfenlei": "自然风光",
    "feiyongyusuan": 200
    }
    ],
    "total": 50,
    "currPage": 1,
    "totalPage": 5
    }
    }


    前端页面结构

    页面目录

    templates/front/
    ├── index.html # 前台首页
    ├── admin/ # 后台管理
    │ ├── index.html # 后台首页(jQuery版)
    │ ├── dist/ # Vue版后台
    │ │ └── index.html
    │ └── pages/ # 后台子页面
    │ ├── lvyouluxian_list.html
    │ ├── lvyouluxian_add.html
    │ ├── yonghu_list.html
    │ └── …
    ├── pages/ # 前台页面
    │ ├── lvyouluxian_list.html # 路线列表
    │ ├── lvyouluxian_detail.html # 路线详情
    │ ├── sightinfo_list.html # 景点列表
    │ ├── login.html # 登录页
    │ ├── register.html # 注册页
    │ └── …
    └── assets/ # 静态资源
    ├── css/
    ├── js/
    ├── img/
    └── layui/ # Layui组件

    前端技术特点

  • Layui组件库:用于后台管理界面的表格、表单、弹窗等
  • Vue.js:部分后台页面使用Vue构建
  • jQuery:DOM操作和Ajax请求
  • 响应式设计:适配不同屏幕尺寸

  • 系统部署与运行

    环境要求

    环境版本要求
    Python 3.6+
    MySQL 5.5+
    Hadoop 3.3.0
    JDK 1.8+

    安装步骤

  • 安装Python依赖
  • pip install -r requirements.txt

  • 配置数据库
    修改config.ini文件中的数据库连接信息

  • 初始化数据库

  • python init.py initsql

  • 启动Hadoop
  • # 启动HDFS
    start-dfs.cmd
    # 启动YARN
    start-yarn.cmd

  • 运行项目
  • python manage.py runserver 0.0.0.0:8000

    运行脚本

    系统提供批处理脚本:

    • 安装.bat:安装依赖包
    • 运行.bat:启动项目
    • 初始化hive数据库.bat:初始化Hive

    答辩常见问题解答

    Q1:系统使用了哪些大数据技术?

    答:系统主要使用了以下大数据技术:

  • Hadoop HDFS:分布式存储景点数据和评论数据
  • Hadoop MapReduce:对数据进行统计分析(评分统计、热度统计、分布统计等)
  • Spark:预留的机器学习功能(聚类、分类、回归)
  • Hive:数据仓库工具,支持SQL查询
  • Q2:协同过滤算法是如何实现的?

    答:系统采用基于用户的协同过滤算法:

  • 首先收集用户的收藏行为数据,构建用户-路线矩阵
  • 使用余弦相似度计算用户之间的相似程度
  • 找到与当前用户最相似的K个用户
  • 推荐这些相似用户收藏但当前用户未收藏的路线
  • 核心代码位于Lvyouluxian_v.py的lvyouluxian_autoSort2函数和cosine_similarity函数。

    Q3:MapReduce分析的具体实现是什么?

    答:

  • 数据准备:从MySQL导出CSV文件,上传到HDFS
  • Mapper阶段:读取CSV数据,提取需要统计的字段
  • Reducer阶段:对相同key的数据进行聚合(计数或求和)
  • 结果输出:将统计结果以JSON格式保存
  • 具体包括:

    • group_mapper.py/group_reducer.py:分组计数统计
    • value_mapper.py/value_reducer.py:数值求和统计

    Q4:系统是如何实现用户认证的?

    答:系统使用Token机制进行用户认证:

  • 用户登录成功后,系统使用Base64编码用户信息生成Token
  • 前端将Token保存在localStorage或cookie中
  • 后续请求在HTTP Header中携带Token
  • xauth中间件拦截请求,解码Token验证用户身份
  • 验证成功后将用户信息存入session供后续使用
  • Q5:数据库设计有哪些表?关系是什么?

    答:主要数据表包括:

    表名说明
    yonghu 用户表
    lvyouluxian 旅游路线表
    luxianfenlei 路线分类表
    sightinfo 景点信息表
    commentinfo 评论信息表
    news 旅游资讯表
    storeup 收藏表
    discusslvyouluxian 路线评论表

    关系:

    • 用户(yonghu)可以发布路线(lvyouluxian)
    • 用户可以收藏路线、景点、资讯(storeup)
    • 用户可以对路线发表评论(discusslvyouluxian)
    • 路线属于某个分类(luxianfenlei)

    Q6:系统有哪些创新点?

    答:

  • 大数据技术应用:使用Hadoop MapReduce对海量景点数据进行统计分析
  • 协同过滤推荐:基于用户收藏行为实现个性化路线推荐
  • 动态路由生成:自动扫描视图文件生成API路由,减少重复代码
  • 模块化设计:采用中间件模式处理认证和参数,代码结构清晰
  • Q7:系统有什么可以改进的地方?

    答:

  • 推荐算法优化:可以结合基于内容的推荐,实现混合推荐策略
  • Spark MLlib应用:可以使用Spark机器学习库进行更复杂的分析
  • 实时推荐:引入实时流处理框架实现实时推荐
  • 前端优化:完全迁移到Vue/React框架,提升用户体验
  • 安全性增强:引入JWT标准Token,增加密码加密强度
  • Q8:如何理解BaseModel基类的page方法?

    答:page方法是分页查询的核心方法:

  • 参数解析:从params获取page(页码)、limit(每页数量)、sort(排序字段)、order(排序方向)
  • 条件构建:根据参数构建Django ORM的filter条件
  • 模糊查询:处理%符号的模糊查询需求
  • 分页实现:使用Django Paginator进行分页
  • 数据转换:将QuerySet转换为字典列表便于JSON序列化
  • Q9:中间件的作用是什么?

    答:

  • xparam中间件:预处理请求参数,将请求体数据存入session
  • xauth中间件:拦截需要认证的请求,验证Token有效性
  • 执行顺序:请求依次通过Security→Session→Common→xparam→xauth→View
  • Q10:如何扩展新的数据表?

    答:

  • 在models.py中定义新的模型类,继承BaseModel
  • 设置表名__tablename__和字段定义
  • 运行数据库迁移创建表
  • 创建对应的视图文件{tablename}_v.py
  • 系统会自动扫描生成对应的API路由

  • 总结

    本系统是一个完整的旅游路线推荐平台,融合了传统Web开发与大数据技术。核心亮点包括:

  • Django框架提供稳定的Web服务
  • Hadoop MapReduce实现大数据统计分析
  • 协同过滤算法实现个性化推荐
  • 模块化架构便于维护和扩展
  • 希望本文档能帮助您顺利完成答辩!


    支持一对一定制项目哈!

    赞(0)
    未经允许不得转载:171主机测评 » 基于Spark的旅游路线推荐系统代码讲解文档
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址