Стив Уилкс, соучредитель и технический директор Striim, обсуждает необходимость архитектуры данных с «приоритетом потоковой передачи» и проводит демонстрацию платформы 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. Для систем, которые пытаются анализировать эти данные, это мало что значит. Но если вы объедините их с контекстом и скажете, что устройство XYZ — это датчик на конкретном двигателе конкретной машины, у вас появится больше контекста. Включив эти данные, вы сможете значительно улучшить потоковую обработку. Вы можете выполнять потоковую аналитику, коррелируя данные, объединяя их из нескольких различных потоков и ища совпадения. Например, используя веб-логи и сетевые логи, вы пытаетесь объединить их по IP-адресу и ищете события, произошедшие с обеих сторон за последние 30 секунд. Такая корреляция относится к сложной обработке событий (Complex Event Processing), которая ищет последовательности событий во времени, соответствующие определенному шаблону.
То есть, если происходит это, затем это, затем это, и это важно, вы можете провести статистический анализ аномалий и интегрироваться со сторонними системами машинного обучения. Мы также можем генерировать оповещения, запускать внешние системы и создавать очень информативные потоковые панели мониторинга для визуализации результатов вашей аналитики. Любые данные, собранные изначально, а также результаты обработки и аналитики могут быть доставлены куда угодно, причем в рамках одного приложения можно отправлять данные во множество различных целевых систем. Вы можете передавать их в корпоративные и облачные базы данных, файлы, 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. Но обычно на этом все не заканчивается.
Необработанные данные могут быть не очень полезны. Один из фрагментов данных здесь — это идентификатор продукта. И этого, вероятно, недостаточно. Поэтому сначала мы извлечем различные поля, включая идентификатор местоположения, идентификатор продукта, количество запасов и т. д. Это таблица мониторинга запасов, и мы только что превратили ее из необработанного формата в набор именованных полей. Это облегчит работу с ними в дальнейшем. Вы видите, что структура теперь сильно отличается от того, что мы видели в потоке данных. Если мы затем захотим добавить дополнительный контекст, мы сможем объединить эти данные с чем-то еще. Итак, сначала мы настроим это так, чтобы вместо записи необработанных данных в 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, потому что он разделен разделителями, и разделитель может быть любым, не обязательно запятой. Сохраним это. Теперь у нас есть компонент, который будет записывать данные в файл. Если мы развернем и запустим это, то начнем создавать файл с данными в режиме реального времени.









