Retour à l'accueil

Flink JOIN en temps réel : tables 4+2 TB

L'article décrit l'implémentation d'un INNER JOIN en temps réel entre les tables users et domains de téraoctets sur Apache Flink. Table API est utilisé pour lire CDC depuis Kafka et KeyedCoProcessFunction avec MapState pour la relation un-à-plusieurs. Mise à jour de résultat O(1) assurée.

JOIN en temps réel de tables 6 TB sur Flink : implémentation complète
Advertisement 728x90

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 :

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|

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

Advertisement 728x90

Lire ensuite