# Implementacja inner join dla tabel w Apache Flink z obsługą operacji CRUD
W systemach przetwarzania danych w czasie rzeczywistym opartych na Apache Flink inner join między strumieniami wymaga uwzględnienia wszystkich operacji CRUD pochodzących z Debezium. Poprzednie podejście ignorowało usuwania i aktualizacje, co prowadziło do niepoprawnych wyników w relacjach one-to-many. Nowa implementacja parsuje operacje i poprawnie zarządza stanem.
Struktura wiadomości Debezium określa logikę: pole op wskazuje typ operacji, before i after zawierają dane przed/po zmianie.
Parsowanie wiadomości Debezium
Debezium generuje zdarzenia w Kafka o stałej strukturze:
{
"op": "(c|r|u|d)",
"source": { ... },
"ts_ms": "...",
"before": [Data, null],
"after": [Data, null]
}
Semantyka pól:
c(create):before=null,after=Datar(read/snapshot):before=null,after=Datau(update):before=Data,after=Datad(delete):before=Data,after=null
Mapper wyciąga aktualne dane w zależności od op.
Mapowanie modeli danych
Klasa Domain deserializuje Row z Flink, uwzględniając typ operacji:
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');
}
}
Logika wyboru danych:
- Dla
c/r— bierzemyafter. - Dla
d— zawszebefore. - Dla
u—after(szczegóły w następnej części).
Podobnie dla modelu User.
KeyedCoProcessFunction do inner join
Główna logika w InnerJoinFunction używa state do przechowywania powiązanych rekordów według klucza (user_id):
ValueState<User>dla bieżącego użytkownika.MapState<Integer, Domain>dla domen użytkownika.
Przy nadejściu User:
- Aktualizujemy state.
- Iterujemy wszystkie domeny i generujemy Output dla każdej.
Przy nadejściu Domain:
- Zapisujemy w MapState.
- Jeśli użytkownik istnieje, generujemy Output.
Pełny kod:
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
}
}
Ograniczenia implementacji
Obecne podejście poprawnie obsługuje CRUD, ale zawodzi przy zmianie klucza join (user_id) w aktualizacji. State jest powiązany z kluczem lokalnie, retransmisja do nowej kluczowej partycji nie następuje.
Rozwiązanie — w następnej części.
Co ważne:
- Parsowanie
opz Debezium określa wybórbefore/after. - Flaga
deletew modelach zapewnia usunięcie z witryny. - One-to-many: iteracja po MapState generuje wszystkie kombinacje.
- State-backend jest kluczowy dla wydajności przy dużych wolumenach.
- Zmiana klucza join wymaga specjalnej obsługi.
— Editorial Team
Brak komentarzy.