> For the complete documentation index, see [llms.txt](https://shimin-huang.gitbook.io/doc/llms.txt). Markdown versions of documentation pages are available by appending `.md` to page URLs; this page is available as [Markdown](https://shimin-huang.gitbook.io/doc/datawarehouse/fang-an-shi-jian/ji-yu-flink-de-shi-shi-shu-cang-jian-she.md).

# 基于Flink的实时数仓建设

## 概述

### 建设目的

**解决由于传统数据仓库数据时效性低解决不了的问题**

* 面向主题的
* 集成的
* 相对稳定的
* 处理上一次批处理流程到当前的数据

### 实时数仓架构

#### Lambda架构

![img](https://2758483936-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2FcIRMOn4u9hdIi6agitCf%2Fuploads%2Fgit-blob-052802a99ba9cfe2a1111456dbf217d9b8a0042f%2FLambda%E6%9E%B6%E6%9E%84.jpg?alt=media)

#### Kappa架构

![](https://2758483936-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2FcIRMOn4u9hdIi6agitCf%2Fuploads%2Fgit-blob-03f83f61fe37a0e3508512f7a6d0221c59f6c72e%2FKappa%E6%9E%B6%E6%9E%84.jpg?alt=media)

#### 实时OLAP架构

![](https://2758483936-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2FcIRMOn4u9hdIi6agitCf%2Fuploads%2Fgit-blob-639cd7c9c5f1f6172bfb93148b9dc235727dbc0e%2F%E5%AE%9E%E6%97%B6OLAP%E6%9E%B6%E6%9E%84.jpg?alt=media)

## 流式Join

### 纬表join

#### 将纬度表加载到内存关联

**方案1**

* 通过flink的RichFlatMapFunction的open方法，一次性将纬度表的数据全部加载到内存中，后续在每条流的消息去内存中关联。
* 优点是实现简单，但是仅支持小数据量纬度表，更新纬度表需要重启任务
* 适用于纬度表小、变更频率低、对变更及时性要求低。

**方案2**

* 通过Distributed Cache分发本地纬度表文件到task manager后加载到内存关联
* 使用env.registerCachedFile注册文件
* 实现RichFuntion在oepn方法中通过RuntimeContext获取cache文件，解析和使用文件数据
* 适用于纬度表小、变更频率低、对变更及时性要求低。

**方案3**

* 理论外部缓存来存储维度表，在将外部缓存维度表加载到内存中使用
* 纬度更新反馈到结果有延迟，一般是从外部缓存倒入内存的延迟问题

#### 广播维度表

* 实现方式
  * 将维度表数据发送到kafka作为广播原始流S1
  * 定义状态描述符MapStateDescripitor。调用S1.broadcast，获取broadCastStream S2
  * 调用非广播流S3.connect(S2)，得到BroadcastConnectedStream S2
  * 在KeyedBroadcastProcessFunction/BroadcastProcessFunction实现关联处理逻辑，并作为参数调用S4.process()
* 优点:纬度变更可即时更新到结果
* 缺点:数据保存在内存中，支持维度表数据量较小
* 适用于实时感知维度变更，维度数据可以转换为实时流的场景

#### Temporal Table

* 适用于changelog流，存储各个时态数据的变化

![](https://2758483936-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2FcIRMOn4u9hdIi6agitCf%2Fuploads%2Fgit-blob-352e11ee45f2e6bbc12f973f2afa69739125dda0%2FTemporaltable.jpg?alt=media)

#### 维度表join方案对比

![](https://2758483936-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2FcIRMOn4u9hdIi6agitCf%2Fuploads%2Fgit-blob-a1bd51b80305dbb81d433f2031676ee25eb62d94%2F%E7%BB%B4%E5%BA%A6%E8%A1%A8join%E6%96%B9%E6%A1%88%E5%AF%B9%E6%AF%94.jpg?alt=media)

### 双流Join

#### Join-Regular join

```sql
# 双流join
SELECT * FROM Orders INNER JOIN Product ON Orders.productId=Product.id
```

* 仅支持有界流和等号连接

#### join-Interval join

* 限定join的时间窗口，对超出时间范围的数据清理，避免保留全量State，支持processing time和event time

```sql
SELECT *
FROM Orders o,
	Shipments s
WHERE
	o.id=s.orderId
	AND s.shiptime BETWEEN o.ordertime
	AND o.ordertime + INTERVAL '4' HOUR
```

#### join-Window join

* 将两个流中有相同key和处在相同window的元素做join

![](https://2758483936-files.gitbook.io/~/files/v0/b/gitbook-x-prod.appspot.com/o/spaces%2FcIRMOn4u9hdIi6agitCf%2Fuploads%2Fgit-blob-0f0f0d03bf779d6c65380ad6f2e73860123c559c%2Fwindowjoin.jpg?alt=media)

## 实时数仓问题解决

### 大State问题

* 可以在数据接入时通过`ROW_NUMBER`函数对数据流去重，然后在进行join。

```sql
-- 去重
create view view1
select *(
select *,row_number()over(partition by id order by proctime()desc) as rn from s1)
where rn =1;

create view view2
select *(
select *,row_number()over(partition by id order by proctime()desc) as rn from s2)
where rn =1;

insert into dwd_t
select view1.id,view2.name
view1 left outer join view2 on view1.id=view2.id
```

### 多流join优化

* 将多流通过union all合并，把数据错位拼接到一起，后面加一层Group By，相当于将Join关联转换成Group By

### 回溯历史数据

* 采用批流混合的方式来完成状态复用，基于Blink流处理来处理实时消息流，基于Blink的批处理完成离线计算，通过两者的融合，在同一个任务里完成历史所有数据的计算
* 将实时的流和存储在olap系统的总的离线数据进行union all，完成消息的回溯
