跳转至

数据仓库

数仓分层

五层架构 数据流向:ODS(原始数据)→ DWD(清洗明细)→ DIM(维度关联)→ DWS(公共汇总)→ ADS(应用指标) 1. 操作数据层(ODS - Operational Data Store):原始数据,不做修改 2. 明细数据层(DWD - Data Warehouse Detail):对 ODS 层数据进行清洗、过滤、维度退化,按照业务过程建模,保存各业务过程最小粒度的明细事实数据。 3. 维度层(DIM - Dimension):基于维度建模理论构建,存放一致性维度表,保存维度对象的描述属性信息,供事实表关联使用。 4. 汇总数据层(DWS - Data Warehouse Service / Summary):基于上层应用需求,以主题对象为建模驱动,对 DWD 明细数据按照不同维度进行聚合汇总,生成公共统计粒度的汇总表,供上层ADS复用。 5. 应用数据层(ADS - Application Data Service):存放各业务方最终需要的统计指标结果,直接供外部应用(报表、大屏、接口)查询使用。

数仓分层目的

  1. 简化复杂需求,每一层只处理单一的步骤
  2. 以空间换时间,提高数据的复用性
  3. 通过规范数据分层,开发一些通用的中间层数据来减少重复开发,提高开发效率
  4. 让数据结构更加清晰,方便数据的定位和理解
  5. 可以用于数据血缘追踪,当数据出现问题后,可以快速定位问题并清楚它的危害范围
  6. 可以通过分层管理来锁定数据位置,快速进行修复,隔离故障避免对其他数据表的污染
  7. 架构弹性化,业务变更时会通过视图保持兼容,屏蔽业务变更的影响

数据仓库建模的方法和优缺点

主要有两种,ER模型,适合按主题分类整合数据,但是不能直接用于分析决策,好处是冗余更少,缺点是在大规模数据跨表分析中容易造成多表关联问题,大大降低执行效率【优化方案是Data Vault模型和Anchor模型】。 另一种就是维度模型,目的是为分析决策的需求服务,也是主流的数据建模方法,维度模型主要分为三类,最常用的是星型模型,核心是以事实表为中心,所有的维度表会直接连接在事实表上(类似星星);衍生出了雪花模型,不同的地方就在于雪花模型的维度表可以再连接其他的维度表(类似与MySQL的3NF模型);星型模型还衍生出了星座模型,不同于其他两种方法,它基于多张事实表并且会共享维度信息(因为,企业中不会只有一张事实表,所以一般就是使用星座模型来进行维度建模)

维度建模中表的类型

维度建模中表的类型分为两种,第一种是维度表,主要是对业务对象的描述信息(例如乘客信息表,司机信息表之类),形式上包含单一的主键列和一些对该主键的描述信息,也因此维度表会很宽;第二种是事实表,主要是对业务过程的描述(例如下单,支付之类),每个事实表都包含了数值型的度量值,维度外键和退化维度,退化维度的属性被存储到事实表中,减少关联,通常事实表会比较大

什么时候会退化维度

当事实表中包含的维度外键数量过多,或者维度外键的关联关系比较复杂的时候,就会退化维度。退化维度的目的是为了减少关联,提高查询效率。例如,一个订单表中包含了多个维度外键,每个维度外键都指向一个维度表,那么就会退化维度,将维度表中的属性存储到事实表中,减少关联。

数据模型怎么建设的

数据建模的流程:自底向上和自顶向下 。自顶向下就是从业务需求到数据模型再到具体实现,用来确定数据的口径、字段定义 ;自底向上是从具体数据到维表再到主题域,构建维表进行数据维度的划分。例如数仓五层架构的构建,整体架构就是自顶向下设计的,但是具体到维度表和事实表的构建又可以采用自底向上的方式逐步迭代。

业务过程和主题域的关系

一个主题下可以有多个业务过程,但是一个业务过程一般只属于某一个主题

事实表的类型

事实表有三种类型:事务事实表、周期快照事实表、累积快照事实表

事实表的设计过程

一共有五步,分别是选择业务过程,声明粒度,确定维度,确定事实,冗余维度

