Разработка моделей данных для обработки инженерных данных в режиме реального времени

Понимание обработки инженерных данных в реальном времени

В отличие от пакетной обработки, где данные собираются в течение периода, а затем обрабатываются оптом, обработка в реальном времени требует субсекундных задержек. Это различие имеет решающее значение в инженерных случаях использования, таких как прогнозное обслуживание, где задержка в анализе данных вибрации от турбины может привести к катастрофическому отказу; или в управлении интеллектуальными сетями, где колебания напряжения должны быть исправлены в течение миллисекунд, чтобы предотвратить отключения.

Для удовлетворения этих требований модели данных должны быть разработаны с глубоким пониманием скорости, разнообразия и объема данных. Считывания датчиков с устройств Интернета вещей (IoT) часто приходят к миллионам событий в секунду, каждое из которых содержит временные метки, идентификаторы и несколько измерений. Модель данных должна эффективно захватывать этот поток, минимизировать накладные расходы на хранение и обеспечивать быстрый поиск для последующей аналитики и оповещения.

Ключевые проблемы включают обработку неупорядоченных данных, управление событиями с поздним прибытием и обеспечение семантики обработки точно один раз, когда дубликаты не могут быть допущены. Хорошо разработанная модель данных абстрагирует эти сложности, обеспечивая чистый интерфейс для инженеров для запроса и визуализации данных в режиме реального времени.

Основные принципы проектирования моделей данных в системах реального времени

Разработка модели данных для инженерных данных в реальном времени требует балансирования компромиссов между несколькими основными принципами. Эти принципы определяют решения по схеме проектирования, двигателям хранения и шаблонам запросов.

Масштабируемость и эластичность

Модель данных должна масштабироваться горизонтально, чтобы вместить растущие объемы данных без ухудшения производительности. Это часто включает в себя разделение данных по нескольким узлам. Например, данные временных рядов могут быть разделены по временному диапазону или по хэшу идентификатора датчика. Эластичность позволяет системе автоматически добавлять или удалять узлы при изменении нагрузки, что особенно важно в инженерных средах, где всплески данных происходят во время экспериментов или наращивания производства.

Низкая задержка чтения и написания путей

Приложения в реальном времени требуют как операций записи, так и чтения для завершения в течение миллисекунд. Структуры данных, которые поддерживают только записи приложений, такие как деревья слияний с лог-структурой (LSM), распространены в базах данных, таких как InfluxDB или TimescaleDB. Для считываний модель должна поддерживать эффективное сканирование диапазона в окнах времени и точечные поиски для конкретных состояний устройства. Стратегии индексации, такие как использование индекса на основе времени в сочетании с индексом тегов для метаданных устройства, необходимы.

Последовательность и целостность данных

В инженерных условиях точность данных не подлежит обсуждению. Модель данных должна обеспечивать соблюдение ограничений по согласованности, таких как обеспечение того, чтобы показания температуры попадали в заданный диапазон. Стратегии разрешения конфликтов, такие как выигрыши в последней записи или векторы версий, применяются, когда данные поступают из нескольких источников. Однако возможная согласованность часто приемлема для контрольных приборных панелей, в то время как сильная согласованность обязательна для контуров управления, которые непосредственно приводят в действие машины.

Гибкость в использовании эволюционирующих схем

Инженерные проекты часто добавляют новые датчики, изменяют частоту выборки или вводят новые типы измерений. Жесткая, заранее заданная схема ломается при изменении данных. Гибкие модели данных, такие как подходы к схеме на считывании (например, с использованием JSONB в PostgreSQL или динамических столбцов в Кассандре), позволяют инженерам принимать данные без изменения схемы хранения. Альтернативно, используя базу данных временных рядов с гибкой моделью меток и полей (например, InfluxDB) обеспечивает хороший баланс между производительностью и адаптивностью.

Выбор правильных структур данных и двигателей хранения

Выбор структур данных напрямую влияет на способность системы обрабатывать данные в режиме реального времени. Ниже приведены наиболее часто используемые структуры в инженерных моделях данных, а также их компромиссы.

