ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

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

用户画像标签存储四库协同:Hive、MySQL、Hbase、ES架构与避坑实践 简介面向大数据工程师与用户画像系统设计者的实战型技术方案聚焦标签数据存储这一核心环节系统讲解Hive、MySQL、Hbase、Elasticsearch四种数据库在画像系统中的定位与选型思路解决海量标签数据如何分层存储、高效查询与同步的难题。资源为1个PDF文档大小仅1.34MB内容围绕标签存储架构展开详解四种数据库的分工Hive承担离线批量计算与结果集存储MySQL管理元数据与结果校验Hbase服务在线高并发读取Elasticsearch支撑标签快速检索与人群分析同时涵盖用户标签表、标签聚合表、人群计算表等核心表结构以及基于Sqoop从Hive到Hbase、MySQL的跨库同步流程和两套数据一致性校验方案。读者可据此掌握不同数据库的适用场景理解离线批量计算与在线实时查询的存储分工并借鉴工程上的数据一致性校验思路。已有211人学习适合正在构建或优化用户画像平台的数据工程师、架构师参考。1. 用户画像系统解决方案标签数据存储不只是选库问题跑完一晚上 Spark 作业5000 万用户的标签结果落在 HDFS 上接下来怎么存、存在哪、怎么被线上圈人服务读到才是用户画像系统解决方案里最容易翻车的环节。很多人以为选个数据库把标签塞进去就完事实际上一套能支撑营收增长的画像系统标签数据存储是拆成四份的Hive 做底仓、MySQL 管元数据和校验、Hbase 扛线上服务、Elasticsearch 接多维透视。这篇就按这个分工把表怎么建、分区怎么设、同步怎么校验、踩过哪些坑一次讲透。适合正在搭画像 ETL、或者想把标签存储从单库改成多库协同的数据开发和数仓工程师。2. Hive 做标签底仓tag 表双层分区、tagmap 聚合与人群表 join2.1 Hive 在整个画像体系里的定位用户画像的数据链路里Hive 存的是所有标签相关数据的计算结果集。原因很直接跑用户标签的作业计算量非常大基本用 MapReduce 或者 Spark 作为执行引擎算完的结果落到 HDFS 上。这个过程在 MySQL、Hbase、MongoDB 这些数据库里是没法完成的——它们要么扛不住这个吞吐量要么根本不适合做批量计算结果的落地。所以 Hive 里存的是三类表用户标签明细表tag 表、标签聚合表tagmap 表、人群计算结果表。这三类表的分工完全不同tag 表记录每个用户每天被打上了哪些标签tagmap 表把分散的标签聚合到用户维度人群表则是按运营规则圈出来的用户集合。下面逐个拆。2.2 tag 表设计日期分区 标签类型分区先看 tag 表的核心字段标签 id、用户 id、标签权重。标签权重这个字段很重要它表示用户和标签之间的关联强度后面做用户分群、相似人群扩展都是靠它。表的分区设计是双层data_date数据日期和tagtype标签主题。为什么 tag 表要按日期分区这个好理解每天跑批数据按天隔离。关键是为啥还要按标签主题再分一层分区。原方案里写得明白为了 ETL 调度方便可以同时计算多个标签插入该表。如果没有 tagtype 这一层分区十几类标签作业同时跑全部写同一个日期分区作业之间会相互阻塞出问题还要全表扫描排查。有了 tagtype 二级分区每个标签作业只写自己的分区互不干扰。建表的 Hive DDL 大概是这样的CREATE TABLE dw.profile_tag_userid ( tag_id STRING COMMENT 标签id, user_id STRING COMMENT 用户id, tag_weight DOUBLE COMMENT 标签权重关联强度 ) PARTITIONED BY ( data_date STRING COMMENT 数据日期如20180421, tagtype STRING COMMENT 标签主题如userid_all_paid_money ) STORED AS ORC;这里分区字段data_date和tagtype是 Hive 的分区列不占用普通字段名额。实际查询时指定分区能直接走 HDFS 目录裁剪不扫全表。ORC 是列式存储压缩比高适合标签这种宽表场景。如果你的集群 Parquet 用的多换掉STORED AS ORC就行分区逻辑不变。向表里插入测试数据时一次插入同时指定两个分区值INSERT OVERWRITE TABLE dw.profile_tag_userid PARTITION (data_date20180421, tagtypeuserid_all_paid_money) SELECT tag_id, user_id, tag_weight FROM tmp_paid_money_result WHERE data_date 20180421;这里的意思是把当天算好的累计消费金额标签写入对应分区。INSERT OVERWRITE是覆盖写跑批作业天然幂等即使当天作业重跑也不会在分区里留下脏数据。跑批时每个标签作业提交自己的 Spark 任务写自己的 tagtype 分区不会互相抢锁。实际存储路径和分区一一对应hdfs://master:9000/root/hive/warehouse/dw.db/profile_tag_userid/data_date20180421/tagtypeuserid_all_paid_money这条路径把库名、表名、分区层级都暴露出来了。dw.db是库名profile_tag_userid是表名后面两级就是分区目录。看到这个路径就知道这张表的数据日期和标签类型排查问题时非常直观。关键点userid 维度和 cookieid 维度各做了一套一模一样的表。也就是说登录用户按 user_id 打标签未登录用户按 cookieid 打标签两套数据分开存储、分开聚合。做跨端行为打通时再把两套映射关系关联起来。2.3 tagmap 表把分散在分区的标签聚合到同一行tag 表解决了标签存储问题但引入一个新问题一个用户每天会被打上十几个标签散在不同 tagtype 分区下。运营想看一眼这个用户的全貌得跨全部分区去扫。这就轮到 tagmap 表出场。tagmap 表的思路是把同一个用户身上的所有标签聚合到一行。从dw.profile_tag_userid到dw.profile_user_map_userid这个转换过程实际是把多行标签记录合并成用户的标签集合。聚合的执行命令类似这样INSERT OVERWRITE TABLE dw.profile_user_map_userid PARTITION (data_date20180421) SELECT user_id, concat_ws(,, collect_set(tag_id)) AS tag_list, sum(tag_weight) AS total_weight FROM dw.profile_tag_userid WHERE data_date 20180421 GROUP BY user_id;这里collect_set把同一个用户的多个标签 id 收集成数组concat_ws再拼成逗号分隔的字符串sum(tag_weight)给出用户的标签权重总和。聚合后一行就是一个用户的完整标签集合查一个人的画像只需要扫一行数据。collect_set是去重的如果标签作业重复跑了多次同一个标签不会在聚合结果里出现两遍。如果你要保留标签的累计次数可以用collect_list配合count使用。注意聚合时想清楚是保留最新权重还是累加权重这直接影响后续做用户倾向性判断的结果。cookieid 维度同样有一套聚合逻辑两套表结构完全对称。做跨端打通时以 userid 维度的表为准cookieid 维度的标签按映射关系归并过来。2.4 人群表join 出手机号推给外呼中心人群表记录的是运营圈选的结果用户 id、人群名称 id、推送到的业务系统。人群表本身不存标签它存的是这批用户是哪个运营活动圈出来的、要推到哪个系统。实际落地时人群表要和订单表、收货信息表做 join才能把用户 id 变成可触达的手机号。执行过程大致是这样SELECT t1.user_id, t2.order_id, t3.phone FROM dw.user_group_info t1 JOIN dw.order_info t2 ON t1.user_id t2.user_id JOIN dw.user_receive_info t3 ON t2.order_id t3.order_id WHERE t1.group_id 20240401_act001 AND t1.data_date 20240401;这段 SQL 里人群表 join 订单表拿到订单编号订单表再 join 收货信息表拿到手机号。这样圈定人群之后就能得到一批可外呼的运营用户和他们的手机号推给外呼中心做电话触达。这里有个隐蔽的坑如果一个用户在活动期间下了多笔订单join 出来的手机号会重复后面避坑章节专门说。Hive 这部分的核心结论tag 表按日期和标签类型双层分区解决了大批量跑批的写入问题tagmap 表把分散标签聚合到用户维度解决了查询效率问题人群表 join 订单和收货信息解决了业务触达问题。三层表从存储到应用构成了完整的 Hive 底仓。3. MySQL 在画像里的三个角色元数据、校验位、业务回读3.1 MySQL 不存明细只存三类数据Hive 存了几亿行标签明细MySQL 存什么答案是三类数据画像标签的元数据、结果集的校验信息、同步到业务系统的数据。原方案里写得很清楚MySQL 在画像项目中不承担明细存储它做的是轻量级、高并发、需要事务保障的活。先理解为什么 MySQL 干不了明细存储用户标签表几百亿行MySQL 单表几千万行就开始吃力加上 JOIN、聚合查询会更难。但 MySQL 的强项是点查快、支持事务、锁机制完善。所以画像系统里凡是要做校验判断状态标记配置读取的事情都在 MySQL 里做。3.2 画像标签的元数据运营端标签树的字典元数据表维护的是标签本身的描述信息标签 id、名称、主题、一级分类、二级分类、标签描述。运营人员在画像产品化的界面上看到的标签树、标签说明全部是从这张表读的。没有这张表标签 id 就是黑匣子运营看到的只是一串数字不知道这个标签代表什么。元数据表的设计上有个实操细节一级分类和二级分类不要用自增 id建议直接用有业务含义的编码。因为标签树在运营端是层级展示的编码本身带层级关系查询和排序都方便。而且标签上线、下线、变更描述都算配置变更建议在表里加一个status字段控制上下线不要物理删除。提示元数据表的数据量很小几百上千条但它被读取的频率极高。运营端每次打开标签树、每次圈人筛选都要查。这类表一定要走 MySQL 缓存索引建在分类字段上别建在描述字段上。3.3 结果集校验量级监控与标志位这是 MySQL 在画像里最容易被忽视、但最不能省的角色。Hive 作业每天跑完怎么知道结果对不对靠 MySQL 里的校验表。校验表记录三类指标当日该标签覆盖的用户量当日该标签与昨日相比的波动比例当日该标签覆盖的用户占当日活跃用户的比例这三个指标分别回答三个问题这个标签今天出了多少数据和昨天比有没有剧烈波动在整体活跃用户里的渗透率是否正常波动比例的阈值一般设在 ±20%超过就报警。比如某个标签昨天覆盖 1000 万用户今天突然只有 200 万那大概率是作业挂了或者上游数据出了岔子。更关键的是标志位机制。校验表里放一个标志位用于判断某些任务是否需要继续执行。比如 Hive 到 Hbase 的同步作业先更新校验表如果校验通过就置为 1同步程序看到 1 才继续往下跑校验失败置为 0下游任务不执行避免把错误数据推到线上。这个设计把数据质量检查和任务调度解耦了。调度系统不用关心数据质量只看标志位。校验规则后续调整只需要改校验程序调度链路完全不用动。3.4 同步到业务系统Python 脚本和 Sqoop 二选一客服系统用的是关系型数据库 MySQL画像系统圈定用户后需要把待运营的用户推送到客服系统。数据从 Hive 到 MySQL常见做法是 Sqoop也可以写 Python 脚本拉取。我一般会直接用 Python 脚本因为可以灵活处理字符集、类型转换和增量逻辑。脚本思路是这样的import pymysql from pyhive import hive # Hive 连接参数 hive_conn hive.Connection( hosthive-server, port10000, usernameetl_user, databasedw ) # MySQL 连接参数 mysql_conn pymysql.connect( hostmysql-server, userapp_user, password******, databaseuser_profile, charsetutf8mb4 ) # 从 Hive 读取人群结果 hive_cursor hive_conn.cursor() hive_cursor.execute( SELECT user_id, tag_id, tag_weight FROM dw.profile_tag_userid WHERE data_date 20240401 AND tagtype userid_all_paid_money ) # 批量写入 MySQL mysql_cursor mysql_conn.cursor() batch_size 5000 batch [] for row in hive_cursor: batch.append(row) if len(batch) batch_size: mysql_cursor.executemany( INSERT INTO app_user_tags (user_id, tag_id, tag_weight) VALUES (%s, %s, %s) ON DUPLICATE KEY UPDATE tag_weight VALUES(tag_weight), batch ) mysql_conn.commit() batch [] if batch: mysql_cursor.executemany( INSERT INTO app_user_tags (user_id, tag_id, tag_weight) VALUES (%s, %s, %s) ON DUPLICATE KEY UPDATE tag_weight VALUES(tag_weight), batch ) mysql_conn.commit() hive_cursor.close() mysql_cursor.close() mysql_conn.close()脚本逻辑分三步先从 Hive 读出当天某标签的明细结果再按 5000 条一批写入 MySQL最后用ON DUPLICATE KEY UPDATE处理重跑时的幂等。批量大小这个参数值得调MySQL 单次批量 insert 太大容易锁表超时太小效率低实际业务里 3000 到 5000 是甜点区间。用 Sqoop 的话就是一条命令的事情适合走定时调度sqoop export \ --connect jdbc:mysql://mysql-server:3306/user_profile \ --username app_user --password ****** \ --table app_user_tags \ --export-dir /user/hive/warehouse/dw.db/profile_tag_userid/data_date20240401 \ --input-fields-terminated-by \001 \ --update-mode allowinsert \ --update-key user_id--update-key user_id指定主键去重--update-mode allowinsert表示已存在就更新、不存在就插入。注意 Hive 默认字段分隔符是\001Sqoop 导入时要写清楚不然字段全会串列。同步完成后别急着走人回 MySQL 数一下行数和源表对比这个习惯能挡掉一大半数据事故。4. Hbase 存储与圈人服务映射建表、调度链路与两套校验方案4.1 从 Hive 映射 Hbase一个建表语句打通两层Hbase 在画像系统里承担的角色是线上服务广告系统、push 消息系统直接读它。为什么不用 Hive 直接提供线上服务因为 Hive 查询延迟在秒级到分钟级线上广告和 push 根本等不起。Hbase 按 rowkey 点查是毫秒级正好匹配这个场景。Hive 数据同步到 Hbase最常见的做法是建一张 Hive 映射表让 Hive 能直接读写 Hbase 表CREATE TABLE dw.hbase_tag_sync ( user_id STRING, tag_id STRING, tag_weight DOUBLE ) STORED BY org.apache.hadoop.hive.hbase.HBaseStorageHandler WITH SERDEPROPERTIES ( hbase.columns.mapping :key, cf:tag_id, cf:tag_weight ) TBLPROPERTIES ( hbase.table.name profile:user_tags );这里hbase.columns.mapping是核心参数:key表示 Hbase 的 rowkey 用user_id剩下的字段映射到列族cf下的tag_id和tag_weight两列。cf是列族名可以改成你集群上实际的名字。hbase.table.name指定 Hbase 里的真实表名不写的话默认和 Hive 表同名。建好映射表后往映射表插入测试数据INSERT OVERWRITE TABLE dw.hbase_tag_sync SELECT user_id, tag_id, tag_weight FROM dw.profile_tag_userid WHERE data_date 20240401 AND tagtype userid_all_paid_money;这条语句触发的是一个 MapReduce 作业把 Hive 分区里的数据批量写入 Hbase 表。作业执行过程中不需要额外的同步工具Hive 和 Hbase 直接通过 HBaseStorageHandler 交互。写完直接查 Hbasescan profile:user_tags, {LIMIT 10}rowkey 设计决定查询性能。用user_id做 rowkey优点是单个用户的所有标签在同一行点查一个用户毫秒级返回但要注意热点问题如果 user_id 是递增数字写入会集中在少数 Region 上。实际工程上可以在 user_id 前面加盐比如String.format(%02d, hash(user_id) % 100)拼上去让数据均匀分布在各个 Region。4.2 圈人服务的离线与实时链路Hbase 在工程上最典型的应用场景是圈人服务。整体流程是这样运营人员根据规则圈出人群及对应的标签集先存入 MySQLSpark 作业读取 MySQL 的标签集信息计算对应的人群计算得到的人群写入 Hive 分区表最后每天或实时将该分区数据同步到 Hbase供广告系统、push 消息系统读取。这个流程可以做成离线模式走 ETL 做 T1也可以做成实时模式实时计算出对应用户群推到 Hbase。两种模式的区别主要在 Spark 作业的触发方式离线走定时调度实时走消息队列或阿里云 DataWorks 的实时节点。工程上有一个不能省的环节因为灌入 Hbase 的数据直接应用在线上、直接反馈给用户所以 Hive 同步到 Hbase 时必须做校验防止同步过程中出问题——比如 Hive 数据有 5000 万条同步到 Hbase 只剩 1000 万条这种情况一旦发生就是线上事故。4.3 两种校验机制临时表重命名与状态表标志位原方案给了两种解决方案我都实际落地过分别说下适用场景。方案一临时表重命名。Hive 到 Hbase 同步数据后先在 Hbase 建立一张临时表校验临时表和对应 Hive 表的数量差异。差异在可接受范围内把临时表重命名为正式表。这个方案的优点是正式表永远只有完整数据线上不会读到半成品缺点是重命名那一刻有短暂的服务空窗期另外要维护双份表空间。# 伪代码同步到临时表 hive -e INSERT OVERWRITE TABLE dw.hbase_tag_sync_temp SELECT ... # 校验数量 hbase org.apache.hadoop.hbase.mapreduce.RowCounter profile:user_tags_temp # 校验通过后重命名 hbase shell EOF disable profile:user_tags alter profile:user_tags_temp, {NAME cf}, tableName profile:user_tags EOF方案二状态表标志位。Hive 同步后直接把数据写入正式表同时在 Hbase 建立一张状态表用于记录状态位。校验正式表和 Hive 表的数量差异在可接受范围内才写入状态表。接口请求时只读取状态表里最近日期的记录如果状态位不存在就说明当前数据未就绪拒绝读最新数据线上继续用前一天的。这个方案的巧妙之处在于把数据同步和数据可用彻底分开。同步程序只管写校验程序只管标记状态接口按状态位取数。Hbase 同步异常时不会写入状态表也完全不影响线上数据的读取——因为接口读的永远是状态位指向的那张表状态位没更新接口就还是读旧数据。两套方案的实际取舍临时表方案适合数据量适中、重命名窗口能够接受的服务状态表方案适合 7×24 小时在线、完全不能接受服务中断的业务。我自己在做 push 消息系统时用的状态表方案因为推送服务凌晨也在发消息任何服务空窗都是不可接受的代价。注意不管用哪种方案校验的可接受范围不要拍脑袋写死。常见做法是设置百分比阈值比如差异率在 1‰ 以内放行超过就告警置失败。直接比较绝对数值在数据量线性增长阶段会频繁误报这个参数建议做成配置项而不是写死在代码里。5. 上线前必须过的避坑清单标签量级、同步丢数与状态位失效5.1 Hive 同步 Hbase 后少了几千万行现象Hive 表里 5000 万条数据同步到 Hbase 后 scan 只有 1000 万条量级差了五倍。原因最常见是 rowkey 冲突导致数据互相覆盖。多个标签作业同时往同一个 rowkey 写数据Hbase 按列族内最新时间戳取值后写的覆盖先写的。另一个常见原因是 Hive 源表里有重复 user_id同一 user_id 对应多条标签记录写入 Hbase 时同一行被反复覆盖。解决先看源表有没有去重SELECT count(*) FROM (SELECT DISTINCT user_id FROM ...) t对比总行数和 distinct 行数如果源表有重复同步前先按 user_id 聚合。同步完成后用 RowCounter 对比 Hbase 实际行数和 Hive 去重行数差异在阈值内才算通过。从那以后我对每个同步作业强制加这一步没再被这个坑绊倒。5.2 量级一致但用户 ID 错位现象Hive 和 Hbase 行数完全一致校验也通过了但线上查某个用户时发现他身上的标签是另一个人的。原因同步程序是按分区扫描的但写入 Hbase 时 rowkey 拼接出了问题。比如 user_id 是 string写入时Bytes.toBytes(user_id)错误地用了user_id 或者中间经历了类型转换导致数字被截断两个不同 user_id 映射到了同一个 rowkey。解决校验时不要只看行数要做用户维度的抽样比对。hbase 里随机取 1000 个 rowkey回到 Hive 里查对应 user_id 的标签逐个字段比对。这个步骤我在同步脚本里用一个临时表实现每天自动跑抽样结果落 MySQL。5.3 状态位没按最近日期过滤接口读到过期数据现象状态表里有多天记录接口按固定日期去查某天同步失败后没写状态位接口读到的还是几天前的数据运营看到的用户画像明显过期。原因状态表只写了当天是否成功这个标志位接口查的时候没用MAX(data_date)过滤而是写死了日期参数。调度系统那天没跑同步接口依然按当天的日期去查自然读不到记录。解决所有接口查状态表时必须按日期倒序取第一条确保读到的永远是最新可用的数据。状态表里加一个is_ready字段只有校验通过才置 1接口端过滤is_ready 1。5.4 人群表 join 订单表手机号重复推送现象外呼中心反馈同一个人被呼叫了五六次用户投诉。原因人群表 join 订单表是一对多关系——一个用户下过多笔订单join 出来就是多行再 join 收货信息表每笔订单都带出一个手机号自然重复。SQL 本身没报错但业务上触达重复了。解决join 完成后按 user_id 去重保留一条记录。用ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY order_time DESC) rn取 rn1或者更简单粗暴的做法直接SELECT DISTINCT user_id, phone。前者能控制取哪条后者代码更短实际场景我优先用ROW_NUMBER因为还能顺便过滤掉无效手机号。5.5 MySQL 元数据没跟上运营端标签树断层现象新增了一批标签Hive 里数据跑完了运营端标签树里却看不到圈人功能用不了新标签。原因标签元数据是配置数据新增标签需要走一步元数据同步。但这个步骤经常被遗忘因为 Hive 作业不依赖它也能跑效果是延迟暴露的——等到运营要用才突然发现。解决把元数据同步放在 ETL 链路的最前面并和标签作业解耦。新增标签时先在 MySQL 元数据表插入记录Hive 作业启动前先检查元数据表是否有对应标签 id没有就直接失败。这个习惯等于给标签配置上了保险漏配的情况在作业阶段就会被拦截而不是拖到运营使用阶段。问题现象原因解决Hbase 丢数据5000 万变 1000 万rowkey 冲突、源表未去重源表去重 RowCounter 校验用户标签错位行数一致但内容错rowkey 拼接或类型转换出错抽样比对明细字段状态位失效接口读到过期数据未按日期倒序过滤状态表强制MAX(data_date)过滤手机号重复同一用户被重复外呼join 一对多未去重ROW_NUMBER()按用户去重标签树断层运营端看不到新标签元数据同步被遗忘作业启动前强校验元数据6. Elasticsearch 接住多维透视同步验证与查询技巧6.1 ES 在画像里的定位按标签条件反查人群Hbase 擅长的是给定 user_id查出这个人的全部标签但运营经常要反过来问过去 30 天浏览过母婴类目、且消费金额在 5000 以上的用户有哪些这种按标签组合条件反查人群的场景Hbase 就吃力了它没有倒排索引只能全表扫描。Elasticsearch 的分布式倒排索引正好解决这个问题用户标签查询、人群多维透视分析这类对响应时间要求高的场景我会优先把数据同步到 ES。ES 索引设计上有个关键点doc 的_id直接用 user_id标签字段全部设为keyword类型。这样单用户点查走_id人群筛选走 terms 组合查询两种查询模式都很快。千万别把标签设成text类型会被分词器拆得乱七八糟精确匹配直接失效。6.2 同步 ES 后的三面对数验证Hive 同步到 ES 后我会强制走一遍三面对数Hive 源表的 count、ES 的_count、MySQL 校验表的记录数。三面对不上就报警绝不直接放给线上。# 查看 ES 索引总文档数 curl -X GET http://es-server:9200/user_profile_tags/_count # 按标签组合条件查人群 curl -X GET http://es-server:9200/user_profile_tags/_search -H Content-Type: application/json -d { query: { bool: { filter: [ {term: {tags: userid_all_paid_money}}, {range: {paid_amount: {gte: 5000}}} ] } }, size: 0 }这个查询里term精确匹配标签range过滤消费金额size: 0表示只要聚合结果不要明细。返回的total就是满足条件的人群规模运营圈人前先看这个数比在 Hive 里跑一遍快几个数量级。四类数据库的分工这时候就完整了Hive 管大批量计算结果集的批处理MySQL 管元数据、校验、业务回读Hbase 管线上 kv 点查ES 管多维组合检索。每类数据库只干自己最擅长的那件事。从那以后我每次同步 ES 都强制走一遍三面对数再加状态位写入这套校验流程挡掉过好几次同步脚本的隐性 bug希望帮到你。本文还有配套的精品资源点击获取
返回列表