首先是选择业务过程,就是通过对业务流程的整个生命周期进行分析,然后,选择与目标需求有关的业务过程(每一个业务过程可以概括为一个个不可拆分的行为事件),举个例子,滴滴打车的整个业务流程:从乘客点单,平台发布,司机抢单,司机接单,完成订单。就是一个完整的业务过程。而选择业务过程的表现就是根据我们的需求在其中去选择对应的过程(例如选择了乘客点单这个业务过程)

其次是声明粒度,粒度就是确定事实表中一行所表示的业务的细节层次,通常推荐在设计事实表的时候,粒度定义的越细越好,例如交易域下单事务事实表,其对应的业务过程就是 下单,其粒度就是 一笔订单中的一个商品项

然后是确定维度,即确认清楚当前业务过程所涉及到的维度信息。例如交易域下单事务事实表中的:商品、用户、省份、价格、订单状态等多个维度,通俗地说就是“何人在何时何地干了什么”

接下来就是确认事实,即确认当前业务过程中所涉及到的可度量的值,例如订单量,订单金额等,通俗地讲就是“发生了多少,金额是多少”

最后是冗余维度,也叫退化维度。目的是在事实表中冗余一些事实表中使用的常用维度,减少多表之间的关联,提高查询效率

如何确定事实表需要的字段

其实就是事实表设计过程的后三步,声明粒度是确认主键,确定维度是确认外键,确定事实就是确认度量值对应字段

维度表的构建方式

第一种,正常规范维度表,将大表拆分为小表,减少数据冗余,确定是多表连接造成效率低下 第二种,反规范化维度表,将部分维度信息反规范化设计,减少表之间的join,构建宽表,冗余描述性信息,提高查询性能 第三种,在反规范化的基础上构建拉链表,通过在维度表中新增生效日期和失效日期字段来保存状态,用于历史追踪,观察数据怎么变化

维度表的设计原则

第一,维度属性越多越好 第二,维度的值尽可能使用更多的文字描述而不是编码 第三,维度属性被越多表引用越好

ODS层存在的意义

  1. 提供多源数据的统一视图,解决不同源数据的冲突问题
  2. 作为数据仓库的第一层,起到数据缓冲的作用,避免对上层数据仓库的直接冲击,确保数据仓库的稳定和性能
  3. 数据结构与业务系统相对更接近,能够更好地满足业务部门对数据的灵活需求,提升数据应用的效率

举例说明怎么保证数据写入到DWS的幂等性的(怎么保证数据一致性)

[实时] 以 交易域SKU粒度下单各窗口汇总表 为例说明,数据来源是kafka中DWD层的订单明细数据 在 flink 程序中,先由 string 转成 JSONObject,根据唯一键去重后,再转成对应的实体类。 之后设置水位线,并根据省份id分组开窗聚合。最后关联省份维度信息补全字段,写出到 ClickHouse

首先,依赖于Doris的幂等性,对应数据的下游存储即DWS层的目标表,是Doris存储层,在建表的时候使用了AGGREGATE KEY 模型,并将度量列指定为 REPLACE。当Flink 写入的时候,Doris 会根据 AGGREGATE KEY(如 stt、edt、sku_id 等)进行聚合。如果 Flink 因故障重试,重复下发同一条数据(相同 Key),Doris 不会报错或插入重复行,而是直接用新数据替换旧数据。这就保证了即便下游收到重复写入请求,最终表中的数据也是确定且唯一的。

其次,提供计算引擎层(Flink)的一致性:使用抵消机制(Retract)精准去重,按照主键分组,然后使用Keyed State缓存上一条数据,那么中间结果来的时候就输出一次,完整结果来的时候先抵消上一条数据,再累加当前数据。通过先抵消再累加的机制确保了即使流中存在变更,窗口内聚合出的度量值也始终与业务数据库最终状态保持一致。

最后也是最重要的一点就是端到端的精准一次(Exactly-Once):Flink Checkpoint + 两阶段提交。为了保证数据不重不丢,项目启用了 Flink 的 Checkpoint,并配合 Sink 端的事务或幂等特性。对于Kafka Source来说,开启Checkpoint 后,Flink 可以重置偏移量,保证故障重启时不会漏读数据。对于Doris Sink来说,开启两阶段提交(2PC) 或利用 Doris 的 Stream Load 配合 Label 唯一性机制。当 Checkpoint 成功触发时,Doris Sink 才会正式提交这批数据;如果任务失败,未提交的数据则会回滚,结合 Doris 的 REPLACE 模型,重复提交只会覆盖,不会重复累加。

