Вопросы по теме 'flink-streaming'

Порядок временного окна событий потоковой передачи Flink
У меня возникли проблемы с пониманием семантики оконного времени событий. Следующая программа генерирует несколько кортежей с отметками времени, которые используются в качестве времени события, и выполняет простую агрегацию окон. Я ожидаю, что вывод...
3179 просмотров
schedule 20.05.2022

Веб-интерфейс Fllink не отображает записи, полученные в реализации пользовательского источника
Я создал собственный источник для обработки потока журнала во Flink. Программа работает нормально и дает мне желаемые результаты после обработки записей. Но когда я проверяю веб-интерфейс, я не вижу счетчиков. Ниже приведен снимок экрана:...
1731 просмотров
schedule 27.04.2023

Apache Flink 1.0.0. Проблемы миграции, связанные со временем события
Я попытался перенести какую-то простую задачу на версию Flink 1.0.0, но это не удалось со следующим исключением: java.lang.RuntimeException: Record has Long.MIN_VALUE timestamp (= no timestamp marker). Is the time characteristic set to...
2914 просмотров
schedule 08.03.2023

Запрос данных из Apache Flink
Я собираюсь перейти с собственного потокового сервера на Apache Flink. Одна вещь, которая у нас есть, - это интерфейс DRPC, подобный Apache Storm, для выполнения запросов к состоянию, содержащемуся в топологии обработки. Так, например: у меня есть...
1050 просмотров
schedule 14.04.2024

flink Stream NoSuchMethodError: org.apache.flink.api.common.ExecutionConfig.setRestartStrategy
java.lang.NoSuchMethodError: org.apache.flink.api.common.ExecutionConfig.setRestartStrategy (Lorg / apache / flink / api / common / restartstrategy / RestartStrategies $ RestartStrategyConfiguration;) в com.WriteIntoKafka.main...
1234 просмотров
schedule 22.04.2022

Почему используется только один экземпляр GlobalWindow?
Посмотрите на этот пример : // We create sessions for each id with max timeout of 3 time units DataStream<Tuple3<String, Long, Integer>> aggregated = source .keyBy(0) .window(GlobalWindows.create())...
477 просмотров
schedule 06.02.2023

Несоответствие ввода: ожидается тип кортежа при попытке выбора в PatternStream
У меня проблемы с тестированием новых функций Flink 1.0.0. Я возился с CEP, и мне еще не удалось запустить простой демонстрационный код: val pattern : Pattern[TrafficEvent, _] = Pattern.begin[TrafficEvent]("start") val patternStream =...
318 просмотров
schedule 02.05.2023

Как выполнить timeWindow () для String DataStream во Flink?
Я хочу создать временное окно для потоковой передачи данных в Apache Flink. Мои данные выглядят примерно так: 1> {52,"mokshda",84.85} 2> {1,"kavita",26.16} 2> {131,"nidhi",178.9} 3> {2,"poorvi",22.97} 4> {115,"saheba",110.41}...
976 просмотров
schedule 29.03.2023

Использование Apache Flink и RxJava
В настоящее время я использую Apache flink и использую RxJava внутри него, мои вопросы: использование обоих из них уместно? потому что мои операции flink всегда являются функциями карты, и внутри них я интенсивно использую Rx, например, беру кортежи...
737 просмотров
schedule 24.05.2023

не удалось найти неявное значение для параметра доказательства
Я пишу простую задачу по подсчету слов, но продолжаю получать эту ошибку: could not find implicit value for evidence parameter of type org.apache.flink.api.common.typeinfo.TypeInformation[String] [error] .flatMap{_.toLowerCase.split("\\W+")...
3425 просмотров
schedule 08.06.2022

Flink timeWindow получает время начала
Я рассчитываю количество (сумма 1) по временному окну следующим образом: mappedUserTrackingEvent .keyBy("videoId", "userId") .timeWindow(Time.seconds(30)) .sum("count") Я хотел бы добавить время начала окна...
2813 просмотров
schedule 20.04.2023

Разница между Apache Flume и Apache Flink
Мне нужно прочитать поток данных из какого-то источника (в моем случае это поток UDP, но это не имеет значения), преобразовать каждую запись и записать ее в HDFS. Есть ли разница между использованием для этой цели Flume или Flink ? Я знаю,...
4074 просмотров

ClassNotFoundException: org.apache.flink.streaming.api.checkpoint.CheckpointNotifier при использовании темы kafka
Я использую новейшие банки Flink-1.1.2-Hadoop-27 и flink-connector-kafka-0.10.2-hadoop1. Потребитель Flink выглядит следующим образом: StreamExecutionEnvironment env=StreamExecutionEnvironment.getExecutionEnvironment(); if (properties...
7334 просмотров

Использование Grok в потоковой передаче flink
Flink Pipeline выглядит следующим образом: читать сообщения (строку) из темы кафка. сопоставление с образцом через преобразование grok в формат json. Агрегации по временному окну по извлеченному из json полю. Ниже приведен код для...
299 просмотров

Программа потоковой передачи Flink работает правильно со временем обработки, но не дает результатов со временем события
Обновление добавлено env.getConfig().setAutoWatermarkInterval(1000L); не устранил проблему. Думаю, проблема в другой части моего кода. Итак, во-первых, немного больше предыстории. Программа использует поток JSON смешанных типов сообщений...
1770 просмотров
schedule 04.11.2023

Окно, не достигающее своей длины окна
Я пробовал примеры работы с окнами flink, и чтобы проверить время окна, я добавил метку времени к событию потока. И я обнаружил, что продолжительность окна меньше, чем длина окна. Также, если бы мне пришлось использовать скользящее окно и изменить...
76 просмотров
schedule 23.03.2022

Flink записывает SingleOutputStreamOperator в два файла вместо одного
Я пытаюсь flink для проекта на работе. Я дошел до того, что обрабатываю поток, применяя окно подсчета и т. д. Однако я заметил странное поведение, которое не могу объяснить. Создается впечатление, что поток обрабатывается двумя потоками, и вывод...
739 просмотров
schedule 23.04.2023

является ли поток flink неизменным?
Мне нужно разделить поток данных с помощью flink. первый с именем "myDs" - содержит повторяющиеся данные второй, названный "goodDataStream", должен фильтровать дубликаты. частичный код: goodDataStream = myDs .filter( new...
137 просмотров
schedule 16.02.2023

мигающее окно для отметки времени
У меня есть поток данных вроде Eventname, Event id, Start_time ( Time Stamp) .. здесь я хочу применить преобразование окна к последнему полю Start_time , которое имеет отметку времени, мое требование такое, как будто я хочу взять данные за...
640 просмотров

Потоковая передача flink создать файл (csv или текст) во временном окне
Я новичок в флинке, у меня есть трансформация, предположим, что это val supportTask= customSource .map( line => line.split(",")) .map( line =>...
686 просмотров