首页 / 资讯中心 / 文章详情

用户画像标签存储架构:Hive/MySQL/Hbase/ES四库协同实践

用户画像标签存储架构:Hive/MySQL/Hbase/ES四库协同实践 ★ FEATURED ARTICLE
简介面向大数据开发工程师与数据仓库工程师的用户画像标签数据存储完整方案PDF适合正在规划或优化画像平台的团队参考。资源围绕Hive、MySQL、Hbase、Elasticsearch四种存储引擎展开讲清各自在画像场景中的定位与分工并给出用户标签表、标签聚合表、人群计算表的字段设计与分区策略。压缩包内共1个PDF文件大小约1.34MB内容全部集中在文档中。目前已有211人学习使用。文档覆盖从Hive向Hbase、MySQL同步标签数据的完整流程包括Sqoop落地方式、数据量校验机制以及将圈定人群推送给广告系统、Push系统、客服系统等业务方的调度思路。其中穿插tag表、tagmap表等真实表结构示例与执行过程说明可帮助读者对照设计自己的画像存储层规避表结构混乱与数据同步丢失等常见问题。1. 用户画像系统的标签数据存储为什么同一套数据要拆进四个数据库很多刚接触用户画像系统的人第一反应是“标签不就是一张用户表加几个字段吗MySQL 就够了吧”。真正跑过生产环境的人都不会这么想。用户画像系统的核心是标签数据存储一套标签从离线计算到线上服务至少要经过 Hive、MySQL、Hbase、Elasticsearch 四个库每个库承担完全不同的角色少一个都会出问题。这个结论不是设计文档里拍脑袋定的而是由数据量、查询时延、写入频率和业务系统的接入方式共同逼出来的。本文拆解的这份解决方案把四种数据库的存储定位、表结构设计、同步链路和校验机制都讲得很清楚适合正在搭画像平台的数据工程师、数仓开发以及准备把标签服务线上化的后端同学。新手可以照表结构复现熟手可以直接拿走里面的同步校验方案。2. Hive 主存标签结果集tag 表、tagmap 表与人群计算表的建表细节2.1 为什么画像的“数据底座”必须放 Hive画像相关的数据有一个共同特点数据量极大且计算逻辑复杂。一个中型的电商平台每天活跃用户百万级每个用户身上挂几十个标签再按时间分区累积单日新增数据就是千万甚至亿级行。这种计算量下标签的生成作业跑的是 MapReduce 或者 Spark结果写 HDFS。MySQL 扛不住这么大的写入Hbase 虽然能扛写入但跑批量作业时没有 Hive 的 SQL 生态方便。所以 Hive 在画像系统里的定位是“所有标签相关计算结果集的默认落点”包括用户标签表、标签聚合表、人群计算结果表全部先落在 Hive 数仓里。Hive 存储的一个关键设计是分区。方案里明确提到tag 表按“日期 标签主题”双分区设计。日期分区解决的是每天全量标签快照的隔离问题标签主题分区解决的是 ETL 调度时同时计算多个标签的并行插入问题。这里有一个细节值得注意如果只按日期分区那么同一天内不同主题的标签写入时需要反复动态分区或覆盖写调度上非常被动加了标签主题分区后每个标签作业可以独立往自己的分区里写互不干扰某个标签的计算失败也只会影响它自己的分区重跑成本极低。2.2 tag 表一条用户一个标签一行记录用户标签表tag 表是画像系统最底层的表记录的是“用户 id、标签 id、标签权重”的最小粒度对应关系。以下建表语句是这套方案的核心结构几乎所有标签结果都会落到这种格式里CREATE TABLE dw.profile_tag_userid ( user_id STRING COMMENT 用户ID, tag_id STRING COMMENT 标签ID, tag_weight DOUBLE COMMENT 标签权重, tag_type STRING COMMENT 标签类型, data_date STRING COMMENT 数据日期分区 ) PARTITIONED BY (data_date STRING, tag_type STRING) STORED AS ORC;向 Hive 插入测试数据的常见做法是INSERT OVERWRITE TABLE dw.profile_tag_userid PARTITION (data_date2024-05-20, tag_typeuserid_all_paid_money) SELECT user_id, userid_all_paid_money AS tag_id, sum(pay_amount) AS tag_weight FROM dwd_order_detail WHERE data_date2024-05-20 GROUP BY user_id;这段 SQL 的逻辑是从订单明细表按用户汇总支付金额把汇总结果作为标签权重写入 tag 表分区字段 tag_type 取值为标签主题名。注意这里用了 INSERT OVERWRITE 而不是 INSERT INTO因为同一分区每天重复计算时需要保证幂等。tag_type 分区字段还有一个非常实用的功能不同类型的标签如消费能力标签、活跃度标签、偏好标签可以并行向同一张表的不同分区写入不需要锁表也不会互相覆盖。2.3 tagmap 表把同一个人身上的标签聚合成一条tag 表的粒度是一用户一标签一行但查询时往往是“一个用户身上的全部标签”。如果每次查询都扫 tag 表全分区性能非常差。于是方案里设计了 tagmap 表将同一个用户的所有标签聚合到一行用 Map 或者拼接字符串的方式存储。结构类似CREATE TABLE dw.profile_user_map_userid ( user_id STRING COMMENT 用户ID, tag_map MAPSTRING, DOUBLE COMMENT 标签ID到权重的映射, data_date STRING COMMENT 数据日期分区 ) PARTITIONED BY (data_date STRING) STORED AS ORC;聚合执行的思路是从 tag 表按 user_id 分组把 tag_id 和 tag_weight 收集成 Map。常见做法是写一个 HiveQL用 collect_list 和 map 函数组合构造 Map或者直接用 Spark 的 mapFromEntries 函数处理。聚合的目的不是减少数据量而是把“用户视角”的查询从扫描 N 行变成扫描 1 行。这个表是后续人群圈选和画像查询的主力表。2.4 人群计算表从标签圈人到业务系统的数据出口人群表记录的是圈人结果字段包括用户 id、人群名称 id、推送到的业务系统。这个表的典型特征是它需要关联订单表和用户收货信息表得到用户的手机号等联系信息推送给外呼中心或短信系统。核心 join 逻辑是人群表先 join 订单表拿到订单编号再 join 收货信息表拿到手机号。三步 join 下来本质上是从“标签筛选”过渡到“运营触达”。这个表的设计要点是它除了存 user_id还必须冗余业务系统需要的字段。因为下游系统不一定能通过 user_id 反查用户信息直接把手机号、订单号冗余在人群结果表里同步到业务库时可以减少一次 join。方案中给到的表结构里人群名称 id 和业务系统标识是必备字段业务系统标识决定了这条数据后面同步到 MySQL 还是 Hbase。3. MySQL 管元数据与校验画像系统的“控制面”搭建3.1 MySQL 在画像系统里的三个职责MySQL 在画像系统里不存明细标签数据它管的是三类东西画像标签的元数据、结果集校验信息、同步到业务系统的数据。一句话概括Hive 是数据面MySQL 是控制面。元数据维护着标签的 id、名称、主题、一级二级分类、标签描述等。一个标签从需求提出到上线先在元数据表里登记然后才进入 Hive 计算流程。结果集校验信息包括当日标签覆盖用户量、当日与前一日波动比例、当日标签覆盖用户占活跃用户比例、任务是否继续执行的标志位。这些校验表的存在是为了避免 0 点调度任务跑完之后没人检查是否产出正常就直接同步线上。3.2 元数据表与校验表的字段设计元数据表结构设计上至少要包含这些字段字段名类型说明tag_idvarchar(64)标签唯一 ID与 Hive 中 tag_id 对齐tag_namevarchar(128)标签名称tag_themevarchar(64)标签主题与 Hive 分区对应level1_categoryvarchar(64)一级分类level2_categoryvarchar(64)二级分类tag_desctext标签描述包括口径和计算逻辑ownervarchar(64)负责人用于标签运维结果集校验表的核心字段是标签 id、数据日期、覆盖用户量、波动比例、校验状态。其中波动比例的计算逻辑是当日覆盖用户量 - 前日覆盖用户量/ 前日覆盖用户量超过阈值时就写入告警标志位。这个标志位在调度系统里会被读取决定下游任务是否继续执行。很多团队的调度依赖只用“前一个任务是否成功”但画像场景更严谨的做法是“前一个任务成功且数据质量校验通过”否则会带着错误数据一路同步到线上。3.3 从 Hive 同步 MySQLSqoop 命令与 Python 脚本的取舍同步到业务系统这一步方案里给了两个路径Sqoop 同步和 Python 脚本同步。Sqoop 适合一次性的、整表的同步命令示例sqoop export \ --connect jdbc:mysql://mysql-host:3306/user_profile \ --username root --password secret \ --table customer_push_list \ --export-dir /user/hive/warehouse/dw.db/customer_push_list \ --input-fields-terminated-by \001 \ --columns user_id,phone,order_id,campaign_id \ --batch \ --update-mode allowinsert \ --update-key user_id,campaign_id这段命令的逻辑说明从 Hive 的 customer_push_list 表导出数据到 MySQL 的 customer_push_list 表字段分隔符是 Hive 默认的 \001update-mode 设置为 allowinsert 表示已存在的主键更新、不存在则插入。这个参数很关键——业务系统推送场景里同一个用户可能被多个活动圈中以 user_id 和 campaign_id 作为联合主键才能避免数据互相覆盖。日常维护中我更倾向于写一个 Python 脚本同步。因为 Sqoop 的缺字段、字符集问题在画像这种多口径场景里比较常见Python 脚本可以同步过程里做数据量对比、格式转换和异常重试。脚本的核心逻辑是查 Hive 当日分区总数查 MySQL 目标表当前总数两数对不上就发告警不执行写入。4. Hbase 与 Elasticsearch圈人结果入库与在线查询的两种路径4.1 为什么圈人结果要推 Hbase线上服务的时延要求画像系统产品化之后运营人员圈完一个人群这批人需要进入广告系统、push 消息系统等线上服务。这类业务系统读数据的特点是QPS 高、单次查询要求毫秒级返回、数据量从几十万到几千万不等。Hive 完全扛不住这种查询压力MySQL 在千万级数据 高并发场景下也容易出现连接打满。Hbase 的随机读写能力在这个场景里是刚需。同步链路一般是这样运营在页面上配置规则规则和圈出的人群标签集存 MySQLSpark 作业读取 MySQL 中的标签集信息去 Hive 的标签表计算具体人群计算结果写入 Hive 当日分区然后再由同步作业把该分区数据写入 Hbase 对应表。4.2 Hive 映射 Hbase 表的两种建表方式同步 Hive 数据到 Hbase常见做法是创建 Hive 到 Hbase 的映射表。第一种方式是启动 Hive创建一张映射到 Hbase 的 Hive 外部表然后向该映射表插入测试数据Hbase 侧会自动建表。第二种方式是直接在 Hbase 侧建表然后用 Hive 关联查询。我实际落地更推荐第一种因为可以直接复用 Hive 的 SQL 逻辑插入语句执行起来就是一个 MapReduce 作业跑完数据自然落在 Hbase 里。如果非要用 Hbase 原生 API 写批量导入常见做法是使用 Hbase 的 BulkLoad 工具先让 MapReduce 作业生成 HFile再加载到 Hbase 表。这个方式跳过了 WAL 写入速度比直接 put 快很多适合千万级以上的初始导入。但注意 BulkLoad 只适合一次性大批量导入不适合频繁的小批量同步因为生成 HFile 和加载 HFile 的过程需要额外的 HDFS 空间操作不当容易把 region 弄得不均匀。4.3 Hbase 同步的两种数据校验方案这是整套方案里工程含金量最高的一段。因为灌入 Hbase 的数据直接应用到线上反馈到用户那里任何数据问题都会直接暴露所以同步必须加校验。方案里给了两种做法。第一种Hive 同步 Hbase 后先在 Hbase 里建一个 temp 临时表数据写入临时表再校验临时表和 Hive 表的数据量差异。如果差异在可接受范围内把 Hbase 临时表 rename 成正式表。这利用了 Hbase 改表名的低成本特性但代价是同步期间线上读不到当天最新数据。第二种Hive 同步 Hbase 后直接写入正式表同时建立一张状态表。同步完成触发校验校验通过后在状态表写入当天日期和校验状态。线上接口请求时只读取状态表中最近日期的数据。如果同步异常状态表不更新线上继续读取前一天的数据。这种方案的优点是不影响线上读取缺点是数据质量有问题时当天数据可能已经暴露给线上一段时间了所以校验任务必须紧跟在同步任务之后立刻执行。我在生产上更倾向第二种方案配合告警同步完成立刻校验数量对不上马上钉钉告警DBA 介入前线上可能只有几分钟的脏数据窗口但总比长时间没有新鲜数据要好。4.4 Elasticsearch 在画像里的真实定位不是替代 Hbase方案里对 Elasticsearch 的描述很务实一个开源的分布式全文检索引擎近乎实时地存储、检索数据扩展性好可以处理 PB 级别数据。对于用户标签查询、用户人群计算、用户群多维透视分析这类对响应时间要求较高的场景可以考虑选用 Elasticsearch。实际工程中ES 在画像系统里主要干两类事。第一类是标签明细的快速查询比如运营输入一个用户 id 或 cookie立刻看到这个人的全部标签这个场景用 ES 的基于 user_id 的查询毫秒级返回第二类是人群的多维透视分析比如圈选了“近 30 天有购买行为且客单价高于 500 元”的人群后想看这批人的年龄分布、城市分布ES 的聚合查询aggs非常擅长这种场景一个 JSON 查询就能返回多维统计结果而 Hive 跑这种即席分析至少要分钟级。ES 存储画像数据的索引结构设计上常见做法是每个标签一个字段或者用 nested 对象存标签数组。前者适合标签数量少、结构固定的场景后者适合标签动态扩展的场景。索引设计时一定要给 tag_id 字段设置 keyword 类型否则分词后无法精确匹配。5. 标签数据同步避坑四个我实测过的生产环境问题5.1 Hive 同步 Hbase 后数据量对不上现象Hive 表统计有 5000 万条同步到 Hbase 后 count 只有 1000 万条。原因排查下来最常见的是 rowkey 设计冲突。Hbase 的 rowkey 如果只取 user_id那么同一用户的多条标签记录会互相覆盖导致数据行数变少。解决方式分两步第一确认目标表的 rowkey 是否包含区分字段同步前先预估同一 user_id 最多有几条记录用 user_id 标签类别或时间戳拼 rowkey第二同步完成后先查 Hbase 中目标表的行数再查源 Hive 表行数两者偏差超过 1% 时触发告警而不是等业务方反馈数据缺失。5.2 分区字段写错导致标签数据覆盖现象某天标签计算结果异常发现大量用户标签丢失。原因同一个 tag 表按日期和标签主题分区写 SQL 时把分区字段 tag_type 写成了另一个主题的值插入后覆盖了同主题前一天的数据。尤其是使用了 INSERT OVERWRITE 的分区表一旦分区条件写错不会有任何报错数据就静默覆盖了。解决所有分区表写入脚本里强制校验 data_date 和 tag_type 是否为当天预期的值脚本里写死日期变量不允许用 current_date 之类的动态值上线前先 SELECT COUNT(*) 看一眼目标分区今天是否有数据有数据立刻停。5.3 MySQL 同步时主键冲突导致任务卡死现象Sqoop 同步 Hive 数据到 MySQL跑了一小时后任务失败报 duplicate entry。原因MySQL 目标表主键设置不合理Hive 侧同一 user_id 出现重复导致 MySQL 插入冲突。解决在 Sqoop 命令里加 update-key 参数把 sync 模式改为 upsert同时检查源 Hive 表是否有重复有重复就在导出 SQL 里先去重。这条经验也适用于 Python 脚本同步——脚本里必须写 INSERT ... ON DUPLICATE KEY UPDATE 而不是 INSERT。5.4 校验标志位与数据写入不在同一事务里现象状态表显示当天数据校验通过但线上接口读到的还是旧数据。原因Hbase 数据先写入正式表再写状态表两步之间没有原子性。如果中间有短暂的时间窗口状态表里是当天最新日期但 Hbase 正式表还在写入过程中线上的查询就可能读到半新半旧的数据。解决把写入顺序反过来先写状态表再写 Hbase 正式表或者利用 Hbase 的 timestamp 版本控制线上读取时只取 status 表标记的日期之前的数据。更稳妥的方案是直接用 Hbase 的 coprocessor 或异步批处理框架把状态更新和数据写入包装在同一个任务里失败时统一回滚。6. 数据量校验脚本与临时表双写基线不重跑全量作业的妥协方案前面讲了那么多表结构和同步链路最后分享一个我一直在用的具体技巧如何高效做 Hive 到 Hbase 的日常数据校验而不用每次重跑全量作业。先说基线。画像系统里Hive 中的用户标签表每天都在增量更新但 Hbase 中的线上标签表通常是全量覆盖逻辑。全量覆盖的问题在于每次同步都要把全部用户的数据推一遍随着用户量增长同步时间越来越长。我的习惯是建一张基线表Hive 里维护一张最基础的“离线用户全量表”记录 user_id、最近活跃日期、核心标签摘要。同步 Hbase 时只同步当天有变化的用户没变化的用户继续沿用前一天 Hbase 里的数据。这样把全量同步变成了增量同步同步时间从两小时降到二十分钟。增量同步带来的新问题是怎么确认 Hbase 里的最终数据等于 Hive 全量数据的期望结果我用的校验策略是双层比对import happybase from pyhive import hive # 1. 从 Hive 查询当天应覆盖的总用户数 conn hive.Connection(hosthive-server, port10000) cursor conn.cursor() cursor.execute( SELECT count(*) FROM dw.profile_tag_userid WHERE data_date 2024-05-20 ) hive_cnt cursor.fetchone()[0] # 2. 从 Hbase 查询当天同步后实际覆盖的用户数 pool happybase.ConnectionPool(size3, hosthbase-server) with pool.connection() as hbase_conn: table hbase_conn.table(profile_tag_user) # 当天同步的用户都会带上日期前缀 rowkey hbase_cnt 0 for _ in table.scan(row_prefixb2024-05-20): hbase_cnt 1 # 3. 计算偏差率超过阈值则写告警状态不更新状态表 diff_ratio abs(hive_cnt - hbase_cnt) / max(hive_cnt, 1) if diff_ratio 0.01: print(f数据量偏差率达到 {diff_ratio:.2%}, 触发告警)这段脚本的逻辑说明先从 Hive 拿到当天的应覆盖用户总量再根据 rowkey 前缀从 Hbase 统计实际写入量偏差率超过 1% 就触发告警。这种做法比单独依赖 Sqoop 日志或 Hbase 的计数统计更可靠因为它是基于 SQL 语义和数据落盘结果的双向验证。参数上要关注两个点。一个是 Hbase 的 scan 超时设置用户量大时全表 scan 会很慢可以把 row_prefix 设计成按天加用户 id 哈希前缀这样每天的同步记录分布在固定前缀下scan 范围可控。另一个是校验的触发时机不要等 Hbase 同步作业完成之后立刻校验建议加一个三分钟的延迟确认 Hbase 的写入已经完成且 region 没有在分裂合并时再校验减少误报。这个方案不完美的地方在于临时表双写会多耗一倍 Hbase 存储空间。所以我在生产里不会天天用临时表只在版本升级、计算口径调整、Hbase 集群扩容后做一次全量双写校验。日常增量同步就靠状态表 数据量偏差告警两条线兜底。从那以后我每次搭画像系统都强制走一遍“Hive 明细落分区、校验表盯波动、同步链路加状态位、线下脚本做抽检”这条流程少一步心里都不踏实。数据量对不上这种问题一旦漏到线上排查成本是同步成本的几十倍。这套方案里最值得学的不是某个具体的 SQL而是那套“两种校验方案并行”的思路——宁可多写一个临时表也别让脏数据在线上裸奔。希望帮到你。本文还有配套的精品资源点击获取
阅读完成 · 觉得有帮助?
咨询建站