Введение

Постановка проблемы

С чего все началось: Google нужно было индексировать данные для поиска по веб-страницам.

Что такое большие данные?

  • такой объем данных, который не помещается в ОЗУ одного компьютера
  • такой объем данных, который собирается быстрее чем один компьютер успевает их обрабатывать

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

Более формально характеристики больших данных описывают моделью 5V:

  • Volume (объем) — данные измеряются терабайтами и петабайтами
  • Velocity (скорость) — данные поступают и должны обрабатываться с высокой скоростью
  • Variety (многообразие) — данные разнородны по структуре: структурированные, полуструктурированные, неструктурированные
  • Veracity (достоверность) — данные могут быть неполными, противоречивыми, содержать шум
  • Value (ценность) — итоговая ценность извлекаемой из данных информации для бизнеса

Конечный результат обработки практически не зависит от одного конкретного элемента данных.

По способу обработки данные принято делить на:

  • Batch processing (пакетная обработка) — данные накапливаются, а затем обрабатываются большими порциями с определенной периодичностью (например, раз в сутки); характерна высокая пропускная способность при относительно большой задержке (latency)
  • Stream processing (потоковая обработка) — данные обрабатываются практически сразу по мере поступления, событие за событием или небольшими микропакетами; задержка минимальна, но пропускная способность обычно ниже, чем у batch-обработки

Также стоит различать два класса систем по характеру нагрузки:

  • OLTP (Online Transaction Processing) — системы, оптимизированные под большое количество коротких транзакций (вставка/обновление отдельных записей), типичный пример — реляционные СУБД в бэкенде приложения
  • OLAP (Online Analytical Processing) — системы, оптимизированные под сложные аналитические запросы к большим объемам исторических данных (агрегации, срезы), типичный пример — хранилища данных (data warehouse)

Дополнительную сложность при работе большими данными представляет их хранение данных перед обработкой. Для этого можно использовать:

  • Файлы — неудобно хранить метаданные, возможны проблемы с надежностью хранения, возможны проблемы со скоростью доступа
  • Реляционные СУБД (принципы ACID) — записывают данные медленнее чем нереляционные СУБД, масштабируются хуже чем нереляционные СУБД
  • Нереляционные СУБД (принципы BASE) — возможны проблемы с согласованностью данных

Компромисс между согласованностью, доступностью и устойчивостью к разделению сети распределенных систем описывается теоремой CAP: при возникновении сетевого разделения (partition) система вынуждена выбирать между согласованностью (Consistency) и доступностью (Availability). Именно поэтому многие нереляционные СУБД придерживаются принципов BASE (Basically Available, Soft state, Eventual consistency) вместо строгих ACID-гарантий.

Apache Hadoop

Hadoop — первая платформа для работы с большими данными.

Кластеры Hadoop строятся на недорогом аппаратном обеспечении (Commodity Hardware). Кластер Hadoop готов к выходу из строя небольшого количества вычислительных узлов.

Hadoop не создавался в Google — компания лишь опубликовала научные статьи, вдохновившие его разработку: в 2003 году вышла статья об архитектуре распределенной файловой системы Google (GFS), в 2004 году — статья о модели вычислений MapReduce. Опираясь на эти идеи, Даг Каттинг (Doug Cutting) и Майк Кафарелла (Mike Cafarella) реализовали аналогичную функциональность в рамках открытого проекта поискового робота Nutch, а в феврале 2006 года выделили ее в отдельный подпроект под названием Hadoop (назван в честь игрушечного слона сына Каттинга). Развитием проекта занялась компания Yahoo, куда в 2006 году перешел работать Каттинг. В январе 2008 года Hadoop получил статус самостоятельного top-level проекта фонда Apache.

Для обработки данных в Hadoop применяется MapReduce — модель распределенных вычислений при которой:

  • данные разделяются на большое количество одинаковых элементарных фрагментов
  • элементарные фрагменты обрабатываются параллельно функциями типа Map на узлах кластера Hadoop
  • результаты обработки элементарных фрагментов сводятся (редуцируются) в конечный результат функциями типа Reduce

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

Программисту остается написать на языке Java функции типа Map и Reduce и скомпилировать их в готовые к запуску модули в формате Jar.

Для запуска обработки данных пользователь должен создать задание в котором указывается исполняемый модуль в формате Jar и входные данные.

Apache Hadoop YARN

YARN (Yet Another Resource Negotiator) — компонент Hadoop, который управляет очередностью запуска заданий и распределением заданий узлам кластера.

Hadoop Distributed File System (HDFS)

Для хранения данных Hadoop использует файловую систему HDFS — Hadoop Distributed File System.

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

Благодаря репликации обеспечивается устойчивость распределенной системы к отказам отдельных узлов. Файлы в HDFS могут быть записаны лишь однажды (модификация не поддерживается). Запись в файл в одно время может вести только один процесс.

Организация файлов в пространстве имён — традиционная иерархическая: есть корневой каталог, поддерживается вложение каталогов, в одном каталоге могут располагаться и файлы, и другие каталоги.

HDFS предусматривает наличие центрального узла имён (name node), хранящего метаданные файловой системы и метаинформацию о распределении блоков. Узлы кластера которые непосредственно хранят блоки файлов называются узлами данных (data node).

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