[离线] DWS层保证数据写入幂等性(即数据一致性)的核心机制是 Hive 的 INSERT OVERWRITE 语法(结合分区表) 核心原理:使用 INSERT OVERWRITE 覆盖分区 Hive 的 INSERT OVERWRITE 会先删除目标分区(或表)中的旧数据,再写入新数据。因此,无论任务执行多少次,最终该分区里保留的都是最后一次计算生成的完整数据集,不会产生重复数据。

为什么DWS不直接对接应用,为什么要搞ADS

如果跳过ADS层,直接让应用对接DWS,会导致数据模型臃肿、查询性能下降、公共层稳定性差等问题,最终增加维护成本和业务风险

如果订单状态发生了更新,有新的数据来,怎么去重

如果有新的数据,maxwell 输出的 type 字段信息会由 insert 变成 update,并新添 old 字段保留旧数据

你的项目数据量多大?

模拟的数据,有一个波动的范围,日志数据有50到60G,业务数据有1到2G

怎么判断数据仓库做的好不好

判断数据仓库做的好不好,可以从以下几个维度来评估: 数据质量、性能、模型规范和复用性、数据血缘、成本、业务价值

数据倾斜是什么

数据倾斜指的就是在分布式计算中,数据分布不均匀的现象,就导致了某些节点(Task/Partition)处理的数据量会远高于其他节点。 在系统中的表现为大部分任务都很快就完成了,只有一个或者少数的几个任务执行的很慢,从而造成了多数节点空闲,少数节点长时间运行(被称为 “长尾任务”)

这种情况带来的性能瓶颈就是会导致集群的利用率低,造成资源浪费,甚至最终会出现节点OOM(内存溢出)或者超市,导致任务执行失败

数据倾斜发生的原因

[map task] 在map task端造成数据倾斜的原因可能有三点:一,不可拆分的压缩算法,例如使用了gzip,zip等来对数据文件进行的压缩;二,数据文件大小差异大,例如在计算流中存在某些文件特别大;三,数据存在较多的空值。

[reduce task] 在reduce task端造成数据倾斜的原因主要是由于 shuffle + Join/GroupBy操作后的key分布不均匀,而导致key分布不均匀的原因主要有三种可能:一,信息获取不全的时候被填充了默认值;二,该业务本身存在热点(如特价或热销商品购买量大增);三,存在由同一IP的爬虫等产生的恶意数据。

如何判断数据倾斜

怎么快速判断发生了数据倾斜, 1. [时间]看任务运行状态的监控,观察是否有较长处理时间的长尾任务 2. [数据量]检查各Task的Input Size/Records,发生数据倾斜的节点的数据量可能是平均值的数十倍 3. [日志]检查节点日志,倾斜节点由于处理的数据量过大,就容易频繁发生GC(垃圾回收)或者00M(内存溢出),日志都会记录下来 4. [重试]Task由于数据量过大导致的处理速度过慢,进而被标记为失败再被重试 5. [资源]通过YARN或者k8s的监控可以发现会出现倾斜节点的资源利用率接近100%,而其他节点闲置的情况 6. [业务]对数据业务有了解,知道存在热点实体,例如电商大促中的爆款商品。由于数据来自特定热点的特定分区,就会导致某天的日志里激增

如何解决数据倾斜

解决数据倾斜的办法就是针对数据倾斜的原因进行修复 [map] 1. 可以选择可切分的压缩算法(如lzo,bzip2,snappy等)但是注意,虽然lzo压缩文件是可切片的,但是它的可切片特性依赖于其索引,所以需要手动为lzo压缩文件创建索引 2. 尽量让每个数据文件的大小保持基本一致 3. 过滤掉空值,无效数据等异常数据 4. 通过Map端Join(Broadcast Join)消除shuffle,这种适用于大表Join小表的情况,即小表的数据量远远小于大表,小表的大小可以放入内存,并且小表满足广播条件,不然的话一旦广播的数据量很大的话就容易内存溢出。在Spark中,将小表数据封装为广播变量,再广播给所有的executor节点,然后利用map算子在本地进行join操作,就没必要将数据按key分发到同一节点,从而做到彻底消除shuffle阶段。

