Operadores de Transformación en la API DataStream de Apache Flink

  1. Transformaciones Fundamentales de Flujos

Los operadores de transformación básicos permiten alterar, depurar o multiplicar los registros que transitan por un flujo de datos en Apache Flink.

1.1 Mapeo (Map)

El operador map ejecuta una conversión estricta de uno a uno. Por cada registro ingerido, se emite un único registro transformado. La lógica se define implementando la interfaz MapFunction, donde los parámetros genéricos dictan los tipos de entrada y salida.


public class MapTransformation {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        DataStreamSource<servermetric> metrics = env.fromElements(
            new ServerMetric("node-01", 85.5, 1620000000L),
            new ServerMetric("node-02", 42.0, 1620000005L)
        );

        // Normalizar el uso de CPU de porcentaje (0-100) a escala decimal (0.0-1.0)
        metrics.map(new MapFunction<servermetric double="">() {
            @Override
            public Double map(ServerMetric metric) throws Exception {
                return metric.cpuUsage / 100.0;
            }
        }).print();

        env.execute("Map Example");
    }
}
</servermetric></servermetric>

1.2 Filtrado (Filter)

Mediante el operador filter, se evalúa cada elemento contra una expresión booleana. Solo aquellos que retornan true continúan su camino hacia el siguiente operador, manteniendo inalterado el tipo de dato original. Requeire la implementación de FilterFunction.


public class FilterTransformation {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStreamSource<servermetric> metrics = env.fromElements(
            new ServerMetric("node-01", 95.0, 1620000000L),
            new ServerMetric("node-02", 42.0, 1620000005L),
            new ServerMetric("node-03", 88.5, 1620000010L)
        );

        // Conservar únicamente métricas con uso de CPU crítico (mayor a 90%)
        metrics.filter(new FilterFunction<servermetric>() {
            @Override
            public boolean filter(ServerMetric metric) throws Exception {
                return metric.cpuUsage > 90.0;
            }
        }).print();

        env.execute("Filter Example");
    }
}
</servermetric></servermetric>

1.3 Mapeo Plano (FlatMap)

El operador flatMap permite transformar un elemento de entrada en cero, uno o múltiples elementos de salida. Es ideal para desanidar colecciones o dividir cadenas de texto. Se apoya en la interfaz FlatMapFunction y utiliza un Collector para emitir los resultados.


public class FlatMapTransformation {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStreamSource<logevent> logs = env.fromElements(
            new LogEvent("INFO", "user_login,auth_success,cache_hit"),
            new LogEvent("ERROR", "db_timeout,rollback")
        );

        // Dividir la cadena de etiquetas y emitir cada una como un evento independiente
        logs.flatMap(new FlatMapFunction<logevent string="">() {
            @Override
            public void flatMap(LogEvent log, Collector<string> out) throws Exception {
                for (String tag : log.tags.split(",")) {
                    out.collect(log.level + ":" + tag.trim());
                }
            }
        }).print();

        env.execute("FlatMap Example");
    }
}
</string></logevent></logevent>
  1. Operadores de Agregación

Las agregaciones en Flink requieren que los datos estén agrupados lógicamente para poder consolidar información histórica y actual.

2.1 Agrupación por Clave (KeyBy)

Antes de agregar, es obligatorio particionar el flujo. El operador keyBy distribuye los registros en particiones lógicas basándose en el hash de una clave específica. El resultado es un KeyedStream, estructura indispensable para las operaciones de estado.


public class KeyByExample {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStreamSource<servermetric> metrics = env.fromElements(
            new ServerMetric("node-01", 85.5, 1620000000L),
            new ServerMetric("node-01", 92.0, 1620000005L),
            new ServerMetric("node-02", 42.0, 1620000010L)
        );

        // Particionar el flujo utilizando el identificador del servidor
        KeyedStream<servermetric string=""> keyedMetrics = metrics.keyBy(new KeySelector<servermetric string="">() {
            @Override
            public String getKey(ServerMetric metric) throws Exception {
                return metric.serverId;
            }
        });

        env.execute("KeyBy Example");
    }
}
</servermetric></servermetric></servermetric>

2.2 Agregaciones Simples

Una vez obtenido el KeyedStream, Flink provee funciones integradas como sum(), min(), max(), minBy() y maxBy(). La diferencia entre max y maxBy radica en que el primero solo actualiza el campo evaluado, mientras que el segundo retorna el registro completo que contiene el valor máximo.


public class SimpleAggregation {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStreamSource<servermetric> metrics = env.fromElements(
            new ServerMetric("node-01", 85.5, 1620000000L),
            new ServerMetric("node-01", 92.0, 1620000005L)
        );

        // Obtener el registro completo que posee el pico máximo de CPU por servidor
        metrics.keyBy(m -> m.serverId)
               .maxBy("cpuUsage")
               .print();

        env.execute("Simple Aggregation");
    }
}
</servermetric>

2.3 Reducción (Reduce)