Узел имён и узлы данных снабжаются веб-серверами, отображающими текущий статус узлов и позволяющими просматривать содержимое файловой системы. Административные функции доступны из интерфейса командной строки.

Несмотря на то что HDFS изначально являлся неотъемлемой частью проекта, сегодня Hadoop поддерживает работу и с другими распределёнными файловыми системами. Например, в основном дистрибутиве реализована поддержка Amazon S3 и CloudStore.

В облачных развертываниях HDFS все чаще заменяется (или дополняется) объектными хранилищами — Amazon S3, Google Cloud Storage, Azure Blob Storage. В отличие от HDFS, объектное хранилище не требует поддержки собственного постоянно работающего кластера, разделяет вычисления и хранение (compute/storage separation) и оплачивается по факту использования. Это делает его типовым слоем хранения современного data lake, поверх которого запускаются Spark, Hive или BigQuery.

С другой стороны, HDFS может использоваться не только для запуска MapReduce-заданий, но и как распределённая файловая система общего назначения. В частности, поверх неё реализована распределённая NoSQL-СУБД HBase (подробнее — в разделе про экосистему Hadoop ниже).

Развитие Hadoop

В процессе развития платформы Hadoop было создано большое количество высокоуровневых средств образующих своего рода экосистему:

  • Apache Hive — средство выполнения SQL-подобных запросов к данным
  • Cloudera Impala — средство выполнения SQL-подобных запросов к данным
  • Apache Pig — высокоуровневый скриптовый язык для обработки данных
  • Apache Mahout — Machine Learning поверх MapReduce
  • Apache Sqoop — импорт-экспорт данных из / в базы данных
  • Apache Kafka — распределенная платформа обмена сообщениями (message broker); чаще всего выступает источником данных для систем потоковой обработки, а не самой обрабатывающей системой
  • Apache Flume — средство для сбора и обработки логов
  • Apache Zookeeper — средство координации работы узлов
  • Apache Oozie — средство организации рабочего процесса и управления очередью заданий
  • Apache Ambari — веб-интерфейс для управления и мониторинга
  • Cloudera Hue — удобный веб-интерфейс для работы с кластером Hadoop
  • Apache HBase — средство доступа к данным в виде widecolumn NoSQL БД
  • и другие

Принципиально новым этапом в развитии Hadoop стал Apache Spark — платформа вычислений, которая ускоряет вычисления MapReduce благодаря тому, что хранит данные промежуточных результатов вычислений в ОЗУ узлов, а не в HDFS.

Платформа Spark также может работать в контейнерах под управлением Kubernetes и на кластерах Apache Mesos. Кроме того, Spark может управлять кластером самостоятельно.

На базе Spark работают несколько высокоуровневых программных средств:

  • Spark SQL — средство выполнения SQL-подобных запросов к данным
  • Spark Streaming — средство потоковой обработки данных
  • MLlib — средство для машинного обучения
  • GraphX — средство для вычислений на графах

Программирование на Spark осуществляется на:

  • Java
  • Scala
  • Python

Архитектурные паттерны обработки данных

По мере роста экосистемы сложились типовые архитектурные подходы к построению систем обработки больших данных:

  • Lambda-архитектура — совмещает два параллельных конвейера обработки одних и тех же данных: batch-слой (точный, но с задержкой — например, на Hadoop/Spark) и speed-слой (быстрый, но приближенный — например, на Spark Streaming/Kafka). Результаты обоих слоев объединяются на уровне отображения. Недостаток — необходимость поддерживать и синхронизировать два разных конвейера.
  • Kappa-архитектура — упрощение Lambda: вся обработка (в том числе исторических данных) ведется через единый потоковый конвейер, а пересчет прошлых данных выполняется повторным проигрыванием потока (replay) из брокера сообщений (например, Kafka).

Также различают модели хранения данных:

  • Data Warehouse (хранилище данных) — данные приводятся к заранее спроектированной схеме перед загрузкой (schema-on-write), оптимизировано под аналитические SQL-запросы (пример: Google BigQuery, AWS Redshift)
  • Data Lake (озеро данных) — данные любых форматов и структуры складываются “как есть” в исходном виде, схема применяется уже на этапе чтения (schema-on-read), что дает гибкость, но усложняет контроль качества данных
  • Lakehouse — гибридный подход, добавляющий поверх дешевого хранения data lake (в объектном хранилище) слой с ACID-транзакциями, версионированием схемы и другими возможностями, привычными для хранилищ данных (пример: Delta Lake поверх Spark)

Современные тренды

Современные тренды:

  • развертывание компонентов Hadoop в Kubernetes
  • хранение и обработка больших данные в облачных сервисах

Облачные решения:

ElasticStack

Компоненты:

  • ElasticSearch — распределенная поисковая система
  • Logstash — система сбора и обработки логов
  • Kibana — платформа аналитики и визуализации
  • Beats — коллекция легковесных средств экспорта данных

Часть программных компонентов имеет открытый исходный код.

На базе ElasticStack создано большое количество производных продуктов.

В настоящий момент активно развиваются облачные сервисы на базе продуктов ElasticStack.