задний план
В процессе построения центра обработки данных типичным сценарием интеграции данных является импорт данных MQ (очередь сообщений, таких как Kafka, RocketMQ и т. д.) в Hive для последующего построения хранилища данных и статистики индикаторов. Поскольку MQ-Hive — это первый уровень построения хранилища данных, требования к точности данных и производительности в режиме реального времени относительно высоки.
В этом документе основное внимание уделяется сценарию MQ-Hive, предлагается решение в реальном времени на основе Flink для решения проблем существующих решений в ByteDance и описывается текущий статус использования нового решения в ByteDance.
Существующие решения и болевые точки
Существующее решение в ByteDance показано на рисунке ниже, которое в основном разделено на два этапа:
- Запись данных MQ в файл HDFS через службу дампа
- Затем импортируйте данные HDFS в Hive через пакетный ETL и добавьте разделы Hive.
Болевые точки
- Цепочка задач длинная, и исходные данные должны пройти несколько преобразований, прежде чем попасть в Hive.
- Производительность в реальном времени относительно низкая, а задержка службы дампа и пакетного ETL приведет к задержке вывода окончательных данных.
- Затраты на хранение и вычисления высоки, а данные MQ многократно сохраняются и рассчитываются.
- На основе нативной Java после того, как трафик данных продолжает расти, возникают такие проблемы, как единая точка отказа и несбалансированная загрузка машины.
- Стоимость эксплуатации и обслуживания высока, а существующая инфраструктура, такая как Hadoop/Flink/Yarn, не может быть повторно использована в компании.
- Не поддерживает удаленное аварийное восстановление
Решение реального времени на основе Flink
Преимущество
Ввиду текущих недостатков традиционных решений компании мы предлагаем решение для работы в реальном времени на основе Flink, которое записывает данные MQ в Hive в режиме реального времени и поддерживает время события и семантику Exactly Once. По сравнению со старой схемой преимущества новой схемы заключаются в следующем:
- Разработан на основе потокового движка Flink, поддерживает семантику Exactly Once.
- Более высокая производительность в режиме реального времени, данные MQ напрямую поступают в Hive, без промежуточных вычислительных каналов.
- Уменьшите промежуточное хранилище, все данные процесса будут сохранены только один раз
- Поддерживает режим развертывания Yarn для облегчения миграции пользователей.
- Гибкое управление ресурсами, удобное для расширения, эксплуатации и обслуживания
- Поддержка аварийного восстановления в двух комнатах
Общая структура
Общая архитектура показана на рисунке ниже, который в основном включает три модуля: Источник DTS (Служба передачи данных), DTS Core и DTS Sink. Конкретные функции заключаются в следующем:
- DTS Source имеет доступ к различным источникам данных MQ, поддерживает Kafka, RocketMQ и т. д.
- DTS Sink выводит данные в целевой источник данных, поддерживает HDFS, Hive и т. д.
- DTS Core выполняет весь процесс синхронизации данных, считывает исходные данные через Source, обрабатывает их через DTS Framework и, наконец, выводит данные в цель через Sink.
- DTS Framework объединяет основные функции, такие как система типов, сегментация файлов, Exactly Once, сбор информации о задачах, время события и сбор грязных данных.
- Поддержка режима развертывания Yarn, более гибкое планирование ресурсов и управление ими.
Exactly Once
Платформа Flink может обеспечивать семантику Exactly Once или At минимум Once с помощью механизма Checkpoint. Чтобы реализовать полную связь MQ-Hive для поддержки семантики Exactly-once, также необходимо, чтобы MQ Source и Hive Sink поддерживали семантику Exactly-once. Эта статья реализована через протокол Checkpoint + 2PC Конкретный процесс выглядит следующим образом:
- Когда данные записываются, сторона источника извлекает данные из вышестоящего MQ и отправляет их стороне приемника; сторона приемника записывает данные во временный каталог.
- На этапе моментального снимка контрольной точки сторона источника сохраняет смещение MQ в состоянии, сторона приемника закрывает записанный дескриптор файла и сохраняет текущий идентификатор контрольной точки в состоянии;
- На этапе завершения контрольной точки сторона источника фиксирует смещение MQ; сторона приемника перемещает данные из временного каталога в официальный каталог.
- На этапе восстановления контрольной точки загрузите последний успешный каталог контрольной точки и восстановите информацию о состоянии.Исходная сторона использует смещение MQ, сохраненное в состоянии, в качестве начальной позиции; приемная сторона восстанавливает последний успешный идентификатор контрольной точки и перемещает данные во временную каталог на официальный.
добиться оптимизации
В реальных сценариях использования, особенно в сценариях большого параллелизма, задержка записи HDFS подвержена сбоям, поскольку время ожидания отдельных моментальных снимков задачи истекает или происходит сбой, что приводит к сбою всей контрольной точки. Поэтому при сбое Checkpoint важнее повысить отказоустойчивость и стабильность системы.
Это полностью использует строго монотонно возрастающий идентификатор контрольной точки. Каждый раз, когда выполняется контрольная точка, текущий идентификатор контрольной точки должен быть больше, чем раньше. Поэтому на этапе завершения контрольной точки временные данные, меньшие или равные текущему идентификатору контрольной точки, могут быть Отправлено. Конкретная стратегия оптимизации выглядит следующим образом:
- Временный каталог на стороне приемника — {dump_path}/{next_cp_id}, где определение next_cp_id — это текущий последний cp_id + 1.
- На этапе моментального снимка контрольной точки сторона приемника сохраняет текущий последний cp_id в состояние и в то же время обновляет next_cp_id до cp_id + 1.
- На этапе завершения контрольной точки сторона приемника перемещает все данные, меньшие или равные текущему cp_id во временном каталоге, в официальный каталог.
- На этапе восстановления контрольной точки сторона приемника восстанавливает последний успешный cp_id и перемещает данные, меньшие или равные текущему cp_id во временном каталоге, в официальный каталог.
система типов
Поскольку типы данных, поддерживаемые разными источниками данных, различны, для решения проблем синхронизации данных между разными источниками данных и совместимости различных преобразований типов мы поддерживаем систему типов DTS.Типы DTS могут быть уточнены в базовые типы и составные типы. . Поддерживается вложенность типов. Конкретный процесс преобразования выглядит следующим образом:
- На стороне источника исходный тип данных единообразно преобразуется в тип DTS внутри системы.
- На стороне приемника преобразуйте тип DTS внутри системы в тип целевого источника данных.
- Среди них система типов DTS поддерживает преобразование между различными типами, например преобразование между типом String и типом Date.
Rolling Policy
Сторона приемника записывается одновременно, и трафик, обрабатываемый каждой задачей, отличается. Чтобы избежать создания слишком большого количества маленьких файлов или создания слишком больших файлов, необходимо поддерживать пользовательскую стратегию сегментации файлов для управления размером одного файла. . В настоящее время поддерживаются три стратегии сегментации файлов: размер файла, максимальное время хранения файла без обновления и контрольная точка.
Стратегия оптимизации
Hive поддерживает различные форматы хранения, такие как Parquet, Orc, Text и т. д. Процесс записи данных разных форматов хранения неодинаков, и его можно разделить на две категории:
- RowFormat: основан на одной записи, поддерживает операцию усечения HDFS в соответствии со смещением, например в текстовом формате.
- BulkFormat: запись на основе блоков, не поддерживает операции усечения HDFS, такие как Parquet, форматы ORC.
Чтобы обеспечить семантику Exactly Once и одновременно поддерживать Parquet, Orc, Text и другие форматы, в каждой контрольной точке сегментация файлов является обязательной, чтобы гарантировать, что все записываемые файлы будут завершены, и операция усечения не требуется, когда контрольная точка восстанавливается.
Отказоустойчивость
В идеале потоковые задачи будут выполняться все время без перезапуска, но на практике неизбежно возникнут следующие сценарии:
- Вычислительный движок Flink обновлен, и задачу необходимо перезапустить.
- Объем восходящих данных увеличивается, и параллелизм задач необходимо скорректировать.
- Task Failover
Регулировка параллелизма
В настоящее время Flink изначально поддерживает изменение масштаба состояния. В конкретной реализации, когда Задача выполняет моментальный снимок контрольной точки, смещение MQ сохраняется в ListState; после перезапуска задания мастер заданий равномерно распределяет ListState для каждой задачи в соответствии с параллелизмом оператора.
Task Failover
Из-за влияния внешних факторов, таких как дрожание сети, тайм-аут записи и т. д., Task неизбежно не сможет записать, и более важно, как быстро и точно выполнить Task Failover. В настоящее время Flink изначально поддерживает различные стратегии аварийного переключения задач.В этой статье используется стратегия аварийного переключения региона для перезапуска всех задач в регионе, в котором находится сбойная задача.
Удаленное аварийное восстановление
задний план
В эпоху больших данных точность и характер данных в режиме реального времени особенно важны. В этой статье представлены решения для развертывания нескольких машинных залов и удаленного аварийного восстановления.Когда основное машинное помещение временно не может предоставлять внешние услуги из-за отключения сети, сбоя питания, землетрясения, пожара и т. д., услугу можно быстро переключить на комната аварийного восстановления, и в то же время может быть гарантирована семантика Exactly Once.
Компоненты аварийного восстановления
Общее решение требует совместной работы нескольких компонентов аварийного восстановления. Компоненты аварийного восстановления показаны на рисунке ниже, в основном это MQ, YARN и HDFS, а именно:
- MQ должен поддерживать развертывание в нескольких комнатах.Если основная комната выходит из строя, лидер может быть переключен на резервную комнату для последующего потребления.
- Кластеры пряжи развернуты в основной комнате и резервной комнате для переноса заданий Flink.
- Нисходящая HDFS должна поддерживать развертывание с несколькими машинными залами.При выходе из строя основного машинного зала главный может быть переключен на резервный машинный зал.
- Задание Flink выполняется на Yarn, а серверная часть задачи State сохраняется в HDFS, а многомашинная комната State Backend гарантируется благодаря поддержке HDFS с несколькими машинными комнатами.
Процесс аварийного восстановления
Общий процесс аварийного восстановления выглядит следующим образом:
- В обычных условиях MQ Leader и HDFS Master развернуты в главном компьютерном зале и синхронизируют данные с резервным компьютерным залом. В то же время задание Flink выполняется в основной комнате и записывает состояние задачи в HDFS.Обратите внимание, что это состояние также является режимом развертывания в нескольких комнатах.
- В случае аварии MQ Leader и HDFS Master мигрируют из главного машинного отделения в комнату подготовки к авариям, а задание Flink также переносится в комнату подготовки к авариям и восстанавливает информацию о смещении до аварии через состояние, чтобы предоставить точное Однажды семантика
архив времени события
задний план
При построении хранилища данных логика обработки времени процесса и времени события отличается: для времени обработки данные будут записываться в раздел времени, соответствующий текущему системному времени, а для времени события — запись данных в соответствующий раздел времени. в соответствии со временем производства данных, что также называется архивированием в этой статье.
В реальных сценариях неизбежно будут возникать различные сбои в восходящем и нисходящем направлении, и они будут восстановлены через определенный период времени.Если принята стратегия обработки времени обработки, данные во время аварии будут записаны в восстановленный раздел времени, который в конечном итоге приведет к проблеме дыр в разделах или дрейфа данных; если будет принята стратегия архивирования, она будет записана в соответствии со временем события, и такой проблемы нет.
Так как время восходящих событий данных будет нарушено, и раздел Hive не должен продолжать записываться после его создания, бесконечно архивировать во время фактического процесса записи невозможно, а только в пределах определенного диапазона времени. Сложность архивирования заключается в том, как определить глобальное минимальное время архивирования и как допустить определенный беспорядок.
Глобальное минимальное время архива
Исходная сторона читается одновременно, и задача может одновременно считывать данные из нескольких разделов MQ.Для каждого раздела MQ будет сохранено текущее время архива раздела, а минимальное значение в разделе будет принято в качестве минимальное время архивации задачи, и, наконец, будут взяты данные в задаче Минимальное значение, как глобальное минимальное время архивации.
Обработка вне очереди
Для поддержки неупорядоченных сценариев будет поддерживаться настройка интервала архивирования, в котором глобальный минимальный водяной знак — это глобальное минимальное время архивирования, водяной знак раздела — текущее время архивирования раздела, а минимальный водяной знак раздела — это минимальное время архивации раздела, только когда время события соответствует следующим условиям, будет заархивировано:
- Время события превышает глобальное минимальное время архива.
- Время события превышает минимальное время архивации раздела.
Генерация разделов улья
принцип
Сложность создания раздела Hive заключается в том, как определить, готовы ли данные раздела и как добавить раздел. Поскольку сторона приемника записывает одновременно, и будет несколько задач, записывающих одни и те же данные раздела одновременно, данные раздела можно считать готовыми только после завершения записи всех данных раздела задачи.Решение в этой статье выглядит следующим образом:
- На стороне приемника, чтобы каждая задача сохраняла текущее минимальное время обработки, она должна удовлетворять характеристике монотонного увеличения.
- Когда Контрольная точка завершена, Задача сообщает стороне JM минимальное время обработки.
- После того, как JM получит минимальное время обработки всех задач, он может получить глобальное минимальное время обработки и использовать его в качестве минимального времени готовности раздела Hive.
- Когда минимальное время готовности будет обновлено, вы сможете решить, добавлять ли раздел Hive.
динамический раздел
Динамическое разбиение определяет, в какой каталог раздела записываются данные, на основе значения исходных входных данных, а не в фиксированный каталог раздела. Например, в сценарии date={date}/hour={hour} /app={app}, в зависимости от времени раздела И значение поля app определяет конечный каталог раздела, так что одни и те же данные приложения находятся в одном и том же разделе каждый час.
В сценарии со статическим разделом каждая задача записывает только один файл раздела за раз, но в сценарии с динамическим разделом каждая задача может одновременно записывать в несколько файлов раздела. Для записи в формате Parque данные сначала будут записываться в локальный кеш, а затем пакетно записываются в Hive.Когда задача одновременно обрабатывает слишком много файловых дескрипторов, может возникнуть OOM. Чтобы предотвратить однозадачный OOM, дескриптор файла будет периодически проверяться, и дескриптор файла, который не был записан в течение длительного времени, будет вовремя освобожден.
Messenger
Модуль Messenger используется для сбора информации о состоянии выполнения задания для измерения состояния задания и построения рыночных индикаторов.
Сбор метаинформации
Принцип сбора метаинформации заключается в следующем: на стороне приемника основные показатели задачи, такие как трафик, QPS, грязные данные, задержка записи, эффект записи времени события и т. д., собираются через Messenger и агрегируются. через коллектор мессенджеров. Грязные данные нужно выводить на внешнее хранилище, а индикаторы выполнения задач выводятся в Grafana для отображения масштабных индикаторов.
сбор грязных данных
В сценариях интеграции данных неизбежно столкновение с грязными данными, такими как типы ошибок конфигурации, переполнения полей и несовместимые преобразования типов. Для потоковых задач, поскольку задача будет выполняться все время, она должна иметь возможность подсчитывать трафик грязных данных в режиме реального времени, сохранять грязные данные во внешнем хранилище для устранения неполадок и отбирать выходные данные в журнале выполнения.
Мониторинг рынка
Крупномасштабные индикаторы охватывают глобальные индикаторы и индикаторы отдельных заданий, в том числе успешный трафик записи и число запросов в секунду, задержку записи, ошибочный трафик записи и количество запросов в секунду, а также статистику эффекта архивирования, как показано на следующем рисунке:
будущий план
В компании запущено и продвигается решение для работы в режиме реального времени на базе Flink, в дальнейшем мы сосредоточимся в основном на следующих аспектах:
- Расширенные функции интеграции данных, поддержка доступа к большему количеству источников данных, поддержка определяемой пользователем логики преобразования данных и т. д.
- Озеро данных открыто для поддержки импорта данных CDC в реальном времени
- Унифицированная потоковая пакетная архитектура поддерживает полную и инкрементную интеграцию данных сцены.
- Обновление архитектуры для поддержки большего количества сред развертывания, таких как K8S.
- Улучшить обслуживание, снизить стоимость доступа пользователей
Суммировать
Благодаря постепенной диверсификации и быстрому развитию бизнес-продуктов ByteDance внутренняя универсальная платформа разработки больших данных ByteDance становится все более и более функциональной и предоставляет глобальные решения для интеграции данных в автономном режиме, в режиме реального времени, в полном и поэтапном сценариях. первые несколько сотен задач выросли до десятков тысяч, а ежедневная обработка данных достигла уровня PB.Решение реального времени на основе Flink активно продвигалось и использовалось внутри компании, а старая связь MQ-Hive постепенно заменены.
В настоящее время ByteDance по-прежнему развивается с большой скоростью, мы по-прежнему сохраняем свое первоначальное намерение и движемся вперед. Если вам интересна эта статья или сценарии разработки больших данных, включая, помимо прочего, платформы разработки, системы планирования, интеграцию данных, платформы метаданных и качества данных и т. совместное развитие бизнеса ByteDance Проблемы в области эффективности построения данных и управления данными в стадии разработки. Добро пожаловать, чтобы отправить свое резюме, как в Пекине, так и в Шанхае.Адрес электронной почты: dataplatform-hr@bytedance.com.
использованная литература
-
Real-time Exactly-once ETL with Apache Flink
http://shzhangji.com/blog/2018/12/23/real-time-exactly-once-etl-with-apache-flink/
-
Implementing the Two-Phase Commit Operator in Flink
https://flink.apache.org/features/2018/03/01/end-to-end-exactly-once-apache-flink.html
-
A Deep Dive into Rescalable State in Apache Flink
https://flink.apache.org/features/2017/07/04/flink-rescalable-state.html
-
Data Streaming Fault Tolerance
https://ci.apache.org/projects/flink/flink-docs-release-1.9/internals/stream_checkpointing.html