Dev48
ЯЗЫК
  • О нас
  • Услуги
  • Индустрии
  • Технологии
  • Статьи
  • Контакты
Забронировать звонок
    Главная/Статьи/Kriticheskaya rol arhitektury dannyh s prioritetom potokovoy peredachi
Dev48

© 2026 · All rights reserved.

Критическая роль архитектуры данных с приоритетом потоковой передачи

Источник: Striim

Критическая роль архитектуры данных с приоритетом потоковой передачи

Источник: Striim

Узнайте о критической роли архитектуры данных с приоритетом потоковой передачи и посмотрите демонстрацию платформы Striim для потоковой интеграции и потоковой аналитики.

25 сентября 2026 г.

Стив Уилкс, соучредитель и технический директор Striim, обсуждает необходимость архитектуры данных с «приоритетом потоковой передачи» (streaming first) и проводит демонстрацию платформы корпоративного уровня Striim для потоковой интеграции и потоковой аналитики.

Чтобы узнать больше о платформе Striim, перейдите сюда.

Неотредактированная расшифровка:

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

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

Другая сторона медали заключается в том, что бизнесу нужно конкурировать. Технологии всегда были источником этой конкуренции. Поскольку технологии быстро меняются, а мы получаем все больше и больше данных, компании, которые все больше ориентируются на данные, получают конкурентное преимущество. Это касается почти всех отделов, от инженерии или производства до маркетинга. Поэтому выживание большинства компаний зависит от способности внедрять инновации и использовать новые технологии, а объем самих данных также массово увеличивается. Почти все может генерировать данные: грузовики, холодильники, телевизоры, носимые устройства, медицинские устройства, которые становятся все более портативными. Даже такие вещи, как отслеживание пути еды от фермы до ресторана — все это отслеживается и производит огромные объемы данных, и эти данные растут в геометрической прогрессии.

Компания IDC провела исследование пару месяцев назад, и они оценивают, что сегодняшние 16 зеттабайт данных к 2025 году увеличатся в 10 раз. Около 5% этих данных сейчас являются данными в реальном времени, и к 2025 году этот показатель вырастет до 25%. Под «реальным временем» они подразумевают данные, которые производятся и должны обрабатываться и анализироваться в режиме реального времени. Это 40 зеттабайт (21 ноль) данных, которые нужно будет обрабатывать в реальном времени. И 95% из них будут генерироваться устройствами. Проблема в том, что лишь малая часть этих данных может быть физически сохранена — производится недостаточно жестких дисков, чтобы хранить все это. Итак, если вы не можете их хранить, что вы можете с ними сделать? Единственный логический вывод заключается в том, что вам нужно обрабатывать и анализировать эти данные в оперативной памяти, потоковым способом, близко к месту их генерации. Возможно, вы превращаете необработанные данные — тысячи точек данных в секунду — в агрегированные данные, которые поступают реже, но все еще содержат тот же информационный контент. Вот о чем люди говорят как об обработке на периферии (edge processing), которая действительно предназначена для обработки этих огромных объемов данных, которые, как видят люди, будут поступать в будущем.

И дело не только в данных IoT, которые вызывают рост потоков. Каждая единица данных генерируется, потому что что-то произошло, какой-то вид события: кто-то работал в корпоративном приложении, кто-то делал что-то на веб-сайте или использовал веб-приложение, а машины генерировали логи на основе того, что они делали. Приложения генерируют базы данных, базы данных генерируют логи, сетевые устройства — все генерирует логи. Но все они основаны на том, что происходит, на событиях. Поэтому, если данные создаются на основе событий в потоковом режиме, то они должны обрабатываться и анализироваться в потоковом режиме. Если вы собираете данные пакетами, то вы никогда не придете к архитектуре реального времени и пониманию того, что происходит в реальном времени. Но если вы собираете данные как потоки, то вы можете делать и другие вещи: вы можете выполнять пакетную обработку потоковых данных, вы можете доставлять их куда-то еще, но, по крайней мере, данные должны быть потоковыми.

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

