Retour à l'accueil

Jointure interne dans Flink avec CRUD Debezium

L'article décrit l'implémentation de la jointure interne de deux tables dans Apache Flink avec support des opérations CRUD de Debezium. Traitement détaillé des messages, mappage des modèles, KeyedCoProcessFunction avec état. Limitations lors du changement de la clé de jointure discutées.

Flink : jointure interne avec CRUD de Debezium
Advertisement 728x90

Mise en œuvre de l'inner join entre tables dans Apache Flink avec prise en charge des opérations CRUD

Dans les data marts en temps réel basés sur Apache Flink, les inner joins entre flux nécessitent la gestion de toutes les opérations CRUD provenant de Debezium. L'approche précédente ignorait les suppressions et les mises à jour, entraînant des résultats incorrects dans les relations un-à-plusieurs. La nouvelle implémentation analyse les opérations et gère correctement l'état.

La structure des messages Debezium définit la logique : le champ op indique le type d'opération, before et after contiennent les données avant et après le changement.

Analyse des messages Debezium

Debezium génère des événements dans Kafka avec une structure fixe :

Google AdInline article slot
{
	"op": "(c|r|u|d)",
	"source": { ... },
	"ts_ms": "...",
	"before": [Data, null],
	"after": [Data, null]
}

Sémantique des champs :

  • c (création) : before=null, after=Data
  • r (lecture/instantané) : before=null, after=Data
  • u (mise à jour) : before=Data, after=Data
  • d (suppression) : before=Data, after=null

Le mapper extrait les données actuelles en fonction de op.

Mapping du modèle de données

La classe Domain désérialise une Row de Flink, en tenant compte du type d'opération :

Google AdInline article slot
public class Domain implements Serializable {
  public Integer id;
  public Integer user_id;
  public String domain_name;
  public boolean delete;

  public static Domain fromRow(Row row) {
        Character op = (Character) row.getField(1);

        if (op == null) {
            throw new IllegalStateException("Never should happen, if Debezium feels fine");
        }

        Row domain = (Row) row.getField(op == 'd' ? 0 : 2);

        Integer id = (Integer) domain.getField(0);
        Integer user_id = (Integer) domain.getField(1);
        String domain_name = (String) domain.getField(2);  // ispravleno on String

        return new Domain(id, user_id, domain_name, op == 'd');
    }
}

Logique de sélection des données :

  • Pour c/r — prendre after.
  • Pour d — toujours before.
  • Pour uafter (détails dans la section suivante).

De même pour le modèle User.

KeyedCoProcessFunction pour l'inner join

La logique principale dans InnerJoinFunction utilise l'état pour stocker les enregistrements liés par clé (user_id) :

Google AdInline article slot
  • ValueState<User> pour l'utilisateur actuel.
  • MapState<Integer, Domain> pour les domaines de l'utilisateur.

À l'arrivée d'un User :

  • Mettre à jour l'état.
  • Itérer sur tous les domaines et générer une Output pour chacun.

À l'arrivée d'un Domain :

  • Enregistrer dans MapState.
  • Si l'utilisateur existe, générer une Output.

Code complet :

import org.apache.flink.api.common.functions.OpenContext;
import org.apache.flink.api.common.state.*;
import org.apache.flink.streaming.api.functions.co.KeyedCoProcessFunction;
import org.apache.flink.util.Collector;

import java.io.Serializable;

public class InnerJoinFunction extends KeyedCoProcessFunction<Integer, User, Domain, InnerJoinFunction.Output> {
    private MapState<Integer, Domain> domainsState;
    private ValueState<User> usersState;

    @Override
    public void processElement1(final User user, final Context ctx, final Collector<InnerJoinFunction.Output> out) throws Exception {
        usersState.update(user);
      
        for (final Domain domain : domainsState.values()) {
            out.collect(new InnerJoinFunction.Output(
                    user.id,
                    user.firstname,
                    user.lastname,
                    domain.domain_name,
                    user.delete || domain.delete
            ));
        }
    }

    @Override
    public void processElement2(final Domain domain, final Context ctx, final Collector<InnerJoinFunction.Output> out) throws Exception {
        domainsState.put(domain.id, domain);

        final User user = usersState.value();

        if (user != null) {
            out.collect(new InnerJoinFunction.Output(
                    user.id,
                    user.firstname,
                    user.lastname,
                    domain.domain_name,
                    user.delete || domain.delete
            ));
        }
    }

    @Override
    public void open(OpenContext openContext) throws Exception {
        var usersStateDescriptor = new ValueStateDescriptor<>(
                "users",
                User.class
        );
        var domainsStateDescriptor = new MapStateDescriptor<>(
                "domains",
                Integer.class,
                Domain.class
        );
        usersState = getRuntimeContext().getState(usersStateDescriptor);
        domainsState = getRuntimeContext().getMapState(domainsStateDescriptor);

        super.open(openContext);
    }

    public static class Output implements Serializable {
        public Integer user_id;
        public String firstname;
        public String lastname;
        public String domain_name;
        public boolean delete;
        
        // getters and setters omitted
    }
}

Limites de l'implémentation

L'approche actuelle gère correctement les CRUD mais casse lorsque la clé de jointure (user_id) change lors d'une mise à jour. L'état est lié localement à la clé, sans retransmission vers la nouvelle partition de clé.

La solution est dans la partie suivante.

Points clés :

  • L'analyse de op de Debezium détermine le choix de before/after.
  • Le drapeau delete dans les modèles assure la suppression du data mart.
  • Un-à-plusieurs : l'itération sur MapState génère toutes les combinaisons.
  • Le backend d'état est critique pour les performances avec de grands volumes.
  • Changer la clé de jointure nécessite un traitement spécial.

— Editorial Team

Advertisement 728x90

Lire ensuite