# Implementación de Inner Join para tablas en Apache Flink con soporte para operaciones CRUD
En los data marts en tiempo real basados en Apache Flink, los inner joins entre flujos requieren manejar todas las operaciones CRUD provenientes de Debezium. El enfoque anterior ignoraba las eliminaciones y actualizaciones, lo que generaba resultados incorrectos en relaciones uno-a-muchos. La nueva implementación analiza las operaciones y gestiona correctamente el estado.
La estructura del mensaje de Debezium define la lógica: el campo op indica el tipo de operación, before y after contienen los datos antes y después del cambio.
Análisis de mensajes de Debezium
Debezium genera eventos en Kafka con una estructura fija:
{
"op": "(c|r|u|d)",
"source": { ... },
"ts_ms": "...",
"before": [Data, null],
"after": [Data, null]
}
Semántica de los campos:
c(crear):before=null,after=Datar(lectura/snapshot):before=null,after=Datau(actualizar):before=Data,after=Datad(eliminar):before=Data,after=null
El mapper extrae los datos actuales según el valor de op.
Mapeo del modelo de datos
La clase Domain deserializa una Row de Flink, teniendo en cuenta el tipo de operación:
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');
}
}
Lógica de selección de datos:
- Para
c/r— tomarafter. - Para
d— siemprebefore. - Para
u—after(detalles en la siguiente sección).
De manera similar para el modelo User.
KeyedCoProcessFunction para Inner Join
La lógica principal en InnerJoinFunction utiliza estado para almacenar registros relacionados por clave (user_id):
ValueState<User>para el usuario actual.MapState<Integer, Domain>para los dominios del usuario.
Al llegar un User:
- Actualizar el estado.
- Iterar sobre todos los dominios y generar Output para cada uno.
Al llegar un Domain:
- Guardar en
MapState. - Si existe el usuario, generar Output.
Código completo:
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
}
}
Limitaciones de la implementación
El enfoque actual maneja correctamente las operaciones CRUD, pero falla cuando la clave de unión (user_id) cambia durante una actualización. El estado está ligado localmente a la clave, sin retransmisión a la nueva partición de clave.
La solución está en la siguiente parte.
Puntos clave:
- El análisis de
opde Debezium determina la elección debefore/after. - La bandera
deleteen los modelos asegura la eliminación del data mart. - Uno-a-muchos: la iteración sobre
MapStategenera todas las combinaciones. - El backend de estado es crítico para el rendimiento con grandes volúmenes.
- Cambiar la clave de unión requiere un manejo especial.
— Editorial Team
Aún no hay comentarios.