ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

基于Flask+Hadoop+Hive的股票大数据分析系统设计与部署实践

基于Flask+Hadoop+Hive的股票大数据分析系统设计与部署实践 简介一套基于 Python、Flask、Hadoop 与 Hive 的股票大数据分析系统项目面向计算机相关专业学生的毕业设计、课程设计、工程实训和项目初期演示等场景解决从股票数据采集、离线存储、分布式分析计算到前端可视化展示的完整链路问题。压缩包共六十个文件整体大小约 451KB其中包含 27 个以 Python 编写的后端核心源码文件、9 个前端交互脚本以及页面结构、样式表、配置说明和部署文档等辅助材料目录分层清晰便于按功能模块查阅。系统提供 Flask 系统部署文档、Hive 相关配置说明、日志文件、依赖列表和测试脚本能够帮助读者快速理解从启动服务到查询分析的流程同时代码结构规整具备进一步修改和扩展的基础。项目已经过测试运行已有 315 人学习下载适合作为高分毕设参考也可作为大数据分析入门和全栈开发实践的良好案例。1. 股票大数据分析为什么绕不开 Hadoop 与 Hive 这套组合做股票数据分析的毕设最常见的翻车点不是算法写不出来而是数据量稍微一上来单机 Pandas 直接从内存不够变成跑一次要等到下课。这套基于 Python Flask Hadoop Hive 的股票大数据分析系统把存储和计算下沉到 Hadoop 生态用 Hive 做数据仓库Flask 只负责对外提供查询与展示接口是大数据课设里能直接跑通全链路的参考实现。它解决的核心问题是数据链路太长导致的项目失真很多人把股票分析做成 Excel 加 Pandas 的本地作业完全没体现大数据组件的价值。而 Hadoop 负责分布式存储与计算Hive 把复杂数据处理变成 SQLFlask 让分析师能通过网页看到聚合结果三个角色各管一段。适合的人群很清楚第一次接触 Hadoop 生态、想把大数据组件串成一个完整项目的学生需要课程设计或毕设代码做二次开发的人以及想搞清楚 Flask 应用层和 HiveServer2 到底怎么对接的开发者。整体数据流是行情数据落到 Hive 外部表Hive SQL 完成清洗聚合Flask 拿到结果后以 JSON 返回前端。2. 先拆目录Flask 应用层、控制层与 Hive 数据层的分层逻辑拿到压缩包解压以后目录结构比一般课设要规整得多。viewhandler、control、dbmodel、common、templates、static 各司其职不是把代码堆在几个 Python 文件里。这种分层带来的直接好处是答辩评审问「某个功能在哪里实现」时可以按目录指过去路由在 viewhandler业务在 controlSQL 在 dbmodel。整套架构围绕一个原则展开视图层不写 SQL数据层不写业务控制层做中转。2.1 从目录结构反推系统分层Flask-Hive-master 根目录下viewhandler 对应 Web 路由层内部有 page_blueprint.py、user_blueprint.py、data_blueprint.py 三个蓝图文件分别负责页面跳转、用户认证和数据接口。control 目录是业务控制层basecontrol.py、usercontrol.py、datacontrol.py 把公共逻辑、用户逻辑和数据分析逻辑拆开。dbmodel 目录是数据访问层dbbase.py 封装连接datamodel.py 定义数据模型sql_tpl.py 统一管理 Hive SQL。common 里的 response.py 和 auth.py 提供统一返回格式与登录鉴权。目录/模块分层定位核心文件与职责viewhandlerWeb 路由层页面蓝图、用户蓝图、数据蓝图负责路由与参数接收control业务控制层用户控制、数据控制、基础控制复用业务逻辑dbmodel数据访问层连接封装、数据模型、SQL 模板common通用支撑响应包装、权限校验static / templates前端资源base.html、main.html、header.html 及 css/js把 SQL 和业务逻辑分开的最大好处是排查问题快。前端页面报错先看 viewhandler 对应蓝图返回数据不对再去 datacontrol.py 里看数据处理如果确认是 Hive 侧的问题打开 sql_tpl.py 就能看到全部查询语句。项目里同时存在 user_blueprint.pyc、page_blueprint.pyc 这些字节码文件说明作者是在本地跑过之后才打包的这本身也能佐证代码可用。2.2 控制层为什么单独存在很多入门项目会把数据库查询直接写在 Flask 路由函数里这个项目没这么做。control 层单独存在有两个实际理由一是多个蓝图可能要复用同一套统计逻辑比如页面端和 API 端都调用 DataControl 的同一方法避免复制代码二是控制层不和 Flask 的 request/session 绑定可以脱离 Web 容器单独写测试脚本验证逻辑。这一点在后期维护时价值很大想给 datacontrol.py 加单元测试也不用启动整个 Flask 服务。下面是一段按常见做法整理的三层调用关系和项目里 data_blueprint.py 的写法保持一致# viewhandler/data_blueprint.py 中的路由封装常见做法 from flask import Blueprint, request from control.datacontrol import DataControl from common.response import success, fail data_bp Blueprint(data, __name__) data_bp.route(/stock/kline) def stock_kline(): # 视图层只做参数接收和结果包装 code request.args.get(code, ) start request.args.get(start, ) end request.args.get(end, ) if not code or not start or not end: return fail(参数不完整) data DataControl().get_kline(code, start, end) return success(data)这段代码的逻辑是data_bp 路由接收请求参数DataControl 封装业务查询common.response 统一返回 JSON。参数 code 是六位股票代码start 和 end 是 yyyy-MM-dd 格式。视图函数里不出现任何 Hive 连接代码SQL 和连接都在下层。如果哪天要把数据源从 Hive 换成 ClickHouse只需要改 dbmodel 和 datacontrol路由层完全不用动。2.3 入口与配置bootstrap.py 与 config.py 的分工bootstrap.py 在整个项目中承担应用初始化。常见做法是用 Flask 的应用工厂模式把所有蓝图的注册集中在一个函数里# bootstrap.py 应用工厂入口 from flask import Flask from viewhandler.page_blueprint import page_bp from viewhandler.user_blueprint import user_bp from viewhandler.data_blueprint import data_bp def create_app(): app Flask(__name__) app.config.from_pyfile(config.py) # 页面、用户、数据三类路由分开注册 app.register_blueprint(page_bp, url_prefix/) app.register_blueprint(user_bp, url_prefix/user) app.register_blueprint(data_bp, url_prefix/api) return app if __name__ __main__: create_app().run(host0.0.0.0, port8080, debugFalse)url_prefix 把三组业务路径隔开页面走 /用户相关走 /user数据接口走 /api。这样 Flask 系统部署时只需要保证 config.py 里的 HiveServer2 地址、端口和数据库名正确其他代码不用动。分层的意义在这里体现换部署环境动的是配置文件和 Hive 表而不是视图函数。3. Hive 建表与查询接口从存储设计到 Flask 拿到 JSON股票行情数据的特点是文件多、行数多、字段相对固定这正好是 Hive 擅长处理的场景。Hive 在这里相当于把 HDFS 上的文本文件包装成了一张可查询的二维表SQL 写在前面分布式计算在底层执行。想要让 Flask 接口稳定先把表结构和 SQL 边界划清楚。3.1 股票日线数据的 Hive 表结构设计Hive 里建表第一步要决定内部表还是外部表这个项目必须用外部表。外部表的含义是 Hive 只管理元数据数据文件放在 LOCATION 指定的 HDFS 目录里DROP TABLE 不会删掉原始文件。调试阶段反复重建表不会污染行情数据。建表语句按日线行情设计-- 股票日线行情外部表 CREATE EXTERNAL TABLE IF NOT EXISTS stock_db.stock_daily ( code STRING COMMENT 股票代码, name STRING COMMENT 股票名称, open DOUBLE COMMENT 开盘价, close DOUBLE COMMENT 收盘价, high DOUBLE COMMENT 最高价, low DOUBLE COMMENT 最低价, volume BIGINT COMMENT 成交量, amount DOUBLE COMMENT 成交额 ) PARTITIONED BY (trade_date STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY , STORED AS TEXTFILE LOCATION /warehouse/stock_db/stock_daily;字段类型的选择有讲究股票代码用 STRING 而不是 INT因为 000001 这种前导零在整型里会丢成交量可能超过 int 上限所以用 BIGINT价格用 DOUBLE 够用不需要 DECIMAL 的精确到分场景。PARTITIONED BY trade_date 让查询按日期裁剪分区避免全表扫描。数据文件按天放到对应分区目录后需要让 Hive 感知到新分区。手动方式适合少量日期hive -e ALTER TABLE stock_db.stock_daily ADD PARTITION (trade_date2024-01-05)如果目录很多用自动修复hive -e MSCK REPAIR TABLE stock_db.stock_dailyMSCK 会扫描 HDFS 上 LOCATION 下的子目录把符合分区格式的目录自动注册成分区。这里也能理解 Hive 的本质SQL 层只是外壳真正执行的是被翻译成 MapReduce 或 Tez 的分布式任务数据物理存在 HDFS 上。开发环境用 TEXTFILE 主要是方便用 hdfs dfs -cat 直接检查文件内容生产环境再换 Parquet 压缩格式。3.2 Flask 连接 HiveServer2 的代码与参数Python 连接 Hive 的常见方案是 pyhive它通过 thrift 协议和 HiveServer2 通信不依赖 JDK 里的 JDBC 驱动在 Flask 环境里集成最简单。requirements.txt 里一般已经包含 pyhive手动补齐依赖时这样安装pip install pyhive thrift thrift-sasl sasl连接封装通常放在 dbmodel/dbbase.py 中# dbmodel/dbbase.py 中的连接封装常见做法 from pyhive import hive def get_hive_conn(): # host 指向 HiveServer2 所在节点 return hive.Connection( host127.0.0.1, port10000, usernamehive, databasestock_db, authNONE ) def query_hive(sql, limit100): conn get_hive_conn() try: cursor conn.cursor() cursor.execute(sql) # 从 description 取列名避免硬编码字段 columns [desc[0] for desc in cursor.description] rows cursor.fetchmany(limit) return [dict(zip(columns, row)) for row in rows] finally: conn.close()query_hive 的核心逻辑是每次查询新建连接执行 SQL取列名后把每行转成字典最后在 finally 里关连接防止游标和连接泄漏。取列名用 cursor.description这比硬编码字段名更抗表结构变化。参数的对应关系部署时按这个表核对参数含义伪分布式推荐值hostHiveServer2 所在节点 IP127.0.0.1portHiveServer2 默认端口10000username连接用户启动 hiveserver2 的系统用户database默认数据库stock_dbauth认证方式NONE生产按 LDAP/KERBEROS 配置auth 参数最容易忽略。伪分布式部署一般用 NONE如果 Hive 配置了 LDAP 鉴权这里不填会一直报认证失败。而且 pyhive 的 Connection 还支持 connect_timeout 参数建议设置 10 到 30 秒避免 HiveServer2 卡住时 Flask 线程全部挂起。3.3 sql_tpl.py 里的 SQL 模板行转列与列转行的实际用处dbmodel 目录下的 sql_tpl.py 是这套代码里值得细看的设计。把 SQL 集中到模板文件里数据逻辑和业务逻辑彻底分开。常见实现是用 string.Template 做参数替换from string import Template STOCK_DAILY_BY_RANGE SELECT code, trade_date, close FROM stock_db.stock_daily WHERE code ${code} AND trade_date ${start} AND trade_date ${end} ORDER BY trade_date def render_stock_range(code, start, end): # substitute 会检查模板变量是否齐全 return Template(STOCK_DAILY_BY_RANGE).substitute( codecode, startstart, endend )Template 的 substitute 遇到缺失键会直接抛 KeyError这在开发期能尽早暴露参数问题比在 Hive 执行时报语法错误好排查。渲染出来的 SQL 最终交给 query_hive 执行。Hive SQL 里有两个高频操作经常被问行转列和列转行。行转列的场景是前端需要一次拿到多只股票字段用 collect_list 加 concat_ws-- 将同一交易日的多只股票聚合成一行 SELECT trade_date, concat_ws(|, collect_list(code)) AS code_list, concat_ws(|, collect_list(close)) AS close_list FROM stock_db.stock_daily WHERE trade_date 2024-01-05 GROUP BY trade_date;concat_ws 负责把数组拼成竖线分隔字符串collect_list 负责把多行变成数组前端拿到后用 split 还原比传一堆 JSON 数组更省流量。如果反向操作把竖线字符串拆回多行用 lateral view explode() 处理。这类 SQL 直接沉淀在 sql_tpl.py 里下次做报表时可以复用。3.4 查询接口的防御参数校验与 LIMITFlask 接口面向浏览器参数不可信。常见做法是进入 Hive 前做格式强校验SQL 层强制 LIMITdata_bp.route(/stock/list) def stock_list(): code request.args.get(code, ).strip() date request.args.get(date, ).strip() # 参数强校验控制注入风险 if code and (len(code) ! 6 or not code.isdigit()): return fail(股票代码应为6位数字) if date and len(date) ! 10: return fail(日期格式应为yyyy-MM-dd) sql SELECT * FROM stock_db.stock_daily WHERE 11 if code: sql f AND code {code} if date: sql f AND trade_date {date} sql ORDER BY trade_date DESC LIMIT 100 return success(query_hive(sql))这段代码直接拼 SQL 的前提是 code 和 date 都经过了严格格式校验code 被限制为 6 位数字date 被限制为 10 位长度的字符串注入风险基本被挡在外面。海量股票数据场景下不带 LIMIT 的 Hive 查询会把几十万行结果拉回 Flask 进程一次就能打爆内存。固定 LIMIT 100 牺牲一点功能保住接口稳定性。提示如果后续做参数化查询不要指望 Hive 的 PreparedStatement 像 MySQL 一样方便。更稳妥的方向是在应用层做白名单校验再传入 sql_tpl.py 的模板。4. 部署链路从 Hadoop 伪分布式到 Flask 上线压缩包里带的部署说明文档.md 和 Flask系统部署文档.md 已经把步骤拆开了。但文档是死的环境是活的先建立版本意识和启动顺序概念再照文档操作会顺很多。部署时最忌讳跳步Hadoop 没起来就启动 HiveServer2Flask 连不上就疯狂重启最后所有日志都指向同一个根因。4.1 环境核对清单Hadoop 生态最怕版本组合不一致。常见的课设环境里JDK 1.8 Hadoop 2.10.x Hive 2.3.x 的伪分布式组合出现频率较高Hadoop 3.x 与 Hive 3.x 也常见但部署时要注意 Hive 对 Hadoop 版本有最低要求。动手前先执行三条命令确认现状java -version hadoop version hive --version版本输出彼此兼容后再启动服务。伪分布式下 Hadoop 的启动入口在安装目录的 sbin 下start-dfs.sh start-yarn.sh jpsjps 输出里必须看到 NameNode、DataNode、SecondaryNameNode、ResourceManager、NodeManager。少任何一个都要先回看日志否则后面的 Hive 查询一定跑不通。Hive 首次使用前要初始化元数据库默认 Derby 模式直接执行schematool -initSchema -dbType derby nohup hiveserver2 /data/logs/hiveserver2.log 21 Derby 单连接的特性意味着 beeline 和 Flask 应用抢同一个连接时可能互相锁住。如果演示现场要同时开 beeline 和 Flask可以把 metastore 切到 MySQLHive 的连接配置写在 hive-site.xml 里改 connectionURL、driver 和 username 三处即可。4.2 Python 依赖与旧字节码清理把 Flask-Hive-master 放到 Linux 目录下先装依赖再清理字节码缓存cd Flask-Hive-master pip install -r requirements.txt find . -name *.pyc -deleterequirements.txt 里包含 Flask、pyhive、thrift 相关库。之所以要删 .pyc压缩包里带着上一台机器的字节码Python 版本不一致时导入会直接报语法错误或二进制不匹配删掉让解释器重新生成最省事。script 目录里的 runcleanmongocache.py 是清理缓存用的脚本如果部署时不启用缓存模块可以直接忽略。4.3 服务启动顺序与验证启动顺序有强依赖HiveServer2 依赖 HDFS 和 YARNFlask 依赖 HiveServer2。按下面的顺序执行不要颠倒start-dfs.sh start-yarn.sh nohup hiveserver2 /data/logs/hiveserver2.log 21 nohup python bootstrap.py --port 8080 /data/logs/flask.log 21 hiveserver2 启动需要几十秒别急着起 Flask。验证端口和接口ss -lntp | grep -E 8080|10000 curl http://127.0.0.1:8080/api/stock/list?code000001date2024-01-05curl 能返回 JSON 说明整条链路已经打通。如果返回 500先看 flask.log再顺藤摸瓜查 hiveserver2.log。日志里没有报错但接口超时大概率是 HiveServer2 还在初始化等几秒再 curl 一次不要反复重启进程。4.4 高频排错表部署阶段最常见的几个问题列成表对照处理现象原因处理方式pyhive 连接 10000 端口超时HiveServer2 没起来或日志里报错退出看 hiveserver2.log等进程稳定后再连报错 No module named thrift_saslthrift-sasl 未安装pip install thrift-sasl报错 Required field client_protocol is uninitializedthrift 与 pyhive 版本不匹配常见兼容组合是 pyhive 0.6.1 thrift 0.13.0 thrift-sasl 0.4.3查询报 Table not found没切库或分区未注册beeline 里执行 USE stock_db再 MSCK REPAIR TABLE大查询报 Java heap space执行内存不够调大 hive.heapsize降低 fetchmany 的 limit 值thrift 版本问题最隐蔽因为 pip 默认装的高版本 thrift 会和 pyhive 的旧接口冲突。遇到 protocol 报错先把三个库的版本往上面那个组合对齐八成能解决。表找不到的问题则多半是分区目录建了但没执行 MSCK不是表真不存在。5. 再往深走把课设代码改造成能长期维护的 Hive 分析平台5.1 给 Hive 查询加一层结果缓存Flask 每次请求都打到 HiveServer2数据量一大接口就明显变慢。常见做法是在中间加一层 MongoDB 缓存这也解释了项目里为什么会有 runcleanmongocache.py。按 code 加 trade_date 作为 key首次查询把 Hive 结果写进缓存五分钟内直接读缓存返回。Hive 里的行情数据不是秒级更新的这个策略不会带来数据不一致风险但演示流畅度会明显提升。清理脚本的存在说明作者已经在生产思路上做过考虑部署时可以保留这层。5.2 用 Hive 聚合 SQL 给 Flask 减负Flask 侧尽量避免循环查询 Hive把计算下推给 Hive 才是分布式该有的用法。比如周线聚合用 next_day 和 date_sub 组合就能在一条 SQL 里完成-- 周线聚合按周起始日分组计算最高、最低、总量 SELECT code, date_sub(next_day(trade_date, MONDAY), 7) AS week_start, max(close) AS week_high, min(low) AS week_low, sum(volume) AS week_volume FROM stock_db.stock_daily WHERE trade_date 2024-01-01 GROUP BY code, date_sub(next_day(trade_date, MONDAY), 7);next_day(trade_date, MONDAY) 取到本周周一date_sub 再退 7 天得到本周起始日这样按周一分组才能把一周数据归到同一天。计算压力全在 Hive 侧Flask 拿到的是已经聚合好的结果响应速度会比明细查询快很多。5.3 改表结构与改名时的元数据同步Hive 修改表名的 SQL 本身不复杂ALTER TABLE stock_db.stock_daily RENAME TO stock_db.stock_daily_bak;但外部表重命名后HDFS 上的 LOCATION 不会跟随改名。如果目录也手动改了名字必须用 SET LOCATION 同步元数据ALTER TABLE stock_db.stock_daily_bak SET LOCATION /warehouse/stock_db/stock_daily_bak;改完表名记得同步修改 sql_tpl.py 里的模板否则 Flask 查出来的还是旧表名路径会一直报 Table not found。新增字段的场景用 ALTER TABLE ... ADD COLUMNS但旧分区里的数据不会自动回填需要重写对应分区。最后用一条最简单的 SELECT 验证字段和分区都对上再回到前端页面刷新这时候看到的才是真正走通全链路的数据。本文还有配套的精品资源点击获取
返回列表