El operador reduce ofrece mayor flexibilidad al permitir definir una lógica de consolidación personalizada entre el estado acumulado y el nuevo evento entrante, implementando la interfaz ReduceFunction.


public class ReduceAggregation {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        env.socketTextStream("localhost", 9999)
           .map(line -> {
               String[] parts = line.split(",");
               return new ServerMetric(parts[0], Double.parseDouble(parts[1]), Long.parseLong(parts[2]));
           })
           .keyBy(m -> m.serverId)
           .reduce(new ReduceFunction<servermetric>() {
               @Override
               public ServerMetric reduce(ServerMetric accumulated, ServerMetric current) throws Exception {
                   // Conservar el mayor uso de CPU, pero actualizar el timestamp al más reciente
                   double maxCpu = Math.max(accumulated.cpuUsage, current.cpuUsage);
                   long latestTs = Math.max(accumulated.timestamp, current.timestamp);
                   return new ServerMetric(accumulated.serverId, maxCpu, latestTs);
               }
           }).print();

        env.execute("Reduce Example");
    }
}
</servermetric>
  1. Funciones Definidas por el Usuario (UDF)

Flink permite encapsular la lógica de negocio mediante clases personalizadas, anonimato o expresiones lambda, destacando las funciones enriquecidas para manejo de ciclos de vida.

3.1 Clases de Funciones y Lambdas

Se pueden implementar interfaces como MapFunction o FilterFunction directamente. Para hacerlas reutilizables, es buena práctica inyectar parámetros a través del constructor.


public class UDFExamples {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        DataStreamSource<servermetric> metrics = env.fromElements(
            new ServerMetric("node-01", 95.0, 1620000000L)
        );

        // Uso de clase con constructor para umbral dinámico
        metrics.filter(new CpuThresholdFilter(90.0)).print("Class UDF");

        // Uso de expresión Lambda
        metrics.filter(m -> m.cpuUsage > 80.0).print("Lambda UDF");

        env.execute("UDF Examples");
    }

    public static class CpuThresholdFilter implements FilterFunction<servermetric> {
        private final double threshold;
        public CpuThresholdFilter(double threshold) { this.threshold = threshold; }
        
        @Override
        public boolean filter(ServerMetric metric) throws Exception {
            return metric.cpuUsage >= threshold;
        }
    }
}
</servermetric></servermetric>

3.2 Funciones Enriquecidas (Rich Functions)

Las versiones "Rich" (ej. RichMapFunction) proporcionan acceso al contexto de ejecución (RuntimeContext) y métodos de ciclo de vida como open() y close(). Esto es crucial para inicializar conexiones a bases de datos o cargar modelos de machine learning en memoria.


public class RichFunctionDemo {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(2);
        
        env.fromElements(1, 2, 3, 4)
           .map(new RichMapFunction<integer string="">() {
               private String prefix;

               @Override
               public void open(Configuration parameters) throws Exception {
                   // Inicialización de recursos pesados (ej. conexión a Redis)
                   prefix = "Processed-";
                   System.out.println("Subtask " + getRuntimeContext().getIndexOfThisSubtask() + " initialized.");
               }

               @Override
               public String map(Integer value) throws Exception {
                   return prefix + value;
               }

               @Override
               public void close() throws Exception {
                   // Liberación de recursos
                   System.out.println("Subtask " + getRuntimeContext().getIndexOfThisSubtask() + " closed.");
               }
           }).print();

        env.execute("Rich Function Demo");
    }
}
</integer>
  1. Estrategias de Particionamiento Físico

El particionamiento físico dicta cómo se distribuyen los datos entre las instancias paralelas de los operadores en el clúster.

  • Shuffle: Distribución aleatoria uniforme mediante shuffle().
  • Rebalance: Distribución round-robin estricta para balancear cargas desiguales usando rebalance().
  • Rescale: Similar a rebalance, pero opera localmente dentro del mismo TaskManager para evitar tráfico de red, invocado con rescale().
  • Broadcast: Replica cada registro a todas las particiones paralelas subsecuentes mediante broadcast().
  • Global: Dirige todo el flujo a la primera sub-tarea del operador downstream con global().

4.1 Particionador Personalizado

Cuando las estrategias nativas no son suficientes, se implementa la interfaz Partitioner para definir un enrutamiento determinista.


public class CustomPartitioning {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(4);

        DataStreamSource<transaction> txStream = env.fromElements(
            new Transaction("tx1", "merchant_A", 150.0),
            new Transaction("tx2", "merchant_B", 200.0)
        );

        txStream.partitionCustom(new MerchantPartitioner(), tx -> tx.merchantId)
                .print();

        env.execute("Custom Partitioning");
    }

    public static class MerchantPartitioner implements Partitioner<string> {
        @Override
        public int partition(String merchantId, int numPartitions) {
            // Enrutamiento basado en el hash del comerciante para garantizar afinidad
            return Math.abs(merchantId.hashCode()) % numPartitions;
        }
    }
}
</string></transaction>
  1. División de Flujos (Splitting)

