es.davy.ai

Preguntas y respuestas de programación confiables

¿Tienes una pregunta?

Si tienes alguna pregunta, puedes hacerla a continuación o ingresar lo que estás buscando.

Tag: FLINK-STREAMING

Flink SlidingEventTimeWindows no funciona como se espera.

Tengo una ejecución de transmisión configurada como object FlinkSlidingEventTimeExample extends App { case class Trx(timestamp:Long, id:String, trx:String, count:Int) val env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI() val watermarkS1 = WatermarkStrategy .forBoundedOutOfOrderness[Trx](Duration.ofSeconds(15)) .withTimestampAssigner(new SerializableTimestampAssigner[Trx] { override def extractTimestamp(element: Trx, recordTimestamp: Long): Long = element.timestamp }) val s1 = env.socketTextStream(“localhost”, 9999) .flatMap(l => l.split(” “)) .map(l . . . Read more

El estado de checkpoint de Flink siempre está en progreso.

Utilizo el conector de flujo de datos KafkaSource y la función HbaseSinkFunction para consumir datos de Kafka y escribirlos en HBase. Habilito el punto de control de esta forma: env.enableCheckpointing(3000,CheckpointingMode.EXACTLY_ONCE); Los datos en Kafka ya se han escrito correctamente en HBase, pero el estado de los puntos de control en . . . Read more

Creación y envío dinámico de trabajos a Flink

Hola, estoy planeando usar Flink como backend para mi función en la que mostraremos una interfaz de usuario al usuario para crear patrones de eventos de forma gráfica, por ejemplo: múltiples intentos de inicio de sesión fallidos desde la misma dirección IP. Crearemos el patrón de Flink programáticamente utilizando los . . . Read more