Hadoop生态下的城市公共交通时空大数据分析实战 简介一份基于Hadoop架构的城市公共交通大数据时空分析学士学位论文面向计算机科学与技术、软件工程等专业本专科毕业生也适合对大数据处理分析感兴趣的学习者。论文从Hadoop核心组件HDFS与MapReduce入手系统讲解分布式存储与并行计算原理并结合公交GPS数据、刷卡记录等场景完整覆盖数据清洗、预处理、时空聚类、热点识别及流量预测等关键环节。案例部分通过K-means聚类识别公交热门站点、利用时间序列分析预测交通流量并探讨了优化调度与缓解拥堵的应用效果。包内仅有1个docx文档压缩包大小29KB内容包含摘要、相关技术、数据搜集与处理、时空数据存储与处理、应用案例等章节结构清晰。已有269人学习下载对于需要完成类似课题或理解Hadoop落地应用的读者能提供从理论到案例的完整参考。1. 城市公共交通大数据分析为什么第一站选 Hadoop城市公共交通大数据的第一层形态是每天千万条刷卡和 GPS 记录每一条都带着线路、站点、时间和经纬度。真正值得分析的永远是聚合后的时空现象早高峰哪个站点积压某条线全程耗时怎么波动跨区 OD 从哪里到哪里。这种分析本质上是一个先分桶、再关联、最后聚合的批量计算问题Hadoop 正好覆盖这三步HDFS 存原始文件MapReduce/Spark 完成关联Hive 提供 SQL 入口。哪怕集群只有几台机器只要单机装不下全量明细这套“先存后算、分区裁剪、列式扫描”的方案就是比较好的起步选择。下面从数据落盘开始把每层该做什么尽量讲清楚。2. 数据上 Hadoop从原始刷卡记录到可分析的时空分层2.1 Hadoop 开发环境搭建先跑通伪分布式再上集群在搭多台机器的 Hadoop 集群之前我一般会先用伪分布式把数据链路跑通。不是图省事而是伪分布式能暴露 HDFS 权限、Java 版本、临时目录清理这类基础问题这些问题在多节点上排查成本会成倍增加。伪分布式的配置只有三份文件core-site.xml管默认文件系统hdfs-site.xml管数据副本数和目录yarn-site.xml管资源调度。property namefs.defaultFS/name valuehdfs://localhost:9000/value /property property namehadoop.tmp.dir/name value/data/hadoop-tmp/value /propertyfs.defaultFS必须是 hdfs 协议否则 Hadoop 默认读写本地文件Hive 建外表时就会找不到 HDFS 路径。hadoop.tmp.dir里放着 NameNode 的 edits 和 DataNode 的数据块很多教程直接写/tmp机器重启后系统一清理整个集群起不来。伪分布式把dfs.replication设成 1 就够避免每个块都写三份浪费磁盘。启动命令也要分开先start-dfs.sh再start-yarn.sh。启动完用jps看进程至少需要 NameNode、DataNode、ResourceManager、NodeManager 四个。如果 DataNode 起不来优先检查/data/hadoop-tmp的属主和权限以及hdfs-site.xml里有没有重复配置dfs.namenode.name.dir。2.2 按 ODS、DWD、DWS 分层让时空字段始终可追溯Hive 表我习惯分成三层。ODS 层只做装载字段名和源文件保持一致DWD 层做清洗和时空字段标准化DWS 层做聚合。这样做的最大好处是当 DWS 结果和业务预期对不上时可以逐层对比而不是从头重新解析原始文件。CREATE EXTERNAL TABLE ods_card_record ( card_id STRING, bus_line STRING, bus_plate STRING, station STRING, lon DOUBLE, lat DOUBLE, op_time STRING ) PARTITIONED BY (dt STRING) ROW FORMAT DELIMITED FIELDS TERMINATED BY , STORED AS TEXTFILE LOCATION /data/ods/card_record;这张 ODS 表的关键点是PARTITIONED BY (dt STRING)查询时用WHERE dt2026-02-16直接走分区裁剪。如果原始文件按天放在 HDFS把文件挪到LOCATION对应的目录里Hive 就能自动识别分区。DWD 层我会把op_time从字符串转成TIMESTAMP再补一列grid_id。这里要小心转换格式必须和源数据完全一致。公交设备常见格式是yyyyMMddHHmmssHive 里要写unix_timestamp(op_time, yyyyMMddHHmmss)再用from_unixtime转成标准时间。不要用cast(op_time as timestamp)Hive 对字符串时间格式要求严格遇到斜杠日期很容易返回 NULL。2.3 网格化编码把经纬度变成可分组的空间键空间上公交数据最常用的不是裸经纬度而是网格 ID。同一站点周边会存在多个 GPS 点直接按经纬度 group by 会产生大量只有一两条记录的组按网格聚合又需要先把经纬度归到网格。我一般会做一个宽度 0.01 度约 1 公里的空间键CAST(floor((lon - 120.00) / 0.01) AS INT) AS grid_x, CAST(floor((lat - 30.00) / 0.01) AS INT) AS grid_y计算时要先减去城市左下角的基准经度和纬度上面用的是经度 120、纬度 30 附近的城市实际分析以你的数据范围为准。网格大小的选择会直接影响结果粒度网格边长实际约典型用途0.05 度5 公里城市早晚高峰潮汐0.01 度1 公里客流热力图、调度片区0.001 度100 米站点周边精准客流选好边长后grid_x * 10000 grid_y可以合成一列grid_id避免 SQL 里到处写两个字段。但聚合统计时最好还是保留grid_x、grid_y两个原始整数列方便定位异常网格。2.4 数据导入常见坑时区、乱码和脏坐标公交第三方导出的 CSV 常见 GBK 编码用file命令先检查file card_20260216.csv iconv -f GBK -t UTF-8 card_20260216.csv card_20260216_utf8.csv转换后再用hdfs dfs -put上传不要在写 Hive 时临时转码。时区问题比乱码隐蔽车载终端存的是“本地时间”Hive 服务器时区是 UTCfrom_unixtime出来会差 8 小时。处理办法是对时间字段做FROM_UTC_TIMESTAMP(op_time, Asia/Shanghai)。如果原始字符串已经带08:00Hive 解析时会自动转换这时不能再套一层否则会重复偏移。判断清楚源格式后再写清洗 SQL。3. 时空分析核心用 Hive SQL 完成轨迹清洗、OD 与热力3.1 漂移点过滤窗口函数是清理 GPS 轨迹的常用方式GPS 轨迹点里最头疼的是漂移点。公交车在桥下、隧道或高楼间车载 GPS 会短暂跳到几百米外这些点不清理画出来的轨迹会横穿城市。判断漂移的常见办法是看相邻两点之间的距离和时间差距离很大且时间间隔只有几秒基本都是漂移。用 Hive 的 LAG 窗口函数把前一个点搬到同一行WITH gps_seq AS ( SELECT bus_plate, lon, lat, record_time, LAG(lon) OVER (PARTITION BY bus_plate, trip_id ORDER BY record_time) AS prev_lon, LAG(lat) OVER (PARTITION BY bus_plate, trip_id ORDER BY record_time) AS prev_lat, LAG(record_time) OVER (PARTITION BY bus_plate, trip_id ORDER BY record_time) AS prev_time FROM dwd_gps_record WHERE dt 2026-02-16 ) SELECT * FROM gps_seq WHERE prev_lon IS NULL OR (ABS(lon - prev_lon) ABS(lat - prev_lat)) 0.01 OR unix_timestamp(record_time) - unix_timestamp(prev_time) 30;这里trip_id是“车次”的标识没有的话可以用bus_plate 发车时间拼出来。位移阈值 0.01 度大约 1 公里时间阈值 30 秒过滤太多说明阈值太小过滤太少再把时间阈值压到 20 秒。这个 SQL 的逻辑是保留首点、位移小于 1 公里的点、时间间隔超过 30 秒的点把短时间大位移的漂移点排除。参数建议值说明位移阈值0.01 度经度差绝对值加纬度差绝对值时间阈值30 秒超过该时间算正常停留trip_id无则用 bus_plate发车时间避免窗口跨车次窗口函数中的unix_timestamp要求参数时区已经校准否则差值计算仍是对的但整体时间偏移会导致你无法和站点班次对齐。3.2 OD 分析两个上车点之间该怎么定义一次出行OD 分析要回答乘客在哪一站上车、哪一站下车。对于通勤 OD最直接的方式是取同一张卡在同一天内的连续两次上车记录。下面用 LAG 模拟“上一站”WITH card_seq AS ( SELECT card_id, station, op_time, LAG(station) OVER (PARTITION BY card_id, dt ORDER BY op_time) AS origin_station, LAG(op_time) OVER (PARTITION BY card_id, dt ORDER BY op_time) AS origin_time FROM dwd_card_record WHERE dt 2026-02-16 ) SELECT origin_station, station AS dest_station, COUNT(*) AS od_cnt FROM card_seq WHERE origin_station IS NOT NULL AND (unix_timestamp(op_time) - unix_timestamp(origin_time)) 7200 GROUP BY origin_station, station;这个查询把两小时内的连续两次乘车算作一次 OD。为什么用连续两次而不是第一次和最后一次因为很多乘客一天会换乘多次取第一次和最后一次会把中间换乘跳过去。两小时窗口按城市规模调整大城市换乘窗口可能到 3 小时小城市 1.5 小时就够。调整窗口后重点观察 OD 表里是否存在“明显不可能步行到达”的换乘比如直线距离几十公里且时间窗口很短那基本是卡号跨日复用或设备号重复。3.3 站点时段聚合用 DWS 层承接大屏请求DWS 层是给大屏和报表用的。建表时我会把dt和hour都作为分区CREATE EXTERNAL TABLE dws_station_hourly_flow ( station_id STRING, direction STRING, flow_cnt BIGINT ) PARTITIONED BY (dt STRING, hour STRING) STORED AS PARQUET LOCATION /data/dws/station_hourly_flow;STORED AS PARQUET在这里很重要。DWS 层查询通常只读少数列列式存储让扫描更少加上 Snappy 压缩后文件体积比 TEXT 小一半以上。如果 Hive 版本和 Spark 的 Parquet 兼容有问题查询时可能报PendingRolledFile需要在 Spark 侧设置spark.sql.hive.convertMetastoreParquetfalse。写数时注意不能重复累加。一般用INSERT OVERWRITE覆盖当日分区INSERT OVERWRITE TABLE dws_station_hourly_flow PARTITION (dt2026-02-16, hour) SELECT station_id, direction, COUNT(*) AS flow_cnt FROM dwd_card_record WHERE dt 2026-02-16 GROUP BY station_id, direction, hour;动态分区前有两个参数要设置set hive.exec.dynamic.partitiontrue和set hive.exec.dynamic.partition.modenonstrict否则 Hive 默认只按最后一个字段分区SELECT 里的顺序一错就会报错。3.4 热力网格统计group by 后面跟着空间键全域热力网格的统计相对简单INSERT OVERWRITE TABLE dws_grid_flow PARTITION (dt2026-02-16) SELECT grid_x, grid_y, COUNT(*) AS flow_cnt FROM dwd_card_record WHERE dt 2026-02-16 AND lon 0 AND lat 0 GROUP BY grid_x, grid_y;这里多写了lon 0 AND lat 0是因为很多脏数据把经纬度写成 0 或负数。不加这个条件floor((lon - 120)/0.01)会产生一个巨大的负数聚合结果里就会多出一个“错误网格”。如果之前生成了grid_id也可以直接GROUP BY grid_id但大屏要画热力图时最好保留grid_x, grid_y并在外面重建网格中心点SELECT (grid_x 0.5) * 0.01 120.00 AS center_lon, (grid_y 0.5) * 0.01 30.00 AS center_lat FROM dws_grid_flow;中心点计算要加半格否则坐标落在网格左下角地图上看会有明显偏移。4. Hadoop 集群调优时空分析 SQL 跑不动时先改哪里4.1 执行引擎切换Hive 换 Spark 前要改的配置公交时空分析通常以天为粒度跑全量吞吐。Hive 默认的 MapReduce 引擎在窗口函数和多次 Join 时会写大量中间结果到磁盘走两轮 MapReduce 后速度明显下降。如果集群里已经装好 Spark我一般会切到 Spark 引擎set hive.execution.enginespark; set spark.masteryarn; set spark.driver.memory2g; set spark.executor.memory4g;这三行在 Hive CLI 或 Beeline 里执行即可。hive.execution.engine决定作业由哪个引擎执行spark.masteryarn表示让 YARN 分配 Executor而不是在 Hive 进程里起本地 Spark。后者会导致一个 Hive 客户端独占一个 SparkContext多人共用集群时很快把资源占满。切换后有人会报Java heap space大概率是spark.executor.memory给得比原来 MapReduce 的容器还小先看内存配置再去调 SQL。4.2 数据倾斜热点站点打散键的三种做法公共交通天然有热点早高峰的换乘站、晚高峰的商业中心客流能比周边站点高一个数量级。当GROUP BY station_id的某一个键占大头时Spark 的单个任务会等待最热分区完成。常见三种做法两阶段聚合先加盐再聚合最后去盐拆键过滤把最大的几个热点 station 单独拿出来走 mapjoin调大 reducer把 reducers 从默认值提到 300同时减小每个 reducer 的字节数。两阶段聚合最通用下面用 Spark SQL 表示-- 第一阶段按 station 加 10 个随机后缀 INSERT OVERWRITE TABLE tmp_station_skew SELECT concat(station_id, _, cast(rand()*10 as int)) AS rk, count(*) AS cnt FROM dwd_card_record WHERE dt 2026-02-16 GROUP BY concat(station_id, _, cast(rand()*10 as int)); -- 第二阶段去掉后缀聚合 SELECT split(rk, _)[0] AS station_id, sum(cnt) AS cnt FROM tmp_station_skew GROUP BY split(rk, _)[0];注意rand()*10 as int会把后缀控制在 09split(rk,_)[0]要求 station_id 本身不含下划线。如果站点编码里有下划线改用regexp_replace或固定长度截取。这个方案能解决热点倾斜但会多读一遍中间结果只对热点明显且数据量大的场景有意义。对应参数表参数默认意义调整建议mapreduce.job.reduces默认由数据量推断热点明显时固定到 300hive.exec.reducers.bytes.per.reducer1GB调小到 256MB 会增加 reducerspark.sql.shuffle.partitions200大表 group by 时提到 400注意mapreduce.job.reduces在 Spark 引擎下不生效调整时先确认当前引擎。4.3 小文件合并从写入端控制 DWS 层文件数时空分析任务每天跑一次如果每次写入都按小时分区HDFS 上会留下大量小文件。小文件多了NameNode 内存和后续查询的 RPC 都会升高。我一般从写入端控制INSERT OVERWRITE TABLE dws_station_hourly_flow PARTITION (dt2026-02-16, hour) SELECT station_id, direction, COUNT(*) AS flow_cnt FROM dwd_card_record WHERE dt 2026-02-16 GROUP BY station_id, direction, hour DISTRIBUTE BY hour;DISTRIBUTE BY hour让相同 hour 的数据尽量进同一个 reducer这样每个小时分区只生成 1 个文件而不是每个 reducer 都给每个分区各写一份。写入完成后检查目录hdfs dfs -du -h /data/dws/station_hourly_flow/dt2026-02-16如果看到某个小时目录下有几十个 10KB 的文件说明DISTRIBUTE BY没生效回去看写入 SQL 的分区键是否写完整。4.4 借助 YARN 日志定位问题调优不是盲改。任务卡住时先到 YARN 里看状态yarn application -list -appStates RUNNING yarn logs -applicationId application_xxx /tmp/app_xxx.log日志里最关心两类Shuffle Spill和 GC 耗时。如果 spill 出现很多次说明 map 输出的内存不够优先开压缩set mapreduce.map.output.compresstrue; set mapreduce.map.output.compress.codecorg.apache.hadoop.io.compress.SnappyCodec;GC 时间占比高把mapreduce.reduce.memory.mb从默认 1024 提到 2048但不要超过 YARN 容器内存上限。yarn logs里看不到具体 SQL 行号时回 Hive 用explain看执行计划重点检查Reduce Operator Tree里的 group by 是否按预期下了推。5. Hadoop 分析结果上大屏从 Hive 到 ECharts 的最小联动5.1 导出DWS 层到 MySQL 的三种路径大屏一般不会直接查 Hive常见做法是把 DWS 结果同步到 MySQL前端再通过接口读。如果集群里装 Sqoop可以用sqoop export但 Parquet 表的字段类型有时需要手工指定。我更喜欢用 Spark SQL JDBC 写df spark.sql(SELECT dt, station_id, direction, flow_cnt FROM dws_station_hourly_flow WHERE dt2026-02-16) df.write.mode(overwrite).jdbc( urljdbc:mysql://10.0.0.8:3306/bigdata, tablestation_flow_daily, properties{ user: bigdata, password: change_me, driver: com.mysql.cj.jdbc.Driver, rewriteBatchedStatements: true, }, )rewriteBatchedStatementstrue对 JDBC 批量写入很关键不设这个参数时几千行要一条条执行导出会非常慢设了以后耗时能降到十秒内。表结构提前手工建好带好索引Spark 自动建表不会建索引。5.2 ECharts 热力大屏网格数据而非原始点前端拿到的数据一定是聚合后的网格不要直接把 GPS 明细发给浏览器。下面用 ECharts 画热力fetch(/api/grid_flow?dt2026-02-16) .then(r r.json()) .then(rows { const heatmap rows.map(row [row.center_lon, row.center_lat, row.flow_cnt]); myChart.setOption({ geo: { map: china, roam: true }, series: [{ type: heatmap, data: heatmap, pointSize: 15, blurSize: 12 }] }); });center_lon和center_lat必须是由grid_id反算出的网格中心坐标不能用原始 GPS 点。前端拿到原始点后自己聚合数据量大且无法控制大屏会卡死。5.3 三个验证技巧让结果别在最后一公里出错第一对比总量。把 DWS 每小时的flow_cnt累加应该等于当天 ODS 原始记录数去掉过滤后的明细数SELECT dt, sum(flow_cnt) FROM dws_station_hourly_flow WHERE dt2026-02-16 GROUP BY dt;第二抽样验证。取一个站点的原始刷卡数据手工数 30 分钟内的通过量对比 DWS 表里该站该时段的flow_cnt。误差超过 1% 就回查清洗逻辑。第三用explain检查分区裁剪。很多结果偏大是因为 WHERE 条件被谓词下推绕过读了多个分区。一个大屏报表只查一天如果读到的分区数量超过一天就要回看 SQL 里的dt过滤。导出到 MySQL 前把dws_station_hourly_flow按hour做一次ORDER BY这样 MySQL 索引扫描连续前端按时段拉取数据更快。本文还有配套的精品资源点击获取