Separar un flujo en múltiples sub-flujos independientes es una tarea común para enrutar eventos a distintos sistemas de procesamiento.

5.1 Múltiples Filtros vs. Salidas Laterales (Side Outputs)

Aplicar múltiples operadores filter sobre el mismo flujo implica recorrer los datos varias veces, lo cual es ineficiente. La alternativa óptima es utilizar ProcessFunction junto con OutputTag para emitir registros a flujos secundarios en una sola pasada.


public class SideOutputSplitting {
    private static final OutputTag<transaction> HIGH_VALUE_TAG = 
        new OutputTag<transaction>("high-value"){};

    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStreamSource<transaction> transactions = env.fromElements(
            new Transaction("tx1", "m1", 50.0),
            new Transaction("tx2", "m2", 8500.0),
            new Transaction("tx3", "m1", 200.0)
        );

        SingleOutputStreamOperator<transaction> mainStream = transactions.process(
            new ProcessFunction<transaction transaction="">() {
                @Override
                public void processElement(Transaction value, Context ctx, Collector<transaction> out) throws Exception {
                    if (value.amount > 5000.0) {
                        // Enviar a la salida lateral (Side Output)
                        ctx.output(HIGH_VALUE_TAG, value);
                    } else {
                        // Enviar al flujo principal
                        out.collect(value);
                    }
                }
            }
        );

        mainStream.print("Standard Tx");
        mainStream.getSideOutput(HIGH_VALUE_TAG).print("Fraud Alert Tx");

        env.execute("Side Output Splitting");
    }
}
</transaction></transaction></transaction></transaction></transaction></transaction>
  1. Consolidación de Flujos (Merging)

La integración de múltiples fuentes de datos requiere operadores de unión que respeten las semánticas de los tipos involucrados.

6.1 Unión (Union)

El método union() concatena múltiples flujos del mismo tipo en un único flujo. No realiza ninguna combinación de datos, simplemente los fusiona en un solo canal.


public class UnionStreams {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        DataStreamSource<transaction> euStream = env.fromElements(new Transaction("eu1", "m1", 10.0));
        DataStreamSource<transaction> usStream = env.fromElements(new Transaction("us1", "m2", 20.0));
        DataStreamSource<transaction> apacStream = env.fromElements(new Transaction("ap1", "m3", 30.0));

        euStream.union(usStream, apacStream).print("Global Stream");

        env.execute("Union Example");
    }
}
</transaction></transaction></transaction>

6.2 Conexión (Connect) y CoProcessFunction

El operador connect() permite unir dos flujos de tipos diferentes, generando un ConnectedStreams. Para procesarlos conjuntamente, se utilizan funciones como CoMapFunction o CoProcessFunction. Al aplicar keyBy sobre flujos conectados, se garantiza que los eventos con la misma clave de ambos flujos lleguen a la misma instancia paralela, habilitando operaciones de JOIN con estado.


public class ConnectedStreamsJoin {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setParallelism(2);

        DataStreamSource<orderevent> orders = env.fromElements(
            new OrderEvent("order-100", "CREATED"),
            new OrderEvent("order-101", "CREATED")
        );

        DataStreamSource<paymentevent> payments = env.fromElements(
            new PaymentEvent("order-100", "SUCCESS", 250.00)
        );

        ConnectedStreams<orderevent paymentevent=""> connected = orders.connect(payments);

        ConnectedStreams<orderevent paymentevent=""> keyedConnected = 
            connected.keyBy(order -> order.orderId, payment -> payment.orderId);

        SingleOutputStreamOperator<string> joinedResults = keyedConnected.process(
            new CoProcessFunction<orderevent paymentevent="" string="">() {
                
                private transient MapState<string orderevent=""> orderState;

                @Override
                public void open(Configuration parameters) throws Exception {
                    MapStateDescriptor<string orderevent=""> descriptor = 
                        new MapStateDescriptor<>("orderState", Types.STRING, Types.POJO(OrderEvent.class));
                    orderState = getRuntimeContext().getMapState(descriptor);
                }

                @Override
                public void processElement1(OrderEvent order, Context ctx, Collector<string> out) throws Exception {
                    orderState.put(order.orderId, order);
                    out.collect("Order registered: " + order.orderId);
                }

                @Override
                public void processElement2(PaymentEvent payment, Context ctx, Collector<string> out) throws Exception {
                    if (orderState.contains(payment.orderId)) {
                        out.collect("Payment matched for: " + payment.orderId + " | Amount: " + payment.amount);
                    } else {
                        out.collect("Payment received but order missing: " + payment.orderId);
                    }
                }
            }
        );

        joinedResults.print();
        env.execute("Connected Streams Join");
    }
}
</string></string></string></string></orderevent></string></orderevent></orderevent></paymentevent></orderevent>

Etiquetas: apache-flink datastream-api Stream-Processing java flink-transformations

Publicado el 9-27 05:21