Базы данных по временным рядам

Базы данных временных рядов (TSDB) специально построены для хранения и запроса последовательных точек данных, индексируемых временем. Они обычно эффективно сжимают данные с использованием дельта-кодирования и кодирования длины выполнения, снижая затраты на хранение. TSDB также поддерживают политику сбора и хранения данных, которая автоматически собирает или удаляет старые данные. Например, при мониторинге парка ветряных турбин TSDB может хранить сырые данные в течение одной недели, затем образец для почасовых средних для долгосрочного анализа тенденций. Популярные TSDB включают TimescaleDB, InfluxDB и Prometheus.

Магазины с ключевыми ценностями

Хранилища ключевых значений отлично подходят для поиска в режиме реального времени состояния или конфигурации устройства. Они предлагают чрезвычайно низкую задержку для точечных считываний и записи. В моделях инженерных данных ключ часто представляет собой композицию идентификатора устройства и временной метки, в то время как значение представляет собой сериализованную каплю показаний датчиков. Однако хранилища ключевых значений менее эффективны для запросов диапазона на нескольких устройствах или окнах времени. Они лучше всего используются в качестве слоя кэша или для хранения последнего известного состояния каждого устройства.

Поток-процессинговые магазины

Такие технологии, как уплотненные темы Apache Kafka или хранилище состояния Apache Flink, позволяют обрабатывать и хранить данные в самом потоке. Эта архитектура уменьшает потребность в отдельных базах данных, когда основным вариантом использования является аналитика и оповещение в реальном времени. Например, модель данных, реализованная с использованием Kafka Streams, может поддерживать в локальном государственном магазине последние десять минут данных вибрации для каждой машины и вызывать оповещение, когда скользящая средняя превышает порог.

Гибридные подходы

Многие инженерные системы используют гибридную стратегию: используют потоковый процессор для аналитики в реальном времени, TSDB для исторического хранения и хранилище ключевых значений для текущего состояния. Эта архитектура обеспечивает низкую задержку для операционных приборных панелей, а также позволяет проводить глубокий исторический анализ. Модель данных должна определять, как данные перемещаются между этими слоями, часто используя захват данных об изменениях (CDC) или шаблоны двойной записи.

Стратегии проектирования для инженерных моделей данных

Эффективные модели данных для инженерных данных в реальном времени разработаны с конкретными стратегиями, которые учитывают уникальные ограничения домена.

Моделирование устройств и датчиков

Общий подход заключается в моделировании каждого физического устройства или датчика как отдельного объекта, который излучает поток событий измерения. В реляционной модели у вас может быть таблица с метаданными (местоположение, производитель, дата установки) и таблица с временем, типом датчика и значением. Однако в сценариях реального времени таблица измерений может быстро расти на миллиарды строк. Лучшая конструкция заключается в использовании модели временнóго ряда, где каждое измерение хранится в виде строки с временной меткой, идентификатором устройства и полезной нагрузкой пар ключевых значений для разных метрик. Эта структура уменьшает количество таблиц и позволяет эффективно сжимать.

Пример плоской записи измерений:

метка времени: 2025-03-09T14:30:01.234Z, device id: «сенсор-42», метрики: {»температура»: 68.2, «влажность»: 45.1, «давление»: 1013.2}

Нормализация vs. денормализация

Нормализация уменьшает избыточность данных и улучшает производительность записи путем хранения метаданных отдельно. Однако в системах реального времени частое присоединение потока измерений к метаданным устройства может вводить задержку. Денормализация часто предпочтительнее для запросов с горячим путем. Например, включение местоположения устройства непосредственно в строке измерения исключает соединение во время оповещения. Компромиссом является увеличенное хранение и потенциальная несоответствие при изменении метаданных устройства (например, перемещается датчик).

Общим шаблоном является использование нормализованной модели для холодного пути (аналитика) и денормализованной модели для горячего пути (панели приборов реального времени) с синхронными обновлениями для обоих через потоковый процессор.

