数据仓库
项目中用到的源数据怎么来的?数据的吞吐量是多少,每秒有多少条数据(QPS)¶
springboot 程序实时生成,模拟真实的流数据,用配置文件来控制随机范围 springboot 程序模拟的数据分成业务数据和日志数据两大块,分别写入 mysql 和 日志文件 业务数据又分成全量同步和增量同步。其中日志数据和增量业务数据进入实时数仓 日志数据:通过 Flume 监控日志文件,将数据同步到 kafka,topiclog 增量业务数据:通过 maxwell 将数据从 mysql 同步到 kafka,topicdb 这两个 tpoic 就作为实时数仓的ODS层做一个原始数据的保留 2000QPS,模拟的数据,有一个波动的范围,日志数据有50到60G,业务数据有1到2G
实时数仓项目的数据流转是怎样的?(介绍一下实时数仓的构成和数据流转)¶
数据从业务端过来分成业务数据和日志数据两大块,业务数据又分成全量同步和增量同步,其中日志数据和增量业务数据进入实时数仓。 - 日志数据:通过 Flume 监控日志文件,将数据同步到 kafka ;增量业务数据:通过 maxwell 将数据从 mysql 同步到 kafka ,上面这些原始数据作为ODS层 从 kafka 读出数据经过 Flink 程序的加工后, 一部分写回 kafka 作DWD层,其余维度数据写到 HBase 作DIM层 之后这两层的数据会经过 Flink 程序读取出来加工后写入 Doris 作为DWS层, 其中维度数据还会经过 Redis 缓存一天的时间,缓解热点 key 问题。 最后 Doris 的数据经过 SpringBoot 数据接口服务后,得到最终指标数据 成为 ADS层
实时数仓有哪些表¶
- ods :日志表、业务表(odsxxxinc/full)
- dim :商品、优惠券、活动、地区、日期、用户维度表
- dwd:分域
- 用户域:登录、注册
- 流量域:启动、页面、动作、故障、曝光
- 交易域:加购、下单、支付、物流
- 工具域:领取优惠卷、使用优惠卷下单、使用优惠卷支付
- 互动域:点赞、评论、收藏
- dws:最近 1 / 7 / 30 日、历史至今的汇总表
- ads :日活、新增、留存、转化率、GMV等
用户域有哪些表¶
维度表:dimuserinfo(用户维度表,存储用户基础信息) 事实表:dwduserlogin(用户登录事实表,存储用户登录事件),dwduserregister(用户注册事实表) 用户相关业务源表:userinfo(用户主表,Maxwell 同步至 Kafka,分流写入 dimuser_info)
交易,登录建了多少表,放在哪一层¶
一、交易域 1. ODS 层 Kafka 主题:topicdb(MySQL 业务全量 / 增量 binlog 原始数据,包含所有交易业务原始表) 原始 MySQL 交易业务表:activityinfo、activityrule、activitysku、couponinfo、couponrange、skuinfo、spuinfo、financialskucost、订单 / 支付 / 退款类表 2. DIM 层(维度表,存储 HBase,表前缀 dim) 全部交易相关维度,共 14 张: dimactivityinfo、dimactivityrule、dimactivitysku dimbasecategory1、dimbasecategory2、dimbasecategory3 dimbasetrademark dimcouponinfo、dimcouponrange dimfinancialskucost dimskuinfo、dimspuinfo dimbaseprovince、dimbaseregion 3. DWD 层(事务事实表,Kafka,前缀 dwdtrade) 交易域业务过程:加购、下单、取消订单、支付成功、退单、退款成功 对应 6 张 DWD 表: dwdtradecartadd 加购 dwdtradeorderdetail 下单明细 dwdtradeordercancel 取消订单 dwdtradeorderpaymentsuccess 支付成功 dwdtradeorderrefund 退单 dwdtraderefundpaymentsuccess 退款成功
二、用户域 用户域业务过程:注册、登录 1. ODS 层 原始 MySQL 表:userinfo,同步到 Kafka topicdb 2. DIM 层(维度表 HBase) 1 张:dimuserinfo(用户全量维度) 3. DWD 层 dwduserregister,dwduserlogin
项目中的分区表怎么划分的¶
由于 DWS 层所有汇总表都写入 Doris,且需求是“实时统计当日数据并展示”,项目采用了 动态分区(Dynamic Partition) + Range 分区 具体划分是按天进行划分,保存昨天及以后的数据,节省存储空间;提前创建未来3天的分区,避免分区不存在导致写入失败。
DWS的颗粒度是什么?¶
DWS层的粒度是指数据汇总的层次级别,它是基于DWD明细数据按照特定维度组合进行聚合的层级。 常见的粒度包括时间维度(如日、周、月)、业务维度(如用户、商品、地区)以及它们的组合(如用户+日粒度)。 具体设计取决于业务分析需求,需要在查询性能和分析灵活性之间取得平衡。例如在 Doris 表里一行数据表示的是:一个窗口的时间区间内,各个省份下,订单的数量及金额
挑一个最熟悉的讲一下,详细到每层数据怎么来、表都有哪些字段¶
用 最近7/30日各品牌复购率 这个指标举例,属于交易域,与订单相关
数据流如下:
- ODS层:订单、订单明细、订单明细活动关联、订单明细优惠卷关联
- 字段:type变动类型,ts变动时间、data数据、old旧值
- DIM层:商品相关信息的维度表,如省份、分类
- DWD层:下单事务事实表
- 字段:订单、用户、商品、省份、活动、优惠卷等的id
- DWS层:用户商品粒度 订单 最近1日汇总表 => 最近n日汇总表
- 字段:用户id、sku_id和名称、各级分类id和名称
- ADS层:最近7/30日各品牌复购率
- 字段:dt日期、recent_days(7或30)、品牌id和名称、复购率
什么是旁路缓存¶
外部数据源的查询常常是流式计算的性能瓶颈。以本程序为例,每次查询都要连接 HBase,数据传输需要做序列化、反序列化,还有网络传输,严重影响时效性。可以通过旁路缓存对查询进行优化。
旁路缓存模式是一种非常常见的按需分配缓存模式。所有请求都会优先访问缓存,若缓存命中,直接获得数据返回给请求者。如果未命中则查询数据库,获取结果后,将其返回并写入缓存以备后续请求使用。
应当注意两点问题,第一,缓存要设过期时间,不然冷数据会常驻缓存,浪费资源。第二,要考虑维度数据是否会发生变化,如果发生变化要主动清除缓存。
缓存的选型一般会考虑两种: 堆缓存 或者 独立缓存服务(memcache,redis) 堆缓存,性能更好,效率更高,因为数据访问路径更短。但是难于管理,其它进程无法维护缓存中的数据。 独立缓存服务(redis,memcache),会有创建连接、网络IO等消耗,较堆缓存略差,但性能尚可。独立缓存服务便于维护和扩展,对于数据会发生变化且数据量很大的场景更加适用,项目中选择独立缓存服务,将 redis 作为缓存介质。
具体的实现步骤: 1. 查询时从缓存中获取数据。 2. 如果查询结果不为null,则返回结果。 3. 如果缓存中获取的结果为null,则从HBase表中查询数据。 4. 如果结果非空则将数据写入缓存后返回结果。 5. 否则提示用户:没有对应的维度数据
注意:缓存中的数据要设置超时时间,本程序设置为1天。此外,如果原表数据发生变化,要删除对应缓存。为了实现此功能,需要对维度分流程序做如下修改: 维度变更时 如果维度数据的变更类型为insert,则对缓存无影响。 如果维度数据的变更类型为update或delete,则清除缓存。
项目中遇到的挑战性的问题(旁路缓存加异步IO)¶
引入 Redis 作为独立旁路缓存,是为了解决 DWS 层做大表关联(如 SKU 关联数十张维度表)时的性能瓶颈。 - 问题:dws层数据需要将dwd层数据读出来与dim层的维度数据进行join,如果每次都直接查询 HBase,会经历网络 IO、序列化/反序列化,导致吞吐量急剧下降,成为流计算的性能短板
-
为什么选择redis作为旁路缓存: 本地堆缓存(如 Caffeine):虽然快,但无法跨进程管理,且 DIM 层维度数据变更时无法主动通知 DWS 层清除旧缓存。 Redis(独立缓存服务):数据常驻内存,访问延迟极低;并且支持跨任务通信——DIM 层发现维度变更时,可以主动删除 Redis 中的 Key,保证 DWS 层缓存数据的强一致性。
-
旁路缓存: 设置redis作为旁路缓存,查询时会先查redis,redis没有再去查Hbase,且查完之后在redis中加上刚刚查询到的kv值,并设置 TTL 过期时间 DIM 层的缓存失效来保证数据的一致性,如果维度信息发生改变,则必然同步删除redis中旧的维度信息
-
异步IO: 为了不让同步阻塞影响效率,项目使用了Flink的异步IO,编排异步任务。 先异步查询redis,若未命中则异步查询Hbase,查完后异步写回redis,最后返回结果。 这样做使得单个并行子任务可以连续发送多个请求,按照返回的先后顺序对请求进行处理,发送请求后不需要阻塞式等待,省去了大量的等待时间,大幅提高了流处理效率。
刚刚说到异步IO,再详细讲一下¶
因为明细数据可能需要与多个维度数据 join,所以用模板方法的设计模式写好异步IO的函数,从线程池中获取线程,发送异步请求。再到 flink 程序中用 AsyncDataStream.unorderedWait()方法将原本join的串行处理变成并行
为什么要把DIM层存HBase¶
- dim层存放的是维度数据,不能像实时数据流一样存在kafka里,要做长期保存
- 考虑到维度数据的特性和用法是按照key找对应的value,所以考虑使用KV数据库,而熟悉的KV数据库有redis和Hbase
- 因为用redis长期存全部维度信息对内存的压力比较大,所以采用Hbase
拉链表怎么设计的,介绍一下你的拉链表怎么设计的?更新逻辑讲一下。你认为为什么要用拉链表来存这个数据呢¶
在用户维度表运用了拉链表来存储维度数据。 在后面加两个字段用于表示该状态的开始时间和结束时间,没有结束就用9999-12-31表示
變一下你理解的拉链表,你实际怎么运用拉链表的?¶
在记录用户信息的维度表时用到了拉链表 拉链表适合于:数据会发生变化,但是变化频率并不高的维度(即:缓慢变化维) 比如:用户信息会发生变化,但是每天变化的比例不高。如果数据量有一定规模,按照每日全量的方式保存效率很低。比如:1亿用户*365天,每天一份用户信息。(做每日全量效率低)
优点:(1)保留了数据的历史信息;(2)节省存储空间; 缺点:同步和回滚逻辑复杂;
实际运用: 在建表时添加两个字段,分别是开始日期和结束日期,用来表示状态的持续时间,如果是当前的最新状态就用日期最大值9999-12-31表示。 如果要获取某个日期的历史切片的话,通过,生效开始日期<=某个日期 且 生效结束日期>=某个日期,就能够得到某个时间点的数据全量切片。 分区按照过期时间,同一天过期的放在同个分区,最新的数据放在一个分区。
拉链表新数据的更新是怎么样的,装载数据那条sql的思路?¶
首日装载:type='bootstrap-insert',全部放到最新分区 每日装载: 用 截至前一天最新数据的 full join 当天的新增及变化,就是 dim 的 9999-12-31 分区 和 ods 的当日分区 然后利用if函数判断,新数据不为null则用新数据,其余不变,这样得到的是截至当天最新 最后还要 union all 一个旧的数据,并更新 end_date,insert overwrite 的时候要根据 dt 分区
DWD层使用redis时是如何保证缓存与HBASE一致性¶
- 基于事务的两阶段提交:在更新HBase数据时,同时将相关操作记录到一个事务日志表中。然后,Redis的更新操作作为事务的第二阶段。如果Redis更新成功,事务提交;否则,回滚HBase的更新操作。
- 利用消息队列:当DWD层数据发生变化时,将变更消息发送到消息队列。消费者从消息队列中获取消息,先更新HBase,再更新Redis。如果更新Redis失败,可以将消息重新放入队列,进行重试。
- 定期同步:定时任务定期从HBase中读取数据,并与Redis中的数据进行对比。对于不一致的数据,以HBase中的数据为准,更新Redis。这种方式存在一定的时间窗口内数据不一致问题,但实现相对简单。
- 数据版本控制:在HBase和Redis中都记录数据的版本号。当DWD层数据更新时,版本号递增。更新Redis时,先检查版本号,如果版本号不一致,说明数据已过时,需要从HBase重新获取最新数据并更新Redis。