[reduce] 1. 可以通过hashCode(key)%reduce个数来增大reduce并行度,实现简单但不一定有效 2. 加盐。通过给key添加随机数强行打散数据。具体操作如下: 1. 在map阶段先将key加上随机前缀或者后缀,shuffle之后,reduce task先对打上随机前缀或者后缀的key局部聚合,再将key的前后缀去掉后进行全局聚合,就得到了最终结果。关键必须在map端/combine端先进行一次预聚合。 2. 可以先为数据量特别大的key增加随机前后缀,让key分散到不同的task中,再对小表对应的数据进行复制,加上对应的前后缀,确保Join的时候key互相匹配,原本一个Task的任务就被分散到多个task中执行,最后移除前后缀还原key。这种就适用于由于少数key的数据量大导致join的时候发生数据倾斜的情况。 3. 如果出现了数据倾斜的key较多,以至于无法拆分出倾斜数据的情况,就可以选择大表加盐,小表扩容的做法,扩容就是让该表和前缀列表进行笛卡尔积,相当于将小表按照前缀的数量进行复制扩容,每条数据都扩大了N倍,再按照方法2进行处理。

怎么估算数仓需要的存储资源

分两类数据源:MySQL 业务数据、APP 埋点日志,分别统计每日原始数据大小 分层折算:ODS 原样存储,DWD/DIM/DWS/ADS 经列式压缩,日志占存储大头,全量表会重复占用大量空间 叠加约束:按数据保留时长算总裸容量,再乘以 HDFS 三副本,额外预留 20%-30% 磁盘缓冲 优化降存:用户等缓慢变化维度用拉链表替代每日全量,缩短冷数据保存周期,冷热数据分离存储 预留增长:预估 1-2 年业务数据增量,上线后通过 HDFS 命令监控实际磁盘占用校准

指标下沉怎么做的

指标下沉指的是将宏观层面的指标逐步细化到更小的层级,以便进行更深入的分析。 例如,最近30天销售额这个概念,通过指标下沉可以将其划分为以下几个层级
生产线:不同产品线的月销售额
按地区:不同地区的月销售额
按销售平台:不同销售平台的月销售额

口径下沉:指在定义和计算指标时,对数据口径进行进一步的细化和明确。目的是确保数据的一致性和可比性,以避免不同部门或不同场景下的口径不一致导致的数据混淆或误解。

如何设计Java程序实现清洗Kafka中的日志数据(如去除重复、无效字段)

会结合Flink或自定义Consumer实现:
- 去重:使用Flink的KeyedProcessFunction结合状态管理记录唯一标识(如日志ID);
- 无效字段过滤:通过正则匹配或规则引擎(如Aviator)校验字段合法性;
- 数据格式统一:定义Schema,使用Jackson或FastJSON反序列化后标准化输出。

数据仓库和数据库的区别

这个我从几个维度说一下:

  1. 用途不同:数据库是面向事务的(OLTP),主要用来日常业务操作;数据仓库是面向分析的(OLAP),主要用来做数据分析、报表、决策支持。
  2. 数据特点不同:数据库存的是当前的、实时的业务数据;数据仓库存的是历史的、集成的、多源数据,会把各个业务系统的数据整合到一起。
  3. 设计方式不同:数据库用的是三范式建模,尽量减少冗余;数据仓库常用维度建模(星型、雪花模型),允许适当冗余来提高查询效率。
  4. 数据量和查询不同:数据库单表数据量相对小,查询是简单的增删改查;数据仓库数据量很大,查询是复杂的聚合分析。
  5. 数据更新方式不同:数据库是实时更新的,用户操作就会立即改;数据仓库是批量更新的,比如T+1跑数。

实时统计每5分钟的接口调用成功率(成功状态码为2xx),若成功率低于95%触发告警。请描述实现逻辑(可伪代码)。

我用Flink举例说一下实现逻辑: 1. 确定数据来源 2. flink处理逻辑

// 伪代码
DataStream<AccessLog> logStream = env.addSource(kafkaSource).map(JSON::parse);

// 设置水位线,处理30秒乱序
logStream.assignTimestampsAndWatermarks(
    WatermarkStrategy.forBoundedOutOfOrderness(Duration.seconds(30))
);

