在 Apache Flink 上实现大型表实时 JOIN:基本方法
我们需要为 users 表(4 TB)和 domains 表(2 TB)之间的 INNER JOIN 结果启用实时更新。该查询返回 user_id、firstname、lastname 和 domain_name。传统的视图会给数据源(OLTP)带来过载,而通过 CDC 简单物化到 DWH 则无法提供所需的速度。
示例数据:
users:
| id | firstname | lastname |
|----|-----------|----------|
| 1 | Egor | Myasnik |
| 2 | Pavel | Hvastun |
| 3 | Mitya | Volk |
domains:
| id | user_id | domain_name |
|----|---------|-------------|
| 1 | 1 | Approval |
| 2 | 1 | Rejection |
| 3 | 1 | Stoppage |
| 4 | 3 | Cancellation|
备选方案分析
- 直接在 OLTP 上创建视图 — 由于对 OLTP 系统的负载过大,不可接受。
- CDC → Kafka → DWH — 缓解了 OLTP 的负担,但从 DWH 读取的速度与原先相似。额外增加了两层:Kafka 和视图。
两种方法都无法为 JOIN 结果提供实时更新。
Flink 架构
解决方案:CDC 流式数据进入 Kafka → Flink Table API → 有状态 JOIN → sink(MVP 使用控制台)。
Maven 依赖
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-kafka</artifactId>
<version>3.2.0-1.19</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-sql-client</artifactId>
<version>1.19</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-planner-loader</artifactId>
<version>${flink.version}</version>
</dependency>
数据模型
public class User implements Serializable {
public Integer id;
public String firstname;
public String lastname;
// getters/setters
public static User fromRow(Row row) { /* mapper */ }
}
public class Domain implements Serializable {
public Integer id;
public Integer user_id;
public String domain_name;
// getters/setters
public static Domain fromRow(Row row) { /* mapper */ }
}
创建 Kafka 表
tableEnv.executeSql("CREATE TABLE users (" +
"`before` ROW<id: INT, firstname: STRING, lastname: STRING>," +
"`op` STRING," +
"`after` ROW<id: INT, firstname: STRING, lastname: STRING>," +
") WITH (" +
"'connector' = 'kafka'," +
"'topic' = 'users_topic'," +
"'properties.bootstrap.servers' = 'kafka-brokers'," +
"'properties.group.id' = 'users_consumer_group'," +
"'scan.startup.mode' = 'earliest'");
DataStream<User> users = tableEnv.toDataStream(tableEnv.from("users")).map(User::fromRow);
// domains 类似
tableEnv.executeSql("CREATE TABLE domains (...)");
DataStream<Domain> domains = tableEnv.toDataStream(tableEnv.from("domains")).map(Domain::fromRow);
有状态 JOIN 实现
使用 KeyedCoProcessFunction 按 user_id == domain.user_id 对流进行 JOIN:
users
.connect(domains)
.keyBy(
user -> user.id,
domain -> domain.user_id
)
.process(new Join1())
.print();
KeyedCoProcessFunction
public class Join1 extends KeyedCoProcessFunction<Integer, User, Domain, Output> {
private MapState<Integer, Domain> domainsState;
private ValueState<User> usersState;
@Override
public void processElement1(User user, Context ctx, Collector<Output> out) throws Exception {
usersState.update(user);
for (Domain domain : domainsState.values()) {
out.collect(new Output(user.id, user.firstname, user.lastname, domain.domain_name));
}
}
@Override
public void processElement2(Domain domain, Context ctx, Collector<Output> out) throws Exception {
domainsState.put(domain.id, domain);
User user = usersState.value();
if (user != null) {
out.collect(new Output(user.id, user.firstname, user.lastname, domain.domain_name));
}
}
@Override
public void open(OpenContext openContext) throws Exception {
usersState = getRuntimeContext().getState(new ValueStateDescriptor<>("users", User.class));
domainsState = getRuntimeContext().getMapState(
new MapStateDescriptor<>("domains", Integer.class, Domain.class));
}
public static class Output implements Serializable {
public Integer user_id;
public String firstname, lastname, domain_name;
// constructors/getters/setters
}
}
MapState 用于 domains,因为存在一对多关系。状态操作的时间复杂度为 O(1)。
关键点
- 实时更新:JOIN 结果在 CDC 事件到达时即保持最新。
- 有状态处理:users 使用 ValueState(1:1),domains 使用 MapState(1:N)。
- 可扩展性:已在 TB 级数据量上通过 OLAP sink 验证。
- 性能:O(1) 状态访问,低延迟。
- 后续步骤:处理 DELETE/UPDATE、RocksDB、倾斜均衡。
生产环境推荐
- 使用 OLAP sink(如 ClickHouse、Pinot)代替控制台。
- 处理 CDC 中的 op 标志,以正确处理 UPDATE/DELETE。
- 大数据量时使用 RocksDB 作为状态后端。
- 监控键倾斜并分配分区。
- 在集群上测试扩展。
— Editorial Team
暂无评论。