Некоторые вещи, о которых нам говорят наши клиенты: у нас есть устаревшие системы, и под «устаревшими» можно понимать все, что было установлено более года назад, и они не справляются или, по их прогнозам, не смогут справиться с большими объемами данных, которые они ожидают увидеть, с требованиями к низкой задержке, своего рода пониманием данных в реальном времени и способностью быстро внедрять инновации и быстро создавать новые типы анализа, новые типы обработки, дающие новые типы понимания того, что происходит в вашем бизнесе. Они также говорят нам, что мы не можем просто взять и заменить эти системы, и им нужно, чтобы новые и старые системы работали вместе, возможно, с переключением при отказе с одной на другую, пока вы выполняете эту замену. Striim существует уже около пяти лет.

Мы являемся поставщиками платформы Striim, которая обеспечивает потоковую интеграцию и аналитику. Это зрелая платформа. Она находится в промышленной эксплуатации у клиентов уже более трех лет. Среди наших клиентов компании из самых разных отраслей: финансовые услуги, телекоммуникации, здравоохранение, розничная торговля. Мы наблюдаем большую активность в сфере IoT. Striim — это комплексная сквозная платформа, которая выполняет потоковую интеграцию и аналитику в масштабах предприятия, облака и IoT. У нас очень гибкая архитектура, позволяющая развертывать потоки данных для объединения корпоративных систем, облака и IoT. Вы можете развертывать части приложения на периферии, рядом с местом генерации данных. И это не обязательно должны быть только данные IoT. Это могут быть любые данные, генерируемые вблизи источника. Другие части могут работать локально, выполняя некоторую обработку, а остальные — в облаке.

Мы также занимаемся обработкой и аналитикой. Наша платформа позволяет очень гибко развертывать приложения, поскольку они состоят из непрерывного сбора данных в режиме реального времени. Это могут быть данные от устройств, которые вы считаете «реально-временными», например, датчиков, отправляющих события, очередей сообщений и т. д. Или это могут быть файлы, которые обычно ассоциируются с пакетной обработкой. Мы можем считывать данные из конца файла, и по мере записи новых записей в файл, поток немедленно превращает файлы, их ротацию и т. д. в источник потоковых данных. Что касается баз данных, большинство людей воспринимают их как историческую запись того, что произошло в прошлом. Но с помощью технологии под названием Change Data Capture (отслеживание измененных данных) вы можете видеть вставки, обновления, удаления — все, что происходит в этой базе данных в режиме реального времени. Таким образом, вы можете собирать эти данные из базы данных без вмешательства в поток.

Вот и все. Теперь у вас есть поток всех изменений, происходящих в базе данных. Итак, все приложения, созданные на этой платформе, используют ту или иную форму непрерывного сбора данных. Поверх этого вы можете выполнять потоковую обработку в реальном времени с помощью SQL-запросов. Здесь не требуется программирование на Java, C# или JavaScript. Вы можете построить все, используя SQL, что позволяет выполнять фильтрацию, преобразование и агрегирование данных. Да, используя окна данных. Вы можете указать, что произошло за последнюю минуту, отслеживать изменения в данных и отправлять только их и т. д. Также очень важно обогащение данных. Это способность загружать большие объемы справочных данных в память распределенного кластера и объединять их в режиме реального времени с потоковыми данными для добавления дополнительного контекста.

Примером может служить ситуация, когда поступают данные об устройстве: устройство X, Y, Z, значение 1, 2, 3. Это мало что говорит системам, которые пытаются их анализировать. Но если вы объедините это с контекстом и скажете, что устройство X, Y, Z — это датчик на конкретном двигателе конкретной машины, теперь у вас больше контекста. И если вы включите эти данные, вы сможете добиться гораздо лучших результатов при потоковой обработке. Вы можете выполнять потоковую аналитику, которая коррелирует данные, объединяя их из нескольких различных потоков и ища совпадения. Например, используя веб-логи и сетевые логи, вы пытаетесь объединить их по IP-адресу. Вы ищете события, произошедшие с любой стороны за последние 30 секунд. Такая корреляция для обработки сложных событий позволяет искать последовательности событий во времени, которые соответствуют определенному шаблону.

