☰
Spark地铁客流分析全链路实战:HBase+Logstash+Spark SQL
2026/10/3 5:30:08 网站建设 项目流程

简介:本资源是一份面向计算机类本科生的毕业设计实战项目,聚焦城市地铁运营中的客流统计、趋势预测与调度优化问题,以Apache Spark为核心构建端到端大数据分析系统。项目完整覆盖需求分析、Spark分布式计算实现(含Scala/Java双语言代码)、MySQL/HBase数据库设计、ETL流程配置及可视化结果展示,适用于大数据课程设计、毕设选题与分布式系统实践学习。压缩包共194个文件,包含28个Java与17个Scala核心业务代码、17个XML配置及YAML/Properties集群参数文件、88张界面与架构图(PNG/SVG)、4个SQL建表与测试脚本、3份Git规范文档及CSV实测客流数据,整体大小42.6MB,结构清晰、模块可拆解。目前已有259人下载学习,提供开箱即用的Graduation Design主目录,内含完整报告、设计文档、可运行Spark Job及Logstash/Nginx等配套服务配置,助读者快速掌握大数据项目落地全流程。

1. 这不是又一个 WordCount 演示:它用 Spark 处理真实地铁刷卡记录,跑通从 HBase 写入、Logstash 接入、到 Spark SQL 聚合预测的全链路

你手头这份基于spark的地铁大数据客流分析系统.zip,不是教科书里那个“本地单机跑 10 行 CSV 的 Spark 入门 demo”。它是一套在真实毕业设计场景下被反复调试、能扛住日均百万级进出站刷卡数据的轻量级生产级流水线——压缩包里藏着.editorconfig和三个.gitignore,说明作者真写过代码、提交过 Git、踩过 IDE 编码风格冲突的坑;szmc.net-metro.csv是某市地铁公司脱敏后的实测客流原始数据(含时间戳、站点 ID、进出方向、卡类型);hbase.command不是空文件,而是带-n参数的create+put批量导入脚本;logstash-nginx.config明确指向 Nginx 日志采集路径,说明数据源不止 CSV,还有 Web 端埋点;而search.http和szt-api.http两个 Postman 风格的 HTTP 请求文件,直接暴露了后端 API 的/v1/flow/trend和/v1/station/peak接口契约。它解决的不是“怎么装 Spark”,而是“怎么让 Spark 在没 YARN、没 Kubernetes 的实验室服务器上,把 HBase 里的 2TB 历史刷卡记录,按小时粒度聚合出换乘热力图,并支持前端实时拖拽查询”。适合正在写毕设、但被导师一句“要体现工程能力”卡在答辩前两周的计算机本科生,也适合想快速复现一个可演示、可改参数、不依赖云厂商的 Spark 实战案例的转行者。


2. 从 CSV 到 HBase:为什么选 HBase 而不是 MySQL?三步完成地铁刷卡数据建模与批量写入

2.1 地铁数据天然适配 HBase 的三大硬约束:稀疏性、时序性、高写低读

地铁刷卡数据有三个致命特征:第一,稀疏性——每天 24 小时 × 300+ 站点 × 每分钟 100+ 条记录,但单个乘客一天只刷 2~4 次卡,99% 的 (时间, 站点, 卡号) 组合为空;第二,强时序性——所有分析必须带时间窗口(如“早高峰 7:00–9:00 进站量”),且新数据持续写入,旧数据极少更新;第三,写远大于读——每秒写入 500+ 条刷卡记录,但查询通常是按天/周聚合,或按站点查历史趋势。MySQL 的 B+ 树索引在这种场景下会严重碎片化,InnoDB Buffer Pool 频繁刷脏页,而 HBase 的 LSM-Tree 天然为写优化,RegionServer 自动按时间戳切分 Region,配合 RowKey 设计(如stationId_timestamp_cardHash)可实现毫秒级范围扫描。这不是“为了用而用”,是数据模型倒逼存储选型——你在hbase.command里看到的create 'metro_flow', {NAME => 'cf', TTL => 2592000}(TTL=30 天),就是为应对地铁数据生命周期管理做的硬编码决策。

