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

Kafka и Flink дублируют сообщения при перезапуске
Прежде всего, это очень похоже на Кафка снова использует последнее сообщение, когда я перезапускаю клиент Flink , но это не то же самое. Ответ на этот вопрос, похоже, НЕ решает мою проблему. Если я что-то пропустил в этом ответе, перефразируйте...
1961 просмотров

Flink CEP Результаты не распечатаны
Я пытаюсь распечатать строку, если Hello и world были найдены с помощью библиотеки Flink CEP. Мой источник - Kafka и использую консоль-производителя для ввода данных. Эта часть работает. Могу распечатать то, что ввожу в тему. Однако он не...
592 просмотров
schedule 21.04.2024

java.io.NotSerializableException с использованием Apache Flink с Lagom
Я пишу программу Flink CEP внутри реализации микросервиса Lagom. Моя программа FLINK CEP отлично работает в простом приложении Scala. Но когда я использую этот код внутри реализации службы Lagom, я получаю следующее исключение Внедрение...
800 просмотров
schedule 12.07.2022

Когда я запускаю задание flink для сохранения данных в Azure Data Lake, я получаю указанное ниже исключение. Может ли кто-нибудь помочь мне в этом?
I am using concurrend append method from the class Core in Azure to store data to Azure Data lake.Below is the code and the exception which I got.I am getting this exception rarely not always.Could anyone guide me?... public void invoke(String...
543 просмотров

Невозможно использовать серверную часть состояния Flink RocksDB из-за несоответствия порядка байтов
My Flink Job читает из темы kafka и сохраняет данные в серверной части состояния RocksDB, чтобы использовать состояние запроса. Я могу запустить задание и запросить состояние на моем локальном компьютере. Но при развертывании в кластере я получаю...
1292 просмотров
schedule 28.09.2022

warnings.print () печатает события в обратном порядке (последнее событие - первое), за исключением первого события в Apache Flink CEP
Я пытаюсь отфильтровать все временные события, которые> 10 во Flink, используя шаблон ниже, Pattern<MonitoringEvent, ?> warningPattern = Pattern.<MonitoringEvent>begin("first") .subtype(TemperatureEvent.class)...
51 просмотров

Динамическое дросселирование источников flink kafka
Мы используем несколько тем о кафках, но хотим отдать приоритет некоторым из них (~ Качество обслуживания). Согласно тому, что я нашел в Интернете, консенсус заключается в том, чтобы ограничивать не операторы, а источник, а точнее десериализатор...
1171 просмотров
schedule 18.10.2023

Apache Flink: невозможно использовать writeAsCsv() с потоком данных кортежа подкласса
Как рекомендуется здесь: Передовой опыт: присвоение имен большим типам TupleX . Я использую POJO вместо кортежа для своего потока данных. Вот как определил мой POJO: public class PositionEvent extends Tuple8<Integer, String, Integer,...
478 просмотров
schedule 17.11.2022

Повторное использование потока является копией потока или нет
Например, есть ключевой поток: val keyedStream: KeyedStream[event, Key] = env .addSource(...) .keyBy(...) // several transformations on the same stream keyedStream.map(....) keyedStream.window(....) keyedStream.split(....)...
1296 просмотров
schedule 30.12.2023

Запуск нескольких программ Flink в автономном кластере Flink (v1.4.2)
У меня есть автономный кластер Flink на основе Flink 1.4.2 (1 менеджер заданий, 4 слота задач), и я хочу отправить две разные программы Flink. Не уверен, возможно ли это вообще, поскольку в некоторых архивах flink говорится, что менеджер заданий...
741 просмотров
schedule 30.12.2023

Flink - сериализуемый класс (не POJO)
Невозможно сериализовать классы, которые не являются POJO в apache flink? У меня есть служебный класс, который имеет много функций, и я хочу отправить исходной функции объект этого служебного класса, но flink выдает исключение сериализации....
253 просмотров
schedule 15.03.2024

Flink Streaming: разница между TriggerResult.FIRE и TriggerResult.FIRE_AND_PURGE
Я новичок в Flink. У меня есть потоковая программа Flink, которая считает что-то от kafka в 10-секундных окнах сеанса. Вот мой вопрос: Триггер по умолчанию в окнах сеанса - ПОЖАР. Будет ли потоковая передача Flink сохранять в памяти все...
480 просмотров
schedule 26.03.2024

Окно сеанса Flink с триггером onEventTime?
Я хочу создать окно сеанса на основе EventTime в Flink, чтобы оно запускалось, когда время события нового сообщения более чем на 180 секунд превышает время события сообщения, создавшего окно. Например: t1(0 seconds) : msg1 <-- This is the...
906 просмотров
schedule 11.11.2023

Есть ли переменные класса контрольной точки / точки сохранения Flink?
Если приложение Flink запускается после сбоя или обновляется, сохраняются ли переменные класса, которые явно не являются частью KeyedState или OperatorState? Например, BoundedOutOfOrdernessGenerator, описанный в документации Flink, имеет переменную...
58 просмотров
schedule 20.02.2024

Flink: ошибка при запуске scala-shell — не удалось создать DispatcherResourceManagerComponent
Я использую Flink версии 1.10.0, и при запуске оболочки Scala с помощью « start-scala-shell.sh » возникает следующее исключение: Exception in thread "main" org.apache.flink.util.FlinkException: Could not create the...
230 просмотров
schedule 15.11.2023

Параметры кэширования / связывания данных Apache Flink
Это очень широкий вопрос, я новичок в Flink и изучаю возможность его использования в качестве замены существующей аналитической движке. Сценарий состоит в том, что данные, собранные с различного оборудования, получены в виде закодированной строки...
341 просмотров
schedule 18.03.2024

Flink: поддерживает ли Flink абстрактный оператор, который может обрабатывать разные потоки данных с общими полями?
Предположим, у нас есть несколько потоков данных, и у них есть общие черты. Например, у нас есть поток Учитель и поток Студент , и у них обоих есть поле возраст . Если я хочу найти старшего ученика или учителя из потока в реальном времени, я...
62 просмотров
schedule 04.04.2024

Как установить время для контрольной точки в потоковой передаче Apache Flink?
Я запускаю пример Apache Flink с детектором мошенничества с RocksDB в качестве серверной части состояния. Я хочу знать, сколько времени требуется Apache Flink для проверки состояния. Мой подход состоит в том, чтобы распечатать время до и после...
130 просмотров
schedule 08.04.2024

Проблема потребителя Apache Flink Kafka
У меня есть данные в Kafka, я хотел прочитать данные, отправляет ли Kafka данные или нет, фильтровать их и возвращать JSON. // create execution environment StreamExecutionEnvironment env =...
120 просмотров

Управление состоянием в RichFlatMap с ключом и без него.
У меня есть потоковое приложение, подобное этому: DataStream<MyObject> stream1 = source .keyBy("clientip") .flatMap(new MyFlatMapFunction()) .name("Stream1"); //... public...
57 просмотров
schedule 10.05.2024