Apache Flink en 2026: Procesamiento de Streams, Event Time y Preguntas de Entrevista
Guía completa de Apache Flink 2.3 para procesamiento de streams en tiempo real. Watermarks, ventanas, gestión de estado y preparación para entrevistas de data engineering.

Apache Flink 2.3 representa el estándar del procesamiento de streams distribuido con latencia inferior al segundo y garantías exactly-once. A diferencia de los sistemas orientados a batch, Flink procesa datos continuamente a medida que llegan, convirtiéndose en el framework preferido para analítica en tiempo real, detección de fraude y arquitecturas orientadas a eventos.
Flink se distingue de Spark Streaming mediante un verdadero procesamiento de streams: Flink procesa eventos uno por uno con semántica event time, mientras que Spark Streaming procesa micro-batches con processing time por defecto.
Arquitectura de Flink 2.3 para Procesamiento de Streams
Flink se ejecuta sobre una arquitectura distribuida con un JobManager que coordina el trabajo entre múltiples TaskManagers. Cada TaskManager ejecuta task slots que procesan porciones de los operadores paralelos del job. Esta separación permite a Flink escalar horizontalmente mientras mantiene la tolerancia a fallos mediante checkpoints distribuidos.
El modelo dataflow en Flink representa los cómputos como grafos acíclicos dirigidos (DAG). Los datos fluyen desde las fuentes a través de transformaciones hacia los sinks, con cada operador potencialmente ejecutándose en múltiples instancias paralelas.
// Basic Flink streaming application setup
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Enable checkpointing every 10 seconds for fault tolerance
env.enableCheckpointing(10000, CheckpointingMode.EXACTLY_ONCE);
// Configure state backend for large state
env.setStateBackend(new EmbeddedRocksDBStateBackend());
// Define the data source - Kafka in production scenarios
DataStream<String> rawStream = env.addSource(
new FlinkKafkaConsumer<>("events", new SimpleStringSchema(), kafkaProps)
);
// Parse and transform the stream
DataStream<Event> events = rawStream
.map(json -> objectMapper.readValue(json, Event.class))
.assignTimestampsAndWatermarks(
WatermarkStrategy.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((event, timestamp) -> event.getTimestamp())
);La configuración de checkpoint anterior almacena instantáneas del estado distribuido cada 10 segundos. Si ocurre una falla, Flink restaura desde el último checkpoint completado y reproduce los registros desde Kafka.
Event Time vs Processing Time: Entendiendo las Semánticas Temporales
El event time se refiere al momento en que un evento realmente ocurrió, incorporado en los datos mismos. El processing time es cuando Flink procesa el registro. Esta distinción es importante porque los retrasos de red, las entregas desordenadas y los atrasos de procesamiento hacen que el processing time sea poco confiable para operaciones basadas en tiempo.
// Configuring event time with watermarks
public class EventTimeProcessor {
public DataStream<AggregatedMetric> processWithEventTime(
DataStream<SensorReading> readings) {
return readings
// Extract timestamp from the event payload
.assignTimestampsAndWatermarks(
WatermarkStrategy
.<SensorReading>forBoundedOutOfOrderness(Duration.ofMinutes(2))
.withTimestampAssigner((reading, ts) -> reading.getEventTime())
.withIdleness(Duration.ofMinutes(5)) // Handle idle partitions
)
// Window by event time, not wall clock
.keyBy(SensorReading::getSensorId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(new AverageAggregator());
}
}La estrategia forBoundedOutOfOrderness indica a Flink que los eventos pueden llegar con hasta 2 minutos de retraso. Los watermarks avanzan cuando Flink determina que no llegarán más eventos con timestamps anteriores al watermark.
¿Qué sucede con los eventos tardíos en Flink? Por defecto, los eventos que llegan después de que el watermark ha pasado el tiempo final de la ventana son descartados. Se puede configurar latencia permitida con .allowedLateness(Time.minutes(10)) para procesar llegadas tardías, o usar side outputs para capturarlas y manejarlas por separado.
Estrategias de Ventaneo para Analítica en Tiempo Real
Flink proporciona cuatro tipos de ventanas: tumbling, sliding, session y global. Cada una responde a diferentes necesidades analíticas.
// Different windowing approaches for stream processing
public class WindowingStrategies {
// Tumbling windows: fixed-size, non-overlapping
// Use case: hourly aggregations, daily summaries
public DataStream<Summary> tumblingAggregation(DataStream<Transaction> txns) {
return txns
.keyBy(Transaction::getAccountId)
.window(TumblingEventTimeWindows.of(Time.hours(1)))
.sum("amount");
}
// Sliding windows: fixed-size, overlapping
// Use case: moving averages, rolling metrics
public DataStream<Double> slidingAverage(DataStream<Metric> metrics) {
return metrics
.keyBy(Metric::getCategory)
.window(SlidingEventTimeWindows.of(Time.minutes(10), Time.minutes(1)))
.aggregate(new AverageAggregator());
}
// Session windows: activity-based, variable size
// Use case: user sessions, conversation analysis
public DataStream<Session> sessionAnalysis(DataStream<ClickEvent> clicks) {
return clicks
.keyBy(ClickEvent::getUserId)
.window(EventTimeSessionWindows.withGap(Time.minutes(30)))
.process(new SessionBuilder());
}
}Las ventanas de sesión se cierran después de un período configurable de inactividad. Este patrón funciona para el análisis de comportamiento de usuarios donde la duración de las sesiones varía según el engagement.
¿Listo para aprobar tus entrevistas de Data Engineering?
Practica con nuestros simuladores interactivos, flashcards y tests técnicos.
Gestión de Estado y Checkpointing
Flink mantiene el estado de los operadores y el estado keyed durante todo el procesamiento. El estado keyed particiona los datos por clave, permitiendo procesamiento paralelo mientras mantiene los registros relacionados juntos. El estado de operador se aplica a toda la instancia del operador.
// Managing state in a Flink KeyedProcessFunction
public class FraudDetector extends KeyedProcessFunction<String, Transaction, Alert> {
// Keyed state: one value per key (account)
private ValueState<Double> lastAmountState;
private ValueState<Long> lastTransactionTimeState;
private MapState<String, Integer> merchantCountState;
@Override
public void open(Configuration parameters) {
// Initialize state descriptors
lastAmountState = getRuntimeContext().getState(
new ValueStateDescriptor<>("lastAmount", Double.class));
lastTransactionTimeState = getRuntimeContext().getState(
new ValueStateDescriptor<>("lastTime", Long.class));
merchantCountState = getRuntimeContext().getMapState(
new MapStateDescriptor<>("merchantCounts", String.class, Integer.class));
}
@Override
public void processElement(Transaction txn, Context ctx, Collector<Alert> out)
throws Exception {
Double lastAmount = lastAmountState.value();
Long lastTime = lastTransactionTimeState.value();
// Detect suspicious patterns
if (lastAmount != null && lastTime != null) {
long timeDelta = txn.getTimestamp() - lastTime;
// Flag transactions 10x larger than previous within 1 minute
if (txn.getAmount() > lastAmount * 10 && timeDelta < 60000) {
out.collect(new Alert(txn.getAccountId(), "SUSPICIOUS_SPIKE", txn));
}
}
// Update state for next transaction
lastAmountState.update(txn.getAmount());
lastTransactionTimeState.update(txn.getTimestamp());
// Track merchant frequency
Integer count = merchantCountState.get(txn.getMerchantId());
merchantCountState.put(txn.getMerchantId(), (count == null ? 0 : count) + 1);
}
}Este procesador con estado rastrea patrones de transacción por cuenta. El estado persiste a través de los checkpoints, sobreviviendo a fallas sin perder el contexto de detección de fraude.
Flink SQL y Table API para Procesamiento de Streams
Flink 2.3 expande sus capacidades SQL con Materialized Tables para mantenimiento incremental de vistas. La Table API proporciona una interfaz unificada para procesamiento batch y streaming.
-- flink_sql_streaming.sql
-- Create a streaming source table from Kafka
CREATE TABLE orders (
order_id STRING,
customer_id STRING,
product_id STRING,
amount DECIMAL(10, 2),
order_time TIMESTAMP(3),
-- Define watermark for event time processing
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'json',
'scan.startup.mode' = 'earliest-offset'
);
-- Streaming aggregation with tumbling window
SELECT
customer_id,
TUMBLE_START(order_time, INTERVAL '1' HOUR) AS window_start,
COUNT(*) AS order_count,
SUM(amount) AS total_amount
FROM orders
GROUP BY
customer_id,
TUMBLE(order_time, INTERVAL '1' HOUR);El enfoque SQL simplifica el desarrollo para analistas familiarizados con SQL mientras Flink maneja la complejidad del procesamiento de streams internamente.
Flink vs Spark Structured Streaming: Comparación Técnica
Ambos frameworks procesan datos en streaming, pero sus arquitecturas difieren fundamentalmente. Flink procesa registros individualmente con verdadero streaming, mientras que Spark procesa micro-batches. Para comparaciones con Apache Spark, los compromisos entre latencia y consistencia importan en producción.
| Aspecto | Flink | Spark Structured Streaming |
|---|---|---|
| Modelo de Procesamiento | Verdadero streaming | Micro-batch |
| Latencia | Milisegundos | Segundos (intervalo batch) |
| State Backend | RocksDB, HashMaps | En memoria, HDFS |
| Exactly-Once | Nativo con checkpoints | Requiere sinks idempotentes |
| Event Time | Soporte nativo | Soportado desde 2.1 |
| Soporte SQL | SQL streaming completo | Ventaneo limitado |
Cuando se pregunta sobre Flink vs Spark para streaming, conviene enfocarse en la adecuación al caso de uso. Flink sobresale en procesamiento de eventos de baja latencia y patrones de eventos complejos. Spark Streaming es adecuado para organizaciones que ya usan Spark para batch y necesitan procesamiento unificado batch-stream.
Preguntas Comunes de Entrevista sobre Flink con Respuestas
¿Cómo logra Flink la semántica exactly-once?
Flink combina checkpointing con two-phase commit para sinks que soportan transacciones. Durante un checkpoint, Flink captura una instantánea del estado de los operadores y registra los offsets de las fuentes. Para sinks de Kafka, Flink pre-commitea registros a Kafka, completa el checkpoint, luego commitea la transacción. Si ocurre una falla antes de completar el checkpoint, los registros no commiteados son descartados y el procesamiento se reanuda desde el último checkpoint.
¿Cómo se explica la propagación de watermarks en una topología multi-fuente?
Cuando un job lee desde múltiples particiones o fuentes, cada una genera sus propios watermarks basados en los eventos entrantes. El watermark del operador downstream es igual al watermark mínimo de todos los canales de entrada. Esto garantiza que ninguna ventana se cierre prematuramente debido a que una partición rápida avance antes que las más lentas. La configuración de withIdleness() permite avanzar los watermarks cuando algunas particiones dejan de enviar datos.
¿Qué causa backpressure en Flink y cómo diagnosticarlo?
El backpressure ocurre cuando los operadores downstream no pueden seguir el ritmo de los datos upstream. La UI Web de Flink muestra el estado de backpressure por operador. Las causas comunes incluyen:
- Llamadas lentas a sistemas externos (consultas a base de datos, llamadas API)
- Cómputos costosos en funciones map/process
- Paralelismo insuficiente para el volumen de datos
- Operaciones de estado grandes bloqueando el procesamiento
Para resolverlo, se puede aumentar el paralelismo, optimizar operaciones lentas o usar I/O asíncrono para llamadas externas.
¿Cuál es la diferencia entre savepoints y checkpoints?
Los checkpoints son automáticos, incrementales y optimizados para recuperación ante fallas. Flink gestiona su ciclo de vida y elimina automáticamente los antiguos. Los savepoints son disparados por el usuario, son instantáneas completas destinadas a tareas operacionales: desplegar nuevo código, redimensionar el job o migrar entre clusters. Los savepoints persisten hasta eliminación explícita y soportan evolución de esquema.
Desplegando Flink en Kubernetes
El Flink Kubernetes Operator 1.15 simplifica el despliegue con recursos personalizados FlinkDeployment. Maneja el ciclo de vida de los jobs, actualizaciones y escalado.
# flink-deployment.yaml
# Kubernetes deployment for a Flink application
apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
name: fraud-detection-job
spec:
image: flink:1.20
flinkVersion: v1_20
flinkConfiguration:
taskmanager.numberOfTaskSlots: "4"
state.backend.type: rocksdb
state.checkpoints.dir: s3://flink-checkpoints/fraud-detection
execution.checkpointing.interval: "30s"
serviceAccount: flink
jobManager:
resource:
memory: "2048m"
cpu: 1
taskManager:
resource:
memory: "4096m"
cpu: 2
replicas: 3
job:
jarURI: s3://flink-artifacts/fraud-detection-1.0.jar
entryClass: com.example.FraudDetectionJob
parallelism: 12
upgradeMode: savepointLa configuración upgradeMode: savepoint garantiza que el operador tome un savepoint antes de actualizar, preservando el estado entre despliegues.
Optimizando Aplicaciones Flink para Producción
Los despliegues en producción requieren atención al paralelismo, memoria y configuración del state backend. Consultar los patrones ETL y pipelines de datos para consideraciones de integración.
// Production-ready Flink configuration
public class ProductionConfig {
public static void configureForProduction(StreamExecutionEnvironment env) {
// Checkpoint configuration
CheckpointConfig checkpointConfig = env.getCheckpointConfig();
checkpointConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
checkpointConfig.setMinPauseBetweenCheckpoints(30000); // 30 seconds
checkpointConfig.setCheckpointTimeout(600000); // 10 minutes
checkpointConfig.setMaxConcurrentCheckpoints(1);
checkpointConfig.setExternalizedCheckpointRetention(
ExternalizedCheckpointRetention.RETAIN_ON_CANCELLATION);
// State backend with incremental checkpoints
EmbeddedRocksDBStateBackend rocksDB = new EmbeddedRocksDBStateBackend(true);
rocksDB.setDbStoragePath("/tmp/rocksdb");
env.setStateBackend(rocksDB);
// Restart strategy with exponential backoff
env.setRestartStrategy(RestartStrategies.exponentialDelayRestart(
Duration.ofSeconds(1), // Initial delay
Duration.ofMinutes(5), // Max delay
2.0, // Backoff multiplier
Duration.ofHours(1), // Reset backoff after
0.1 // Jitter
));
}
}Los checkpoints incrementales reducen el tamaño de los checkpoints al escribir solo el estado modificado desde el último checkpoint. Esta optimización se vuelve crítica al manejar gigabytes de estado keyed.
¡Empieza a practicar!
Pon a prueba tu conocimiento con nuestros simuladores de entrevista y tests técnicos.
Puntos Clave del Procesamiento de Streams con Apache Flink
- Flink 2.3 procesa eventos individualmente con latencia de milisegundos, a diferencia de los sistemas micro-batch
- La semántica event time con watermarks maneja correctamente datos desordenados, respondiendo a la pregunta común de entrevista sobre eventos tardíos
- El estado keyed particiona datos para procesamiento paralelo mientras mantiene registros relacionados juntos
- Los checkpoints proporcionan garantías exactly-once mediante instantáneas distribuidas y two-phase commit
- El Kubernetes Operator automatiza despliegue, escalado y actualizaciones con preservación de estado basada en savepoints
- Flink es preferible sobre Spark Streaming cuando la latencia sub-segundo o los patrones de eventos complejos son importantes
- Se recomienda configurar RocksDB con checkpoints incrementales para workloads de producción con estado voluminoso
¿Sabrías detectar el bug en Data Engineering?
Un fragmento real, un bug oculto, un intento al día. Sin cuenta para probar.

Escrito por
Anthony Fillion-MailletFundador de SharpSkill
Desarrollador fullstack desde hace más de 10 años. Dirige SharpSkill y responde por todo lo que se publica aquí.
Actualizado el 28 de agosto de 2026
Etiquetas
Compartir
Artículos relacionados

Apache Beam vs Spark en 2026: Pipelines Unificados y Preguntas de Entrevista
Comparación detallada entre Apache Beam y Spark para pipelines de datos en 2026. Portabilidad, windowing, rendimiento y preguntas técnicas de entrevista.

Snowflake en 2026: arquitectura, SQL y preguntas de entrevista para ingenieros de datos
Guía 2026 sobre la arquitectura de Snowflake para ingenieros de datos: cómo se separan el almacenamiento y el cómputo, cómo funcionan los virtual warehouses y las micro-partitions, y las preguntas de entrevista que ponen a prueba la experiencia en producción.

Top 25 Preguntas de Entrevista para Ingenieros de Datos en 2026
Guía completa con las 25 preguntas más importantes para entrevistas de ingeniería de datos en 2026. Incluye SQL avanzado, pipelines ETL/ELT, streaming con Kafka, Spark, orquestación y arquitecturas lakehouse.