2.2 解析szmc.net-metro.csv:字段含义、清洗逻辑与 RowKey 设计原理

先看原始数据结构(取前 5 行):

card_id,station_id,in_out,timestamp,device_id,card_type 1000000001,101,in,2023-08-01 07:15:22,DEV-001,ordinary 1000000001,102,out,2023-08-01 07:28:11,DEV-002,ordinary 1000000002,205,in,2023-08-01 07:32:45,DEV-015,student 1000000003,308,in,2023-08-01 07:41:03,DEV-022,elderly 1000000001,101,in,2023-08-01 18:05:17,DEV-001,ordinary

关键清洗动作(已在src/main/resources/clean_metro.py中固化):

  • 时间标准化:timestamp字段统一转为yyyy-MM-dd HH:mm:ss,并提取hour_of_day(0~23)、is_weekend(布尔)、peak_flag(早/晚高峰标记);
  • 站点归一化:station_id映射到标准站名表(station_map.json),解决同一站点多设备 ID(如DEV-001,DEV-002)导致的重复计数;
  • 进出方向校验:剔除in_out非in或out的脏数据(实际占比约 0.3%,多为设备通信错误);
  • RowKey 设计:采用stationId#yyyyMMdd#HHmmss#cardHash格式(例:101#20230801#071522#e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855),其中#为分隔符,cardHash是card_id的 SHA256 前 16 位,既保证唯一性又避免明文泄露隐私。此设计使scan 'metro_flow', {STARTROW => '101#20230801', STOPROW => '101#20230802'}可精准获取某站全天数据,无需全表扫描。

2.3 执行hbase.command:三步完成建表、导入、验证(附参数详解)

进入 HBase Shell 后,逐行执行hbase.command内容(已去注释):

# 1. 创建表:启用压缩、设置 TTL、预分区(按站点 ID 哈希) create 'metro_flow', {NAME => 'cf', COMPRESSION => 'SNAPPY', TTL => 2592000}, {NUMREGIONS => 16, SPLITALGO => 'HexStringSplit'} # 2. 批量导入 CSV(使用 HBase 自带的 ImportTsv 工具) hbase org.apache.hadoop.hbase.mapreduce.ImportTsv \ -Dimporttsv.columns="HBASE_ROW_KEY cf:station_id cf:in_out cf:timestamp cf:device_id cf:card_type" \ -Dimporttsv.skip.bad.lines=true \ -Dimporttsv.separator=',' \ metro_flow /data/szmc.net-metro.csv # 3. 验证数据量(注意:count 操作在大表上极慢,此处用 get_region_info 替代) echo "scan 'metro_flow', {LIMIT => 5}" | hbase shell

提示:ImportTsv的-Dimporttsv.columns参数必须严格对应 CSV 列顺序,且HBASE_ROW_KEY必须是第一列——这意味着你需先用awk或 Python 脚本将szmc.net-metro.csv转为rowkey,station_id,in_out,...格式。原包中scripts/prepare_hbase_input.py已实现此转换,执行python scripts/prepare_hbase_input.py szmc.net-metro.csv > metro_hbase_input.csv即可。

2.4 验证写入正确性:用get和scan查两条典型记录

# 查一条进站记录(站 101,早高峰) hbase(main):001:0> get 'metro_flow', '101#20230801#071522#e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855' COLUMN CELL cf:in_out timestamp=1700000000000, value=in cf:station_id timestamp=1700000000000, value=101 cf:timestamp timestamp=1700000000000, value=2023-08-01 07:15:22 # 扫描某站某小时全部记录(验证时间范围) hbase(main):002:0> scan 'metro_flow', {STARTROW => '101#20230801#07', STOPROW => '101#20230801#08', LIMIT => 3}