То есть, если происходит это, за ним это, а затем это, и это важно, вы можете провести статистический анализ аномалий и интегрироваться со сторонним машинным обучением. Да, мы также можем генерировать оповещения, запускать внешние системы и создавать очень информативные потоковые панели мониторинга для визуализации результатов вашей аналитики. И любые данные, которые были изначально собраны, результаты обработки, результаты аналитики — все это может быть доставлено куда угодно, и вы можете доставлять данные во множество различных целей в рамках одного приложения. Вы можете отправлять данные в корпоративные и облачные базы данных, файлы, Hadoop, Kafka и т. д. Как новый тип промежуточного ПО, поддерживающий потоковую интеграцию и аналитику, для нас очень важно интегрироваться с вашим существующим программным обеспечением. Поэтому у нас есть множество коллекторов данных и механизмов доставки, которые работают с системами, которые у вас уже могут быть. Это касается систем больших данных, корпоративных баз данных, open-source решений — мы можем интегрироваться со всем этим и делать все это на корпоративном уровне, который по своей сути является кластеризованным, распределенным, масштабируемым, надежным и безопасным как универсальное промежуточное ПО.

Мы поддерживаем множество различных вариантов использования: от интеграции данных в реальном времени и аналитики до создания панелей мониторинга и мониторинга систем. Эти варианты использования охватывают все отрасли и варьируются от построения озера данных и подготовки данных перед их загрузкой до миграции данных и обработки на периферии IoT. Что касается аналитики и шаблонов, то это обнаружение мошенничества, прогнозное техническое обслуживание, борьба с отмыванием денег — вот некоторые из задач, с которыми к нам обращаются клиенты. А если вы хотите создавать панели мониторинга и отслеживать показатели в реальном времени, например, соблюдение SLA или ожиданий, мы занимались такими вещами, как мониторинг качества работы колл-центров, мониторинг SLA. Мы смотрим на работу с точки зрения клиента и можем оповещать, когда что-то работает нештатно. Варианты использования охватывают множество различных отраслей. Здесь много текста, но главный вывод в том, что у нас есть сценарии использования в самых разных отраслях.

Один из примеров — использование Striim для интеграции гибридного облака. Это случай, когда у вас есть локальная база данных, и вы хотите переместить или скопировать ее в облако. Одно дело просто взять базу данных и перенести ее в облако, но при этом можно упустить все, что происходит во время переноса или после него. Поэтому очень важно включить в этот процесс Change Data Capture, чтобы постоянно пополнять вашу гибридную облачную базу данных новой информацией. Используя набор мастеров, вы можете очень быстро создать решение, которое позволит объединить, например, локальную базу данных Oracle и доставлять данные из нее в режиме реального времени, скажем, в Azure SQL DB. Таким образом, у вас будет точная копия локальной базы данных, которая всегда актуальна.

Еще один совершенно другой пример — это использование нас для мониторинга безопасности, когда у вас создается множество различных журналов VPN-шлюзами, брандмауэрами, сетевыми маршрутизаторами, отдельными машинами, по сути, микроконтроллерами. Все, что может создавать журнал, и вы распознаете необычное поведение, чаще всего проявляется при воздействии на несколько систем. Специалисты по безопасности получают огромное количество оповещений из всех этих журналов и систем постоянно. Но многие из них являются ложноположительными. Поэтому цель состояла в том, чтобы определить, что для них является наиболее приоритетным для первоочередного рассмотрения, видя, какая активность происходит и затрагивает несколько объектов. Например, если у вас есть сканирование портов с сетевого маршрутизатора, этот парень просматривает другие вещи. Есть ли какая-то активность на других машинах, которые он просматривал? Хорошо. Выполняют ли они сканирование портов?

