接到一个地图数据分析的需求时我手上只有一张上亿条的用户App定位日志表领导要求统计全城各区域的活跃人数、跨区域流动路线还有商圈热度排名。我第一反应是用Python脚本去处理结果数据全量拉到本地就要很久循环处理经常内存溢出后来我决定改用HiveSQL来做这一整套任务绕开了单机资源瓶颈把ETL、网格化、热度统计、OD流向计算全部用SQL脚本跑完从明细到报表一次成型。这篇文章就把这套“可直接执行的HiveSQL地图数据分析全量脚本”完整拆开讲包括每个脚本的设计原因、完整可运行的SQL、以及在真实集群上踩过的坑。无论你是数据仓库工程师、数据分析师还是刚接触位置数据的后端开发照着这套思路绕开我踩过的坑半天内就能把离线地图分析管道搭起来。1. 拿到亿级定位日志后我先做三条预处理清洗、切分、分区很多初学者拿到位置数据就急着写GROUP BY统计先不说结果可不可信光是跑一遍全表扫描就可能挂掉。我在尝试第一版脚本时就因为没做预处理一个3小时的任务最后只跑出了10万个“孤岛记录”——大量经纬度在境外、uid为空、时间字段还是个字符串。写任何统计SQL之前先把数据治理做成一个公共子查询后面所有脚本都复用这一层。1.1 原始位置表的基本结构字段、分区与隐藏的数据类型陷阱假设我们的原始日志表叫ods_pos_log_daily按日期分区核心字段通常长这样字段名类型说明uidSTRING用户唯一标识也可能是加密IDlngDOUBLE经度WGS84坐标系范围-180到180latDOUBLE纬度范围-90到90location_timeBIGINT毫秒级Unix时间戳provinceSTRING省份citySTRING城市app_idSTRING应用来源这里最容易踩的第一个坑location_time到底是秒还是毫秒。我遇到过上游把毫秒时间戳当秒传给下游的如果统一除以1000处理未来都要返工。正确做法是在建表或预处理时先用from_unixtime(cast(location_time / 1000 as bigint), yyyy-MM-dd HH:mm:ss)转成可读时间然后把转换后的字符串统一作为后续脚本的时间字段。第二个坑是经纬度的标准程度。GPS原始坐标可能是WGS84、GCJ02或BD09混合的如果直接混用后续网格定位会偏向几十到几百米。对于离线分析我们可以在预处理阶段用SQL简单判断如果是国内坐标的合理范围经度73~135、纬度3~53则保留如果是明显不合法的数值直接过滤。第三个坑是分区。原始表按dt字符串分区格式yyyy-MM-dd这是Hive最常规的做法。统计脚本必须强制使用where dt ${hiveconf:bizdate}这样的分区裁剪否则一个不小心就会扫描全表。1.2 清洗脚本过滤海外坐标、重复上报、异常跳点下面这段SQL是我在dwd_pos_log_daily之前加的一层view或临时表它实现了三个过滤坐标范围过滤只保留国内合理经纬度坐标做空值过滤。重复上报过滤同一个用户同一秒内上报多次只保留一条。跳点过滤相邻两条记录之间如果速度超过120km/h判定为GPS漂移或用户乘坐交通工具时的临时噪声这一条后面做停留识别很重要。SET hive.exec.dynamic.partitiontrue; SET hive.exec.dynamic.partition.modenonstrict; SET hive.exec.max.dynamic.partitions.pernode1000; CREATE TABLE IF NOT EXISTS dwd_pos_log_daily ( uid STRING, lng DOUBLE, lat DOUBLE, loc_time_str STRING, loc_ts BIGINT, province STRING, city STRING ) PARTITIONED BY (dt STRING) STORED AS ORC; INSERT OVERWRITE TABLE dwd_pos_log_daily PARTITION (dt${hiveconf:bizdate}) SELECT uid, lng, lat, from_unixtime(cast(location_time / 1000 as bigint), yyyy-MM-dd HH:mm:ss) AS loc_time_str, cast(location_time / 1000 as bigint) AS loc_ts, province, city FROM ( SELECT t.*, row_number() over ( PARTITION BY uid, location_time ORDER BY location_time ) AS rn FROM ods_pos_log_daily t WHERE dt ${hiveconf:bizdate} AND uid IS NOT NULL AND uid ! AND lng BETWEEN 73 AND 135 AND lat BETWEEN 3 AND 53 ) tmp WHERE rn 1;这一段只是把明细落到dwd层过滤了异常坐标和重复数据。真正更严格的“跳点”过滤我会放到后面的网格编码里做因为跳点判断依赖相邻两条记录做距离计算单独提出来会多一次全表自连接先做粗过滤更划算。1.3 正确的日期分区用法别把时间过滤写死在脚本里脚本里所有日期都用${hiveconf:bizdate}代替调用时通过-hiveconf bizdate2025-01-06传入。好处是同一个脚本不用改代码就能每天跑也方便补数据时传任意日期。如果你想跑一个日期区间可以在外层shell脚本循环调用HiveSQL内部只处理单日。这一步里还要提醒一点HiveSQL默认分区裁剪只对静态分区生效。如果你用WHERE dt ${hiveconf:start_date} AND dt ${hiveconf:end_date}这种动态范围性能也许还行但如果你加上PARTITION(dt${hiveconf:bizdate})这种静态分区执行计划会更优。所以我建议核心聚合脚本都采用单日静态分区写INSERT OVERWRITE ... PARTITION(dt${hiveconf:bizdate})。2. 网格化方案为什么不用GeoHash我用纯SQL投影网格编码地图数据分析绕不开“把经纬度变成区域”否则统计结果会散成一堆无意义的点。有几种常见做法直接对经纬度四舍五入、用GeoHash字符串、或者用固定间距的投影网格编码。我最终选择了后者原因下面逐个说清楚。2.1 三种网格方案的对比四舍五入、GeoHash、投影网格直接round(lng,3)/round(lat,3)优点是简单到一句话就能实现缺点是不同纬度下经度方向同样的度数代表的实际距离不同导致网格在高纬度地区被拉扁。而且边界是沿着经线和纬线的做“邻里关系”分析时会出现斜对角距离计算时的扭曲。GeoHash这是最流行的方法可以把经纬度编码成字符串前缀相同的点距离近。但它有几个问题首先是Hive没有内置GeoHash函数生产环境申请UDF麻烦其次是边界不规则同一个区域的编码跨度可能不如投影网格好理解再者是精度选择不够直观几位编码对应多大面积要看表。投影网格编码先把经纬度通过Web墨卡托投影变成平面坐标单位米然后除以网格边长向下取整生成网格ID。这样网格是均匀的正方形跨网格距离、周长面积都符合直觉。最关键的是纯HiveSQL就能实现不依赖任何UDF脚本拿到任何集群都能直接跑。我在实际项目里选了投影网格网格边长为500米和1000米两种商圈分析用500米全城活跃度分析用1000米。2.2 纯SQL实现Web墨卡托投影与网格ID生成Web墨卡托投影公式很简单经度线性映射纬度做对数变换-- 投影公式将经纬度转为Web墨卡托平面坐标单位米 x lng * 20037508.34 / 180; y ln(tan((90 lat) * 3.141592653589793 / 360)) / (3.141592653589793 / 180) * 20037508.34 / 180;HiveSQL写法如下。我们把它封装成一个公共视图后续所有脚本都从这个视图取数保证网格口径一致CREATE VIEW IF NOT EXISTS v_pos_grid AS SELECT uid, dt, from_unixtime(floor(loc_ts / 3600) * 3600, yyyy-MM-dd HH:00:00) AS hour_ts, lng, lat, loc_ts, concat( cast(floor((lng * 20037508.34 / 180) / 500) as bigint), _, cast(floor((ln(tan((90 lat) * 3.141592653589793 / 360)) / (3.141592653589793 / 180) * 20037508.34 / 180) / 500) as bigint) ) AS grid_id_500, concat( cast(floor((lng * 20037508.34 / 180) / 1000) as bigint), _, cast(floor((ln(tan((90 lat) * 3.141592653589793 / 360)) / (3.141592653589793 / 180) * 20037508.34 / 180) / 1000) as bigint) ) AS grid_id_1000, city FROM dwd_pos_log_daily WHERE dt ${hiveconf:bizdate};如果你希望网格更细可以把500换成其他边长比如200米、1000米。由于Web墨卡托在北纬80度以内误差可以接受用它做城市级别分析足够。2.3 网格边长怎么选精度和计算量的取舍网格边长直接决定聚合后的行数和结果粒度。我的经验值是分析场景推荐网格边长预估一个大型城市的网格数量热点区域概览1000米约2000~4000商圈/写字楼分析300~500米约8000~15000街道/店铺级分析100~200米约50000以上网格越小单个网格中包含的用户越少隐私性更好但计算量大得多。500米网格在单日亿级记录下聚合后大约有几万行Hive跑起来没有压力。如果你追求100米级别建议先按500米聚合一次再用窗口函数做近邻处理而不是直接把所有点映射到100米网格。3. 核心统计脚本按小时/日/周聚合全量网格热度预处理和网格编码准备完之后第一个要跑的核心脚本就是“网格活跃度统计”。它输出的是一张全量结果表结构大概是dt, hour_ts, grid_id, pv, uv,活跃用户列表。3.1 输出表结构设计别只存一个汇总值建表时我故意把所有维度和指标都展开这样后续无论是出报表、写BI还是做算法模型都可以直接查这张表CREATE TABLE IF NOT EXISTS ads_grid_active_stat ( dt STRING COMMENT 日期分区, hour_ts STRING COMMENT 小时时间戳, grid_id STRING COMMENT 500米网格ID, pv BIGINT COMMENT 点位访问量, uv BIGINT COMMENT 独立用户数, new_uv BIGINT COMMENT 新增用户数该日第一次在本城市出现, avg_cnt DOUBLE COMMENT 平均每用户上报次数 ) PARTITIONED BY (dt STRING) STORED AS ORC;new_uv单独这个字段可能不是必须的但它是我后来做渠道分析时反查出来的。如果你不是算新增用户可以把这张表简化成dt, hour_ts, grid_id, pv, uv。保留额外字段的好处是不需要每次重新跑原始表。3.2 全量热度计算SQL一条语句从明细到聚合这里用到了公共视图v_pos_grid代码看起来很短但它把前面所有预处理全部串联起来了INSERT OVERWRITE TABLE ads_grid_active_stat PARTITION(dt${hiveconf:bizdate}) SELECT ${hiveconf:bizdate} AS dt, hour_ts, grid_id_500 AS grid_id, count(*) AS pv, count(DISTINCT uid) AS uv, count(DISTINCT if(is_new.uid IS NOT NULL, uid, NULL)) AS new_uv, count(*) / count(DISTINCT uid) AS avg_cnt FROM v_pos_grid v LEFT JOIN ( SELECT uid FROM ( SELECT uid, min(dt) AS first_dt FROM v_pos_grid GROUP BY uid ) t WHERE first_dt ${hiveconf:bizdate} ) is_new ON v.uid is_new.uid GROUP BY hour_ts, grid_id_500;这段SQL里count(DISTINCT if(is_new.uid IS NOT NULL, uid, NULL))是HiveSQL里非常典型的新增用户计算方式用if把非目标记录置为NULL再对NULL做去重天然不计数。如果你一开始写的是count(DISTINCT uid) filter (where ...)注意Hive老版本不支持FILTER语法用if是最兼容的写法。3.3 把脚本包装成每日任务Shell脚本里循环调用单条hive -f命令当然能跑但批量跑多个日期时最好包一层Shell脚本。我在生产环境是这样调用的#!/bin/bash # run_grid_stat.sh bizdate${1:-$(date %F -d yesterday)} for city in beijing shanghai guangzhou; do hive \ -hiveconf bizdate${bizdate} \ -hiveconf city${city} \ -f grid_active_stat.hql done注意HiveSQL里如果用到city过滤要在v_pos_grid视图定义时把city列暴露出来并在外层的WHERE city ${hiveconf:city}过滤。如果你只需要一个城市这个Shell包装可以直接省略。4. OD流向统计用纯SQL构造用户移动轨迹并计算区域间OD网格热度解决的是“哪里人多”可你还要回答“人从哪里来往哪里去”。这就是ODOrigin-Destination分析。传统做法是用MapReduce或Spark写复杂逻辑其实HiveSQL借助窗口函数也能做而且效率不差。4.1 怎么定义OD同一用户连续网格发生变化才算一段出行一条OD记录至少包含四要素用户、出发时间、出发网格、到达时间、到达网格。如果直接把用户的每一条位置记录都当成一次出行数据量会爆炸而且原地抖动会被误判成移动。我的做法是以小时为窗口取这个用户在该小时内的第一个网格和最后一个网格如果两者不同则生成一条OD。CREATE TABLE ads_grid_od_flow ( dt STRING, hour_ts STRING, uid STRING, from_grid STRING, to_grid STRING, outflow_time STRING, inflow_time STRING ) PARTITIONED BY (dt STRING) STORED AS ORC; INSERT OVERWRITE TABLE ads_grid_od_flow PARTITION(dt${hiveconf:bizdate}) SELECT dt, hour_ts, uid, from_grid, to_grid, min_loc_str AS outflow_time, max_loc_str AS inflow_time FROM ( SELECT dt, hour_ts, uid, first_value(grid_id) over (PARTITION BY uid, hour_ts ORDER BY loc_ts) AS from_grid, last_value(grid_id) over (PARTITION BY uid, hour_ts ORDER BY loc_ts ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING) AS to_grid, first_value(loc_time_str) over (PARTITION BY uid, hour_ts ORDER BY loc_ts) AS min_loc_str, last_value(loc_time_str) over (PARTITION BY uid, hour_ts ORDER BY loc_ts ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING) AS max_loc_str, row_number() over (PARTITION BY uid, hour_ts ORDER BY loc_ts DESC) AS rn FROM v_pos_grid ) t WHERE rn 1 AND from_grid ! to_grid;这里有一个容易写错的地方last_value如果不加ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING它默认只取窗口内最后一行而窗口默认是当前行所以会有偏差。我要不是被结果坑过可能到现在还在用max(loc_ts)去取最后一条。4.2 OD聚合输出全城所有网格之间的流向矩阵OD明细表有了下一步就是按出发网格和到达网格聚合得到区域间的流量矩阵INSERT OVERWRITE TABLE ads_grid_od_matrix PARTITION(dt${hiveconf:bizdate}) SELECT from_grid, to_grid, count(*) AS od_cnt, count(DISTINCT uid) AS od_uv FROM ads_grid_od_flow WHERE dt ${hiveconf:bizdate} GROUP BY from_grid, to_grid;如果我需要看商圈之间的通勤情况就再加一个WHERE限定from_grid和to_grid的坐标范围。为了让结果更可读可以在结果表里关联一个网格中心点坐标表把grid_id转成经纬度方便画图。网格中心点可以通过网格ID反算出来也可以在视图里直接把网格中心经纬度保留下来。这点我在后面第6部分再讲。5. 停留点识别与商圈热度挖掘如何从位置日志里提取“人真的呆了一阵子”只看瞬时位置还不够要判断一个区域是不是热门商圈得看“停留时间”和“反复访问”这两个指标。直接统计PV会高估路过的人而停留时长才能区分出来“逛街的人”和“开车路过的人”。5.1 用HiveSQL定义“停留”同一个网格内连续滞留超过30分钟因为我们已经把位置按500米网格归一化了一个最简单且有效的定义是同一个用户同一小时在同一个网格里最早一条记录和最晚一条记录之间的时间差大于等于30分钟并且记录数不少于3条则认为发生了一次有效停留。这种“粗节点”定义对商圈分析完全够用。CREATE TABLE ads_grid_stay_detail ( dt STRING, uid STRING, grid_id STRING, enter_time STRING, leave_time STRING, stay_secs BIGINT, sample_cnt INT ) PARTITIONED BY (dt STRING) STORED AS ORC; INSERT OVERWRITE TABLE ads_grid_stay_detail PARTITION(dt${hiveconf:bizdate}) SELECT dt, uid, grid_id, min(loc_time_str) AS enter_time, max(loc_time_str) AS leave_time, unix_timestamp(max(loc_time_str), yyyy-MM-dd HH:mm:ss) - unix_timestamp(min(loc_time_str), yyyy-MM-dd HH:mm:ss) AS stay_secs, count(*) AS sample_cnt FROM v_pos_grid WHERE dt ${hiveconf:bizdate} GROUP BY dt, uid, grid_id HAVING unix_timestamp(max(loc_time_str), yyyy-MM-dd HH:mm:ss) - unix_timestamp(min(loc_time_str), yyyy-MM-dd HH:mm:ss) 1800 AND count(*) 3;这里HAVING里用了聚合函数unix_timestamp(max(...))在Hive的HAVING子句里是允许的。如果你在集群上遇到解析问题可以把这个表达式放在子查询里再在外层过滤。这个方法有个小瑕疵同一用户一天内分多个时间段在同一个网格停留会被合并成一条超长停留。解决思路是把相邻两条记录的时间间隔大于10分钟作为分割点给每次访问打一个session号这是个进阶题目本文不展开但你需要知道这个特性的存在。5.2 商圈热度指数人数×时长×频次的加权评分有了停留明细热力榜就可以往下走了。我把热力评分定义为热度指数 独立访客数 × 0.5 总停留时长(分钟) / 10 人均停留时长(分钟) × 0.02这个公式是我根据业务经验调的权重可以根据具体情况改。关键点在于既能反映人流规模又能体现深度停留。用SQL算非常快INSERT OVERWRITE TABLE ads_grid_heat_rank PARTITION(dt${hiveconf:bizdate}) SELECT grid_id, count(DISTINCT uid) AS visitor_uv, sum(stay_secs) / 60 AS total_stay_minutes, sum(stay_secs) / count(DISTINCT uid) / 60 AS avg_stay_minutes, count(DISTINCT uid) * 0.5 sum(stay_secs) / 600 (sum(stay_secs) / count(DISTINCT uid) / 60) * 0.02 AS heat_score FROM ads_grid_stay_detail WHERE dt ${hiveconf:bizdate} GROUP BY grid_id ORDER BY heat_score DESC LIMIT 200;如果你要拿这个结果和商圈POI数据叠加可以把网格ID关联到商圈的POI明细表筛选出某个商圈范围内的网格。这一环需要一张“网格-POI映射表”我一般在预处理阶段就生成好。6. 脚本上生产前一定要确认的10个边界条件与优化习惯这套脚本的第一版在测试集群上能跑通不代表放到生产集群上就没问题。我整理了自己踩过的10个坑按重要性从高到低排。6.1 数据倾斜热点网格影响reduce阶段速度激活区域如市中心会聚集大量数据GROUP BY grid_id时热点网格对应的reduce任务可能跑了30分钟其他reduce几十秒就结束了。解决方法有三种把count(DISTINCT uid)改为先按uid, grid_id去重再按grid_id聚合减少单次去重的压力在GROUP BY前先用WHERE按城市分区拆成多个任务或用热点抽样的方式单独计算Top20网格后合并但那样代码复杂度高我这里不推荐一开始就做。6.2 全表扫描的隐患过滤条件必须下推到分区在v_pos_grid视图里我已经把WHERE dt ${hiveconf:bizdate}写死了。但在Hive里视图并不会自动把后续查询的过滤条件下推所以你在外层查v_pos_grid时最好继续保留WHERE dt ${hiveconf:bizdate}。否则某些优化器版本会先扫描所有分区再作过滤等于是全表扫描。6.3 经纬度为0、-9999之类的脏点原始日志里常见的脏点包括经纬度为0、为1、或为负数但明显不对。我们的过滤条件是lng between 73 and 135这个区间只适用于中国大陆。如果项目是海外地图需要把范围改成对应国家的边界。更好的做法是查一份“行政区划边界表”用纬度范围经度范围提前生成合法网格清单再做INNER JOIN过滤。6.4 时区问题时间戳到底按哪个时区解释location_time服务端记录的一般是UTC时间毫秒戳。如果直接用from_unixtime会得到UTC时间统计出来的“小时活跃度”会比本地时间早8小时。所以在清洗阶段我通常使用from_utc_timestamp(from_unixtime(...), Asia/Shanghai)或者简单地在location_time上加上8小时再格式化。不同业务差异很大但这点必须在脚本里写清楚否则会出现凌晨时间整点对不上的情况。6.5 统一网格ID长度使用固定长度字符串生成网格ID时floor(x/500)得到的结果位数不固定可能生成234_567和234_56。如果后面拿字符串拼接或做索引位数不齐会带来排序问题。我的习惯是用lpad(cast(... as string), 6, 0)补零保证所有网格ID长度一致。如果两位数的城市数量超过1000万建议把经度方向的ID直接补到7位。6.6 使用ORC Snappy压缩中间表原始表可能是TextFile或Parquet但我们在dwd层创建的表建议统一用STORED AS ORC同时设置TBLPROPERTIES (orc.compressSNAPPY)。我用一个300GB的TextFile表做过对比转成ORC之后查询数据量减少约75%跑批时间缩短一半以上。6.7 小文件过多会拖垮NameNode和降低查询效率INSERT OVERWRITE会按照reduce数量产出文件。如果小时级聚合结果只有几千行但reduce数量默认是几十个就会产生大量小文件。在跑批前设置SET hive.merge.mapfilestrue; SET hive.merge.mapredfilestrue; SET hive.merge.size.per.task256000000; SET hive.merge.smallfiles.avgsize16000000;这样最终输出文件会尽量合并成256MB左右的合理大小。6.8 中间表清理与生命周期管理OD明细、停留明细都属于中间结果我一般保留7天热度统计表保留30天。可以在每天任务结束后用一条命令清理过期分区ALTER TABLE ads_grid_od_flow DROP IF EXISTS PARTITION(dt${hiveconf:clean_date});别嫌这个动作多余时间长了表分区数会膨胀而且查询时如果你不小心忘了写分区条件会扫出一堆无用的旧数据。6.9 脚本参数化用Shell变量传递日期、城市、网格边长我建议把所有可能变动的参数全部抽到Shell脚本外层。日期、城市、网格边长这三个参数至少是必要的。比如网格边长一变视图里的两个数字就得改如果你把500和1000散落在多个SQL里改起来容易漏。与其这样不如定义视图时就写两个固定常用边长后续要改再单独写新视图。6.10 用幂等设计保护重跑“全量脚本”最大的优点就是可以随时重跑。每个INSERT OVERWRITE都会覆盖当天分区只要传入的日期一致跑多少次结果都一样。这要求你在开发时严禁用INSERT INTO去追加结果否则重跑一遍会出现双份数据。我在代码评审时会把这条作为硬性规定。个人经验把网格计算沉淀成公共视图而不是在每段SQL里重复公式最后分享一个我在实际维护这套脚本三个月后最大的体会网格投影公式和停留判定逻辑一定要沉淀成公共视图或公共表而不是每次在SQL里复制粘贴。我一开始图省事在热度脚本和OD脚本里各写了一遍投影公式后来城市范围从国内扩展到周边国家时坐标范围改了OD脚本忘了同步结果OD结果里出现了大量海外网格花了一个下午才排查出来。自那以后我把所有公共逻辑都收敛到v_pos_grid这个视图里每次口径变更只改一个地方下游脚本全部自动跟着变。这套“可直接执行的HiveSQL地图数据分析全量脚本”核心不在于某一串SQL写得有多炫而在于它把“清洗、网格化、热度统计、OD分析、停留识别”这一整条链路用最少的UDF、最直观的逻辑组合了起来。你拿过去之后最需要改的点无非是表名、城市范围、网格边长和业务打分权重。把这几样替换成你自己的剩下的逻辑可以直接跑这大概就是“全量脚本”最大的价值。