Разделение и заточка

Разделение данных имеет решающее значение для масштабируемости. Разделение на основе времени является наиболее распространенным для данных временных рядов: каждый раздел охватывает определенный временной интервал (например, один час или один день). Это позволяет системе быстро сбрасывать старые разделы и эффективно выполнять запросы диапазона. Разделение на основе идентификатора устройства равномерно распределяет нагрузку по узлам, но это может привести к горячим точкам, если некоторые устройства генерируют гораздо больше данных, чем другие. Комбинация времени и хэша устройства работает хорошо.

Например, в Кассандре ключ раздела может быть , а ключ кластеризации - временная метка.

Индексация для эффективности запросов

Стратегии индексирования должны быть адаптированы к наиболее распространенным шаблонам запросов: «получить все данные для устройства X за последний час» или «найти все устройства, температура которых превышает 100°C в последнюю минуту». Типичным является индекс времени в сочетании с индексом тега устройства. Передовые методы включают использование индекса списка пропусков для баз данных временных рядов или индекса растровых изображений для тегов низкой сердечности. Избегайте переиндексации, поскольку она замедляет запись. Многие TSDB автоматически создают индекс времени на основной столбце временных меток.

Внедрение технологий обработки потоков

Модели инженерных данных в реальном времени часто строятся поверх структур обработки потоков, которые обеспечивают семантику, отказоустойчивость и управление состоянием. Ниже приведены ключевые технологии и то, как они влияют на дизайн модели данных.

Апач Кафка

Кафка выступает в качестве основы для приема данных. Модель данных для тем Kafka должна выровняться с потребителями нисходящего потока. Например, каждый тип устройства может иметь свою собственную тему, или все устройства могут разделять одну тему с разделом на группу устройств. Схема сообщений (например, Avro или Protobuf) включает в себя временную метку, идентификатор устройства и полезную нагрузку метрик. Компактность может быть включена для сохранения только последнего значения для каждого ключа, что полезно для обновлений состояния устройства.

Разъемы Kafka (Kafka Connect) могут подталкивать данные в TSDB или хранилище ключей без дополнительного кодирования.

Apache Flink

Flink обрабатывает потоковые данные с низкой задержкой и поддерживает вычисления с состоянием. Модель данных в Flink определяется типами событий и дескрипторами состояния. Например, для обнаружения аномальных вибрационных паттернов Flink поддерживает состояние, в котором хранятся последние 100 показаний ускорения на устройство. Модель данных должна быть спроектирована так, чтобы минимизировать размер состояния; использовать словари для идентификаторов датчиков и сжимать повторяющиеся поля. Flink также поддерживает обработку событийного времени, поэтому модель данных должна включать временную метку события (не временную метку обработки) для правильного оконного считывания.

Apache Spark скачать

Spark Streaming (или Structured Streaming) обрабатывает данные в микро-пакетах. Модель данных может быть представлена в виде DataFrame или Dataset, со схемами, определенными в коде. В то время как микро-batching вводит более высокую задержку, чем чистая потоковая передача (например, Flink), ее легче использовать для аналитических рабочих нагрузок, которые должны объединять потоки с историческими таблицами. Модель данных должна учитывать механизм контрольных точек, который Spark использует для поддержания точной семантики один раз, который записывает состояние в каталог контрольных точек.

Интеграция баз данных

Потоковые процессоры часто пишут в базу данных в реальном времени. Модель данных должна определять отображение от потока событий к схеме базы данных. Например, задание Flink считывает необработанные данные датчиков от Kafka, применяет некоторую фильтрацию и пишет в InfluxDB с использованием своего линейного протокола. Имена измерений, теги и поля схемы базы данных должны быть разработаны для соответствия запросам, которые будут выполнять панели инструментов. Избегайте слишком большого количества тегов, потому что они могут ухудшить производительность записи; предпочитайте поля для непрерывно меняющихся метрик.

Пример: Модель данных для системы прогнозного обслуживания в реальном времени

