返回首页

Flink 实时 JOIN:表 4+2 TB

本文描述了在 Apache Flink 上 TB 级用户表和域表之间实时 INNER JOIN 的实现。使用 Table API 从 Kafka 读取 CDC,并使用 KeyedCoProcessFunction 配合 MapState 处理一对多关系,确保 O(1) 结果更新。

Flink 上 6 TB 表实时 JOIN:完整实现
Advertisement 728x90

在 Apache Flink 上实现大型表实时 JOIN:基本方法

我们需要为 users 表(4 TB)和 domains 表(2 TB)之间的 INNER JOIN 结果启用实时更新。该查询返回 user_id、firstname、lastname 和 domain_name。传统的视图会给数据源(OLTP)带来过载,而通过 CDC 简单物化到 DWH 则无法提供所需的速度。

示例数据:

users:

Google AdInline article slot

| id | firstname | lastname |

|----|-----------|----------|

| 1 | Egor | Myasnik |

Google AdInline article slot

| 2 | Pavel | Hvastun |

| 3 | Mitya | Volk |

domains:

Google AdInline article slot

| 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

Advertisement 728x90

继续阅读