若get返回COLUMN为空,说明 RowKey 生成逻辑与ImportTsv的HBASE_ROW_KEY列不匹配;若scan返回 0 条,检查STARTROW/STOPROW是否符合字典序(HBase 的 RowKey 是字节序比较,101#20230801#07<101#20230801#079999,但101#20230801#07<101#20230801#070000成立)。


3. Logstash 接入 Nginx 日志:为什么不用 Flume?如何把 Web 埋点日志喂给 Spark Streaming

3.1 选型依据:Logstash 对 Nginx 日志的解析能力远超 Flume 的默认 Source

项目中logstash-nginx.config存在,说明数据源不止刷卡 CSV,还包括 Web 端用户行为日志(如“查询某站今日客流”、“导出周报 PDF”)。Flume 的ExecSource或SpoolingDirSource无法原生解析 Nginx 的combined日志格式(含 IP、URL、状态码、响应时间),而 Logstash 的grok插件可一行匹配:

%{IP:client} - %{USER:ident} \[%{HTTPDATE:timestamp}\] "%{WORD:method} %{URIPATHPARAM:request} HTTP/%{NUMBER:httpversion}" %{NUMBER:response} %{NUMBER:bytes} "%{DATA:referrer}" "%{DATA:agent}"

logstash-nginx.config中的关键配置:

input { file { path => "/var/log/nginx/access.log" start_position => "beginning" sincedb_path => "/dev/null" # 避免重启后重复消费 } } filter { grok { match => { "message" => '%{IP:client} - %{USER:ident} \[%{HTTPDATE:timestamp}\] "%{WORD:method} %{URIPATHPARAM:request} HTTP/%{NUMBER:httpversion}" %{NUMBER:response} %{NUMBER:bytes} "%{DATA:referrer}" "%{DATA:agent}"' } } date { match => [ "timestamp", "dd/MMM/yyyy:HH:mm:ss Z" ] target => "@timestamp" } mutate { add_field => { "event_type" => "web_access" } remove_field => ["message", "timestamp"] } } output { elasticsearch { hosts => ["http://localhost:9200"] index => "metro-web-log-%{+YYYY.MM.dd}" } }

注意:sincedb_path => "/dev/null"是为测试环境简化设计,生产环境应指向持久化路径;@timestamp字段被重写为 Logstash 解析出的时间,而非日志写入时间,确保时序分析准确。

3.2 Spark Streaming 如何消费 Elasticsearch 数据?用es-hadoop连接器直读

src/main/scala/streaming/WebLogAnalyzer.scala中核心代码:

import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ val spark = SparkSession.builder() .appName("WebLogStreaming") .config("spark.es.nodes", "localhost") // ES 地址 .config("spark.es.port", "9200") .config("spark.es.index.read", "metro-web-log-*") // 通配符读取多日索引 .config("spark.es.query", """{"query":{"range":{"@timestamp":{"gte":"now-1h/h","lt":"now/h"}}}}""") // 仅读最近 1 小时 .getOrCreate() // 每 30 秒触发一次微批处理 val webLogDF = spark.read .format("es") .option("es.read.field.as.array.include", "false") .load() // 提取 URL 中的站点 ID(如 /api/v1/station/101/flow) val stationAccessDF = webLogDF .filter($"request".contains("/station/")) .withColumn("station_id", regexp_extract($"request", "/station/(\\d+)/", 1)) .filter($"station_id" =!= "") .groupBy("station_id") .agg(count("*").alias("web_query_count")) stationAccessDF.write .mode("append") .format("jdbc") .option("url", "jdbc:mysql://localhost:3306/metro_db") .option("dbtable", "web_station_query") .option("user", "root") .option("password", "123456") .save()

此代码将 Web 查询行为(如查站 101 客流)与刷卡数据打通,后续可在 Spark SQL 中做关联分析:“某站 Web 查询量激增是否预示该站未来 30 分钟进站量上升?”——这正是毕设答辩时最亮眼的业务洞察点。

3.3 避坑:Logstash 与 Spark Streaming 的时间对齐陷阱

现象原因解决
Spark Streaming 统计的“当前小时 Web 查询量”比 Kibana 图表少 30%Logstash 的@timestamp是日志解析时间,而 Nginx 写日志有延迟(平均 2~5 秒),Spark 读取时now-1h/h窗口漏掉了最后一批日志在 Logstashdatefilter 中加timezone => "Asia/Shanghai",并在 Spark SQL 查询中将@timestamp转为东八区时间:from_utc_timestamp($"@timestamp", "Asia/Shanghai")
Elasticsearch 索引metro-web-log-2023.08.01写满后,Spark 读取时报IndexNotFoundExceptionspark.es.index.read的通配符*默认只匹配存在索引,若某日无日志则索引不存在,通配符失效改用spark.es.resource指定具体索引名,或在 Spark 作业启动前用curl -X GET "localhost:9200/_cat/indices?h=index&s=index"动态获取当日索引列表
Spark 读取 ES 时 OOM(堆内存溢出)ES 返回的_source包含完整日志行(含agent字段的长 User-Agent 字符串),单条记录超 1KB,微批 1000 条即 1MB,Driver 端聚合时内存爆炸在spark.read.format("es")后加.select("client", "request", "response", "@timestamp")显式指定字段,丢弃agent、referrer等非必要字段

4. Spark SQL 核心分析:四类必做客流指标的 SQL 实现与性能调优技巧

4.1 四类指标定义与业务价值:从基础统计到预测前置

指标类型SQL 示例(简化)业务价值毕设得分点
实时进站 TOP10SELECT station_id, COUNT(*) AS in_count FROM flow WHERE in_out='in' AND event_time >= now() - INTERVAL 15 MINUTES GROUP BY station_id ORDER BY in_count DESC LIMIT 10运营调度中心大屏实时展示,发现突发大客流站点展示 Spark Streaming 实时能力
换乘热力图(站间 OD)SELECT a.station_id AS from, b.station_id AS to, COUNT(*) AS transfer_count FROM flow a JOIN flow b ON a.card_id = b.card_id WHERE a.in_out='out' AND b.in_out='in' AND a.timestamp < b.timestamp AND b.timestamp - a.timestamp < INTERVAL 30 MINUTES GROUP BY from, to识别高频换乘组合(如 101→205),指导列车班次加密体现 Join 优化与时间窗口控制
早高峰预测(ARIMA 基线)SELECT station_id, avg(in_count) AS baseline, stddev(in_count) AS std FROM (SELECT station_id, date(event_time) AS dt, COUNT(*) AS in_count FROM flow WHERE in_out='in' AND hour(event_time) BETWEEN 7 AND 9 GROUP BY station_id, dt) GROUP BY station_id为机器学习模块提供统计基线,对比预测偏差展示数据预处理与特征工程
异常波动告警SELECT station_id, dt, in_count, baseline, CASE WHEN ABS(in_count - baseline) > 3 * std THEN 'ALERT' ELSE 'OK' END FROM (...)当某站进站量突增 3 倍标准差,自动推送企业微信告警体现实战闭环能力

4.2 关键 SQL 性能调优:从EXPLAIN到broadcast join的落地步骤

以“换乘热力图”为例,原始 SQL 在 1 亿条记录上运行超 15 分钟。优化路径如下:

Step 1:用EXPLAIN EXTENDED定位瓶颈

EXPLAIN EXTENDED SELECT a.station_id AS from, b.station_id AS to, COUNT(*) FROM flow a JOIN flow b ON a.card_id = b.card_id WHERE a.in_out='out' AND b.in_out='in' AND a.timestamp < b.timestamp GROUP BY from, to

输出中WholeStageCodegen下出现SortMergeJoin,说明 Spark 选择了代价最高的 Shuffle Join。

Step 2:改用 Broadcast Join(因card_id分布倾斜,但station_id维度表仅 300 行)

-- 先缓存维度表(站点信息) val stationDF = spark.read.jdbc("jdbc:mysql://...", "stations", new java.util.Properties()) stationDF.cache() spark.catalog.clearCache() -- 在 SQL 中强制广播 spark.sql( s""" |SELECT /*+ BROADCAST(stationA), BROADCAST(stationB) */ | stationA.name AS from_name, stationB.name AS to_name, COUNT(*) AS count |FROM flow a |JOIN flow b ON a.card_id = b.card_id |JOIN stationDF stationA ON a.station_id = stationA.id |JOIN stationDF stationB ON b.station_id = stationB.id |WHERE a.in_out='out' AND b.in_out='in' | AND b.timestamp > a.timestamp | AND b.timestamp < a.timestamp + INTERVAL 30 MINUTES |GROUP BY from_name, to_name |""".stripMargin)

Step 3:调整 Spark SQL 参数(写入spark-defaults.conf)

spark.sql.adaptive.enabled=true # 启用自适应查询执行(AQE) spark.sql.adaptive.coalescePartitions.enabled=true # 自动合并小分区 spark.sql.autoBroadcastJoinThreshold=50000000 # 广播表阈值调至 50MB(原 10MB) spark.sql.files.maxPartitionBytes=134217728 # 单分区最大 128MB(避免过多小文件)

血泪经验:autoBroadcastJoinThreshold必须大于维度表大小(stationDF.count()× 每行字节数),否则 Broadcast 不生效;用spark.sql("CACHE TABLE stations")比stationDF.cache()更可靠,因后者可能被 GC 清理。

4.3 输出结果到 MySQL:为什么不用 JDBC 直连而用insertInto?

src/main/scala/analysis/ODAnalyzer.scala中写库代码:

odResultDF.write .mode("overwrite") .option("truncate", "true") // 关键!避免 insert 导致主键冲突 .option("batchSize", "10000") // 批量插入,减少网络往返 .jdbc("jdbc:mysql://localhost:3306/metro_db", "od_heatmap", props)

但更推荐用 Spark SQL 的insertInto:

odResultDF.createOrReplaceTempView("od_temp") spark.sql("INSERT OVERWRITE TABLE od_heatmap SELECT * FROM od_temp")

原因:insertInto由 Spark Catalyst 优化器统一规划,可复用已缓存的中间结果;而write.jdbc每次都重新计算 DataFrame,且truncate=true会锁表,高并发时易阻塞。毕设演示时,用insertInto可保证“点击分析按钮 → 3 秒内刷新热力图”的流畅体验。


5. 毕设答辩避坑指南:导师最常问的 5 个问题及满分回答话术

5.1 “Spark 和 Flink 你为什么选 Spark?Flink 不是更适合实时?”

满分答:
“Flink 确实在纯流式场景(如毫秒级风控)有优势,但本项目核心是‘准实时’——客流分析需按小时/天聚合,且大量离线历史数据(HBase 中 2TB 历史记录)需与实时流 Join。Spark Structured Streaming 的 micro-batch 模型,能复用同一套 DataFrame API 处理批流一体,SQL 引擎成熟度高,社区生态(如 MLlib 预测模块)更完善。我们实测:Spark Streaming 处理 10 万 QPS 的刷卡流,端到端延迟稳定在 2.3 秒(P95),完全满足地铁运营‘分钟级响应’要求。”

5.2 “HBase 为什么不直接用 Phoenix 提供 SQL 接口?”

满分答:
“Phoenix 虽提供 SQL,但其二级索引在高并发写入时性能衰减明显,且不支持复杂 Join(如 OD 分析需跨行关联)。本项目选择 HBase 原生 API + Spark RDD/DataFrame,通过精心设计 RowKey(stationId#date#time#hash)和预分区,使 90% 的查询走scan而非全表扫,实测 1 亿行数据scan某站某日耗时 1.2 秒。Phoenix 的 SQL 抽象层反而增加了不可控的开销。”

5.3 “Logstash 采集 Nginx 日志,如果日志量暴增怎么办?”

满分答:
“我们在logstash-nginx.config中已预留弹性:input 使用file插件而非beats,便于横向扩展多个 Logstash 实例;filter 阶段禁用geoip等重 CPU 插件;output 采用elasticsearch的bulk模式(默认 1000 条/批)。压力测试表明:单实例 Logstash 可稳定处理 5000 EPS(Events Per Second),超阈值时只需增加实例数并修改output的hosts列表,无需改代码——这正是微服务架构的伸缩性体现。”

5.4 “你们的客流预测用的是什么算法?准确率多少?”

满分答:
“预测模块分两层:第一层用 Spark MLlib 的LinearRegression做基线(输入:历史 7 天同小时进站量、天气、节假日标志),RMSE=123;第二层用GBTRegressor提升精度(加入 OD 热力图特征),RMSE=89。但更重要的是业务验证:我们将预测结果与 8 月 15 日实际客流对比,早高峰(7–9 点)预测误差 < 8%,完全可用于调度参考。毕设报告第 4.3 节有详细实验表格和残差图。”

5.5 “整个系统怎么部署?需要多少台服务器?”

满分答:
“我们提供三种部署模式:① 单机开发版(1 台 16G 内存服务器,HBase 伪分布式 + Spark Local 模式),适合毕设演示;② 小集群生产版(3 节点:1 Master + 2 Worker,HBase 分布式 + Spark Standalone),支撑日均 500 万刷卡记录;③ 云原生版(HBase on Kubernetes + Spark on K8s),已通过 Helm Chart 封装。部署文档docs/deploy.md中有每种模式的docker-compose.yml和资源配置清单,答辩时可现场演示单机版一键启动。”


6. 最后一道防线:用search.http和szt-api.http验证 API 正确性,以及我每次打包前必做的三件事

6.1 用 Postman 风格.http文件做接口冒烟测试

search.http文件内容(已脱敏):

### 查询某站小时客流趋势 GET http://localhost:8080/api/v1/station/101/flow?start=2023-08-01&end=2023-08-02 Accept: application/json ### 查询全网 OD 换乘矩阵(TOP 50) GET http://localhost:8080/api/v1/od/top50 Accept: application/json ### 触发预测任务(异步) POST http://localhost:8080/api/v1/predict/trigger Content-Type: application/json { "station_id": "101", "horizon_hours": 3 }

执行命令(需安装httpie):

# 测试前确保服务已启动 http --print=h https://localhost:8080/api/v1/station/101/flow start==2023-08-01 end==2023-08-02 # 验证返回 JSON 结构(非空且含 expected_keys) http GET http://localhost:8080/api/v1/station/101/flow start==2023-08-01 end==2023-08-02 | jq '.data[0] | has("hour") and has("in_count") and has("out_count")' # 应输出 true

提示:.http文件中的###是 Httpie 的请求分隔符,每个###块是一个独立请求;==是 Httpie 的 URL 参数语法,等价于?start=...&end=...。

6.2 我每次打包交付前必做的三件事(血泪换来的后悔药)

  1. 清空 HBase 表并重放hbase.command

    echo "disable 'metro_flow'; drop 'metro_flow'" | hbase shell # 重新执行 hbase.command # 验证:hbase(main):001:0> count 'metro_flow', INTERVAL => 60000

    为什么:避免残留测试数据污染演示效果,INTERVAL => 60000让 count 采样而非全扫,10 秒内出结果。

  2. 用spark-sql直连 Hive Metastore 查元数据

    spark-sql --master local[*] -e "SHOW DATABASES; USE metro_db; SHOW TABLES; DESCRIBE flow;"

    为什么:确认 Spark SQL 能识别 Hive 表结构,避免答辩时SELECT * FROM flow报Table not found——这是因hive-site.xml未正确加载导致的玄学错误。

  3. 用jps -l检查进程树,杀掉所有残留 Java 进程

    jps -l | grep -E "(HMaster|HRegionServer|SparkSubmit)" | awk '{print $1}' | xargs kill -9 2>/dev/null

    为什么:HBase 和 Spark 的守护进程常驻后台,若上次运行异常退出,端口(如 HBase 的 16000)被占,新启动必失败。jps -l比ps aux | grep更精准,只杀 Java 进程,不误伤系统服务。

从那以后我每次打包交付前,都强制走一遍这三步——不是怕导师提问,是怕自己在答辩现场点开浏览器,看到Connection refused时冷汗浸透衬衫。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询