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 :
{
"op": "(c|r|u|d)",
"source": { ... },
"ts_ms": "...",
"before": [Data, null],
"after": [Data, null]
}
Sémantique des champs :
c(création) :before=null,after=Datar(lecture/instantané) :before=null,after=Datau(mise à jour) :before=Data,after=Datad(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 :
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— prendreafter. - Pour
d— toujoursbefore. - Pour
u—after(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) :
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
Outputpour 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
opde Debezium détermine le choix debefore/after. - Le drapeau
deletedans les modèles assure la suppression du data mart. - Un-à-plusieurs : l'itération sur
MapStategé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
Aucun commentaire pour le moment.