Volver al inicio

Inner join en Flink con CRUD de Debezium

El artículo describe la implementación de inner join de dos tablas en Apache Flink con soporte para operaciones CRUD de Debezium. Procesamiento detallado de mensajes, mapeo de modelos, KeyedCoProcessFunction con estado. Se discuten limitaciones al cambiar la clave de unión.

Flink: inner join con CRUD de Debezium
Advertisement 728x90

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

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

Semántica de los campos:

  • c (crear): before=null, after=Data
  • r (lectura/snapshot): before=null, after=Data
  • u (actualizar): before=Data, after=Data
  • d (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:

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

Lógica de selección de datos:

  • Para c/r — tomar after.
  • Para d — siempre before.
  • Para uafter (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):

Google AdInline article slot
  • 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 op de Debezium determina la elección de before/after.
  • La bandera delete en los modelos asegura la eliminación del data mart.
  • Uno-a-muchos: la iteración sobre MapState genera 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

Advertisement 728x90

Leer después