Рассмотрим завод с 10 000 машин, каждая из которых оснащена датчиками, измеряющими температуру, вибрацию и скорость вращения. Цель состоит в том, чтобы за 30 минут заранее предсказать сбои и вызвать оповещения об обслуживании.

Модель данных разработана следующим образом:

Эта гибридная модель уравновешивает потребность в оповещениях с низкой задержкой (через обработку потока) с гибким историческим анализом (через базу данных временных рядов). Модель данных остается простой: одна гипертаблица для необработанных данных с индексами, оптимизированными для наиболее распространенной схемы запросов (диапазон времени + идентификатор машины).

Лучшие практики для развертывания производства

Переход от проектирования к производству требует внимания к мониторингу, эволюции схем и управлению затратами.

Монитор и профиль производительность запроса

Используйте инструменты, специфичные для баз данных (например, TimescaleDB's , InfluxDB's query Inspector) для выявления медленных запросов. Монитор записывает пропускную способность и задержку; если записывает всплески задержки, рассмотрите возможность увеличения количества разделов или настройки стратегии уплотнения. Настройте оповещения для тайм-аутов запроса.

План эволюции схемы

Часто меняются схемы инженерных данных. Для управления схемами Avro или Protobuf используются реестры схем (например, реестр сменных схем), для баз данных, поддерживающих эволюцию схем (например, добавление новых полей в колонку JSONB), обеспечивает обратную совместимость. Избегать разрушительных изменений в производственных таблицах; вместо этого добавлять новые столбцы или создавать новые таблицы и асинхронно мигрировать данные.

Оптимизируйте по стоимости

Данные временных рядов могут быть дорогими для хранения при высокой детализации. Внедряйте политику хранения для автоматического удаления данных старше определенного порога. Используйте выборку по времени: храните необработанные данные в течение 7 дней, затем в среднем по одной минуте в течение 30 дней, затем в почасовых средних в течение 1 года. Рассмотрим холодное хранение (например, Amazon S3 Glacier) для архивных данных, которые редко запрашиваются.

Тест с реальными объемами данных

Моделирование ожидаемой скорости передачи данных в среде постановки перед выходом на производство. Измерение распределения задержки (p50, p99, p999) для записи и чтения. Убедитесь, что модель данных может обрабатывать пиковые нагрузки (например, во время запуска машины, когда многие датчики отправляют данные одновременно).

Будущие тенденции в моделировании инженерных данных в реальном времени

Область быстро развивается. Новые тенденции включают использование ускоренных баз данных GPU для аналитики в реальном времени на больших наборах данных и внедрение периферийных вычислений, где модели данных должны работать на устройствах с ограниченными ресурсами. Другая тенденция - интеграция моделей ML непосредственно в конвейер данных, требуя моделей данных, которые могут служить векторами функций и прогнозами наряду с данными необработанных датчиков. Наблюдение и отслеживание линий данных также становятся необходимыми, поскольку инженерным командам необходимо проследить происхождение решения обратно к необработанным данным, которые его проинформировали.

Инженеры должны быть в курсе достижений в потоковой передаче SQL (например, Materialize, RisingWave), которые позволяют анализировать в реальном времени со стандартным SQL, уменьшая потребность в пользовательском коде обработки потоков. Эти инструменты обеспечивают реализацию декларативной модели данных, которая автоматически управляет состоянием и индексами.

Заключение

Проектирование моделей данных для обработки инженерных данных в режиме реального времени является сложной, но полезной задачей. Придерживаясь принципов масштабируемости, низкой задержки, гибкости и согласованности, а также выбирая правильные структуры данных и технологии обработки потоков, инженеры могут создавать системы, которые обеспечивают своевременную информацию и поддерживают непрерывность работы. Ключ заключается в понимании конкретных шаблонов запросов и требований к задержке вашего приложения, прототипа с реальными данными и итерацией на модели по мере развития инженерного ландшафта. Хорошо спроектированная модель данных является основой, на которой построены надежные высокопроизводительные инженерные системы в режиме реального времени.