// 按接口分组 + 5分钟滚动窗口
logStream.keyBy(AccessLog::getInterfaceName)
    .window(TumblingEventTimeWindows.of(Time.minutes(5)))
    .aggregate(
        // 增量聚合:统计成功数和总数
        (acc, log) -> { 
            acc.total++; 
            if (is2xx(log.statusCode)) acc.success++; 
            return acc; 
        },
        (key, window, accs, out) -> {
            // 窗口触发时计算成功率
            double rate = acc.success * 100.0 / acc.total;
            out.collect(new Result(key, rate, window.end()));
        }
    )
    .filter(result -> result.rate < 95.0) // 低于95%过滤出来
    .addSink(alertSink); // 发送告警
3. 告警 可以邮件、钉钉、企业微信,或者写到告警平台。

简单说流程就是:消费日志 → 按接口分组 → 5分钟窗口 → 统计成功数和总数 → 算成功率 → 低于95%告警。

请描述一个你参与的复杂数据清洗场景,如何解决性能瓶颈?

我曾深度参与DWS层交易域SKU粒度下单各窗口汇总表(DwsTradeSkuOrderWindow)的开发。这个场景的数据清洗与聚合极其复杂,主要面临回撤流(Retract Stream)带来的数据重复和大规模维度关联带来的查询性能两大瓶颈。以下是具体的优化实战:

  1. 复杂数据清洗场景:回撤流去重 场景痛点:数据源是DWD层订单明细(dwdtradeorderdetail),该层由订单明细表、订单表、活动表和优惠券表进行 LEFT JOIN 生成。Flink SQL执行LEFT JOIN时会产生回撤流(即相同orderdetailid会先下发一条字段不完整的数据,再下发一条撤回消息,最后下发完整数据)。如果直接聚合,会导致金额(如splittotal_amount)被重复累加,造成数据翻倍。

性能瓶颈:若采用常规的session window或定时器等待“最完整”数据,延迟至少5秒以上,且状态后端需要缓存大量历史数据,内存压力极大。 解决方案(抵消去重):我们采用了文档中提到的“抵消”思路,通过KeyedProcessFunction按orderdetailid分组,维护一个ValueState存储上一条记录。 逻辑:第一条数据到来时直接下发;后续数据到来时,将状态中旧数据的金额取反后下发(用于冲抵上一次的错误累加),然后再下发当前新数据,并更新状态。 效果:该方案无需等待窗口触发,实现了毫秒级实时去重,且状态中仅保留单条数据,TTL设为60秒,彻底解决了回撤流导致的指标失真和状态膨胀问题。

  1. 性能核心瓶颈:多维表关联(维度查询) 场景痛点:清洗后的数据需补充 SKU名称、SPU名称、品牌、一二三级品类等6个维度字段。这些维度存储在HBase中(DIM层)。

性能瓶颈:如果采用同步MapFunction查询HBase,每条数据都要经历网络IO、序列化反序列化,单并行度吞吐量极低。在高流量下,HBase连接成为系统最致命的阻塞点,背压严重,CPU资源大量浪费在等待IO上。 解决方案(异步IO + 旁路缓存):我们严格按照文档设计进行了两层优化: 异步IO(AsyncDataStream):使用Flink的AsyncDataStream.unorderedWait,将同步查询改为异步非阻塞查询。利用CompletableFuture结合Lettuce(Redis异步客户端)和HBase异步连接,单个并行子任务可连续发送成百上千个请求,无需阻塞等待响应,吞吐量提升了3倍以上。 旁路缓存(Redis):在异步查询之上,引入Redis作为一级缓存。查询时优先readDimAsync读取Redis,命中率极高则直接返回;未命中再异步查询HBase,并将结果setex写入Redis(TTL设为2天)。 缓存一致性保障:修改了DIM层的HBaseSinkFunction,当维度表发生update或delete时,同步删除Redis中对应的缓存Key,确保实时数据关联到的永远是最新维度,且不会因缓存雪崩影响性能。

  1. 最终成果 经过上述“抵消去重”+“异步IO”+“旁路缓存”的组合拳,该复杂清洗聚合任务的端到端延迟控制在5秒以内,在处理千万级数据量时,HBase的QPS压力下降了约70%(因Redis分流),且未发生因回撤流导致的数据重复,完美支撑了后续Doris中的多维度即席查询。

评论