# Implementierung eines Inner Joins für Tabellen in Apache Flink mit Unterstützung für CRUD-Operationen
In Apache Flink-basierten Echtzeit-Datamarts erfordern Inner Joins zwischen Streams die Behandlung aller CRUD-Operationen von Debezium. Der vorherige Ansatz ignorierte Löschungen und Updates, was zu falschen Ergebnissen in Ein-zu-Viele-Beziehungen führte. Die neue Implementierung analysiert die Operationen und verwaltet den Status korrekt.
Die Debezium-Nachrichtenstruktur definiert die Logik: Das Feld op gibt den Operationstyp an, before und after enthalten die Daten vor und nach der Änderung.
Parsen von Debezium-Nachrichten
Debezium erzeugt Events in Kafka mit einer festen Struktur:
{
"op": "(c|r|u|d)",
"source": { ... },
"ts_ms": "...",
"before": [Data, null],
"after": [Data, null]
}
Feldsemantik:
c(Erstellen):before=null,after=Datar(Lesen/Snapshot):before=null,after=Datau(Aktualisieren):before=Data,after=Datad(Löschen):before=Data,after=null
Der Mapper extrahiert die aktuellen Daten abhängig von op.
Zuordnung des Datenmodells
Die Domain-Klasse deserialisiert eine Row aus Flink unter Berücksichtigung des Operationstyps:
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');
}
}
Logik zur Datenauswahl:
- Für
c/r—afternehmen. - Für
d— immerbefore. - Für
u—after(Details im nächsten Abschnitt).
Ähnlich für das User-Modell.
KeyedCoProcessFunction für Inner Join
Die Kernlogik in InnerJoinFunction verwendet State, um verwandte Datensätze nach Schlüssel (user_id) zu speichern:
ValueState<User>für den aktuellen Benutzer.MapState<Integer, Domain>für die Domains des Benutzers.
Beim Eintreffen eines User:
- State aktualisieren.
- Alle Domains durchiterieren und für jede ein Output erzeugen.
Beim Eintreffen eines Domain:
- In MapState speichern.
- Wenn Benutzer existiert, Output erzeugen.
Vollständiger Code:
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
}
}
Einschränkungen der Implementierung
Der aktuelle Ansatz behandelt CRUD korrekt, bricht jedoch, wenn sich der Join-Schlüssel (user_id) während eines Updates ändert. Der State ist lokal an den Schlüssel gebunden, ohne Weiterleitung zur neuen Schlüsselpartition.
Die Lösung folgt im nächsten Teil.
Wichtige Punkte:
- Das Parsen von
opaus Debezium bestimmt die Wahl vonbefore/after. - Die
delete-Flagge in den Modellen sorgt für die Entfernung aus dem Datamart. - Ein-zu-viele: Iteration über MapState erzeugt alle Kombinationen.
- Der State-Backend ist für die Performance bei großen Volumen entscheidend.
- Eine Änderung des Join-Schlüssels erfordert spezielle Behandlung.
— Editorial Team
Noch keine Kommentare.