Подключаются ли они к внешним сайтам и загружают вредоносное ПО? Выполняя эту корреляцию в памяти в режиме реального времени, вы можете выявлять угрозы с более высоким приоритетом. Кроме того, предварительно сопоставляя все данные и предоставляя их аналитикам, они могут сразу увидеть необходимую информацию, вместо того чтобы вручную искать ее в куче разных журналов. Это действительно повышает продуктивность аналитиков. Вот еще пара примеров от наших клиентов. Один из них — очень простое перемещение данных в реальном времени, когда данные из баз данных HP NonStop и SQL Server передаются в несколько целевых систем, будь то Hadoop, HDFS, Kafka или HBase, и они используют это как аналитический центр для своих сообществ. По сути, это гарантирует, что где бы они ни хотели разместить данные, они всегда будут актуальными и содержать информацию в реальном времени из других баз данных.

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

Сквозная платформа Striim делает все: от сбора и обработки до доставки аналитики и визуализации потоковых данных. Она проста в использовании, поддерживает язык SQL для построения обработки и анализа, что позволяет создавать и развертывать приложения за считанные дни. Мы соответствуем корпоративному уровню, что означает, что мы по своей сути масштабируемы в распределенной архитектуре, надежны и безопасны. И нас легко интегрировать с выбранными вами существующими технологиями. Это ключевые моменты, которые стоит помнить о том, почему мы отличаемся. Итак, переходим к демонстрации. Первая часть, по сути, покажет вам, как выполнить интеграцию, вместо того чтобы печатать много кода. Мы просто пройдемся по тому, как создать захват измененных данных (CDC) в Kafka, выполнить некоторую обработку и затем доставить данные в другие места.

Это чисто интеграционное решение. Вы начинаете с захвата измененных данных из SQL, в данном случае MySQL, создаете первоначальное приложение, а затем настраиваете получение данных из источника. Мы настраиваем информацию для подключения к MySQL. Когда вы это делаете, мы проверяем, все ли будет работать, и настроен ли у вас захват измененных данных должным образом. Если нет, мы подскажем, как это исправить. Вы выбираете таблицы, из которых хотите собирать измененные данные, и это создаст поток данных. Затем этот поток данных пойдет в Kafka. Мы настроим, как хотим записывать данные в Kafka, то есть настроим конфигурацию брокера, тему и формат данных.

В данном случае мы записываем данные в формате JSON. Когда мы сохраняем это, создается поток данных, который очень прост. В данном случае он состоит из двух компонентов. Мы идем из источника MySQL CDC в средство записи Kafka. Мы можем протестировать это, развернув приложение, это двухэтапный процесс. Сначала вы развертываете компоненты в кластере, а затем запускаете его, и теперь мы можем видеть данные, которые текут между ними. Если я нажму на это, я смогу увидеть данные в реальном времени. Вы видите данные и то, что было до этого. Это, по сути, четыре обновления. Вы также получаете «образ до», так что можете видеть, что именно изменилось. Это данные в реальном времени, проходящие через приложение MySQL. Но обычно на этом все не заканчивается.

Необработанные данные могут быть не очень полезны. Один из фрагментов данных здесь — это идентификатор продукта (product id). И, вероятно, он не содержит достаточно информации. Поэтому первое, что мы сделаем, — это извлечем различные поля из этого, включая идентификатор местоположения, идентификатор продукта, количество запасов и т. д. Это таблица мониторинга запасов, и мы только что превратили ее из необработанного формата в набор именованных полей. Это облегчит работу с ними в дальнейшем. Вы видите, что структура теперь совсем другая. То, что мы на самом деле видим в этом потоке данных. Если мы затем захотим добавить дополнительный контекст, мы сможем объединить эти данные с чем-то еще. Итак, сначала мы просто настроим это так, чтобы вместо записи необработанных данных в Kafka мы записывали обработанные данные. Все, что нам нужно сделать, — это изменить входной поток. Это изменит поток данных. Теперь обработанные данные будут записываться в Kafka.

