Мы создаем приложение Flink с использованием Java, в котором мы читаем два потока данных из двух разных тем Kafka и используем KeyedCoProcessFunction, чтобы найти записи, присутствующие в одном из потоков и отсутствующие в другом. Мы используем две переменные ValueState для хранения состояний каждого из потоков. В методеprocessElement1 мы сохраняем значение второй переменной ValueState и проверяем, не равно ли оно нулю. Если оно равно нулю, мы обновляем первую переменную ValueState объектом, который был передан в качестве параметра методуprocessElement1, и регистрируем службу таймера, чтобы проверить, поступило ли значение для второй переменной или нет. Точно так же мы реализуем методprocessElement2, но здесь мы проверяем, имеет ли первая переменная ValueState значение null или нет, и если это так, мы обновляем вторую переменную ValueState с помощью объекта, который был передан в качестве параметра методуprocessElement2, и мы регистрация службы таймера, чтобы проверить, поступило ли значение для первой переменной или нет.
Мой вопрос: существует ли какой-либо порядок, в котором выполняются функции ProcessElement? Если нет, можем ли мы как-то это контролировать?
Подробнее здесь: https://stackoverflow.com/questions/785 ... on-in-flin