Zurück zur Startseite

Inner Join in Flink mit Debezium CRUD

Der Artikel beschreibt die Implementierung eines Inner Joins zweier Tabellen in Apache Flink mit Unterstützung für CRUD-Operationen von Debezium. Detaillierte Nachrichtenverarbeitung, Modellzuordnung, KeyedCoProcessFunction mit Status. Einschränkungen bei Änderung des Join-Schlüssels werden diskutiert.

Flink: Inner Join mit CRUD von Debezium
Advertisement 728x90

# 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:

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

Feldsemantik:

  • c (Erstellen): before=null, after=Data
  • r (Lesen/Snapshot): before=null, after=Data
  • u (Aktualisieren): before=Data, after=Data
  • d (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:

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');
    }
}

Logik zur Datenauswahl:

  • Für c/rafter nehmen.
  • Für d — immer before.
  • Für uafter (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:

Google AdInline article slot
  • 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 op aus Debezium bestimmt die Wahl von before/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

Advertisement 728x90

Weiterlesen