Tag: APACHE-FLINK
Estoy utilizando la versión 2.34 de Beam y veo que todos los archivos jar de Flink son de la versión 1.13.2. Flink ha lanzado una nueva versión, la 1.13.5. ¿Cómo puedo actualizar a esa versión? No he visto que Beam haya lanzado una nueva versión.
Acabo de comenzar a aprender Flink y estoy probando el caso, “Informes en tiempo real con la API de Tablas“. Cuando ejecuté docker-compose, todos los contenedores funcionaron excepto el jobmanager, que se cerró con el código 2. todos activos cerrado con código 2 Intenté reconstruir y reiniciar, pero no funciona . . . Read more
Actualmente estamos utilizando flink 1.12 en modo de alta disponibilidad en producción. Hay 3 gestores de trabajos (1 líder y 2 en espera). Cuando subo un archivo JAR en uno de los gestores de trabajos, de alguna manera no se refleja en los demás gestores de trabajos. ¿Existe alguna forma . . . Read more
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
Estoy utilizando Flink SQL y el siguiente esquema muestra mis datos fuente (que pertenecen a algunos datos de Twitter): CREATE TABLE `twitter_raw` ( `entities` ROW( `hashtags` ROW( `text` STRING, `indices` INT ARRAY ) ARRAY, `urls` ROW( `indices` INT ARRAY, `url` STRING, `display_url` STRING, `expanded_url` STRING ) ARRAY, `user_mentions` ROW( `screen_name` . . . Read more