Implémentation d'un JOIN en temps réel pour les grandes tables sur Apache Flink : Approche de base
Nous devons activer les mises à jour en temps réel pour le résultat de l'INNER JOIN entre la table users (4 To) et domains (2 To). La requête retourne user_id, firstname, lastname et domain_name. Les vues traditionnelles surchargent la source de données (OLTP), tandis que la matérialisation simple via CDC vers l'entrepôt de données ne fournit pas la vitesse requise.
Exemples de données :
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|
Analyse des alternatives
- Vue directe sur OLTP — inacceptable en raison de la charge sur le système OLTP.
- CDC → Kafka → entrepôt de données — décharge l'OLTP, mais la vitesse de lecture depuis l'entrepôt de données est similaire à l'originale. Ajoute deux couches : Kafka et la vue.
Les deux approches ne permettent pas d'obtenir des mises à jour en temps réel pour le résultat du JOIN.
Architecture Flink
Solution : flux CDC vers Kafka → Table API de Flink → JOIN stateful → sink (console pour le MVP).
Dépendances 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>
Modèles de données
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 */ }
}
Création des tables 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);
// De même pour domains
tableEnv.executeSql("CREATE TABLE domains (...)");
DataStream<Domain> domains = tableEnv.toDataStream(tableEnv.from("domains")).map(Domain::fromRow);
Implémentation du JOIN stateful
Rejoindre les flux par la clé user_id == domain.user_id en utilisant KeyedCoProcessFunction :
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 est utilisé pour les domaines en raison de la relation un-à-plusieurs. Les opérations d'état sont en O(1).
Points clés
- Mises à jour en temps réel : le résultat du JOIN est à jour dès l'arrivée de l'événement CDC.
- Traitement stateful : ValueState pour les utilisateurs (1:1), MapState pour les domaines (1:N).
- Évolutivité : prouvée sur des volumes à l'échelle des téraoctets avec sink OLAP.
- Performances : accès à l'état en O(1), faible latence.
- Étapes suivantes : gérer DELETE/UPDATE, RocksDB, équilibrage du déséquilibre des clés.
Recommandations pour la production
- Utiliser un sink OLAP (ClickHouse, Pinot) au lieu de la console.
- Traiter les drapeaux op provenant de CDC pour une gestion correcte des UPDATE/DELETE.
- Utiliser RocksDB pour le backend d'état avec de gros volumes.
- Surveiller le déséquilibre des clés et répartir les partitions.
- Tester le scaling sur un cluster.
— Editorial Team
Aucun commentaire pour le moment.