Теперь мы добавим кэш — это распределенная сетка данных в оперативной памяти, которая будет содержать дополнительную информацию, которую мы хотим объединить с необработанными данными. Это информация о продуктах. Каждый идентификатор продукта (product ID) имеет описание, цену и другие параметры. Сначала мы создадим тип данных, соответствующий нашей таблице базы данных, и настроим ключ. В данном случае ключом является product ID. Затем мы указываем, как будем получать данные: это могут быть файлы или acfs. Мы воспользуемся средством чтения из базы данных (database reader), чтобы загрузить их из таблицы MySQL. Укажем все параметры подключения и запрос, который будем использовать. Теперь у нас есть кэш с информацией о продуктах. Чтобы использовать его, мы изменим SQL-запрос для объединения с кэшем.

Любой, кто когда-либо писал SQL-запросы, знает, как выглядит оператор JOIN. Мы просто объединяем данные по product ID. Теперь, вместо «сырых» данных, мы получаем дополнительные поля, которые подтягиваются в режиме реального времени из информации о продуктах. Если мы запустим это и снова посмотрим на данные, то увидим дополнительные поля, такие как описание, бренд, категория и цена, которые были получены из другого типа данных и объединены в оперативной памяти. Никаких обращений к базе данных не происходит, поэтому это работает очень быстро. Теперь перейдем к Kafka. Если у вас уже есть данные в Kafka, другой шине сообщений или где-либо еще, например, в виде файлов, вы можете захотеть прочитать их и отправить в некоторые целевые системы. Сейчас мы возьмем данные, которые только что записали в Kafka.

Мы воспользуемся Kafka reader. Найдем его и добавим в качестве источника. Затем настроим свойства для подключения к брокеру, который мы только что использовали. Поскольку мы знаем, что это данные в формате JSON, мы будем использовать JSON-парсер. Он разобьет их на структуру объектов JSON и создаст поток данных. Когда мы развернем и запустим это приложение, оно начнет чтение из топика Kafka. Мы можем просмотреть эти данные и увидеть, что это те самые данные, которые мы записывали ранее, со всей информацией в формате JSON. Вы можете видеть структуру JSON. Однако для других целевых систем, куда мы собираемся записывать данные, структура JSON может не подойти. Что же нам теперь делать?

Мы добавим запрос, который извлечет различные поля из этой структуры JSON и создаст четко определенный поток данных с отдельными полями. Мы напишем запрос для этого, который будет напрямую обращаться к данным JSON, и сохраним его. Теперь, вместо исходного потока данных с JSON внутри, когда мы развернем и запустим его, мы увидим данные. Кстати, именно так вы и будете создавать приложения, постоянно просматривая данные в процессе разработки и добавления новых компонентов. Если мы посмотрим на поток данных сейчас, то увидим отдельные поля, которые у нас были до отправки в Kafka. Не забывайте, что это не обязательно должна быть потоковая передача в Kafka — это может быть что угодно. Если вы делаете что-то вроде того, что мы только что проделали с CDC в Kafka, а затем из Kafka в другие целевые системы, вам не обязательно использовать Kafka посередине — можно просто взять данные CDC и отправить их напрямую в целевые системы.

Теперь мы добавим простую целевую систему, которая будет записывать данные в файл. Для этого мы выберем file writer и укажем нужный формат. Мы запишем это в формате CSV. На самом деле мы называем его DSV, так как он разделен разделителями. Разделителем может быть что угодно, не обязательно запятая. Сохраним это. Теперь у нас есть компонент, который будет записывать данные в файл. Если мы развернем и запустим его, то создадим файл с данными в режиме реального времени.

← Все статьи