Что такое Apache Kafka: как устроен и работает брокер сообщений
Apache Kafka — распределенный брокер сообщений, работающий в стриминговом режиме. В статье мы расскажем про его устройство и преимущества, а также о том, где применяют это ПО.

Apache Kafka — распределенный брокер сообщений, работающий в стриминговом режиме. В статье мы расскажем про его устройство и преимущества, а также о том, где применяют «Кафку».
Что такое брокер сообщений
Главная задача брокера — обеспечение связи и обмена информацией между приложениями или отдельными модулями в режиме реального времени.
Брокер — система, преобразующая сообщение от источника данных (продюсера) в сообщение принимающей стороны (консьюмера). Брокер выступает проводником и состоит из серверов, объединенных в кластеры.
Apache Kafka — диспетчер сообщений, разработанный LinkedIn. В 2011 году был опубликован программный код. В 2012 году Kafka попал в инкубатор Apache, дальнейшая разработка ведется в рамках Apache Software Foundation. Открытое программное обеспечение с разрешительной лицензией написано на Java и Scala.
Изначально «Кафку» создавали как систему, оптимизированную под запись, и создатель Джей Крепс выбрал такое название в честь одного из своих любимых писателей.
Шаги передачи данных
Чтобы понять, как функционирует распределенная система Apache Kafka, необходимо проследить путь данных.
Событие или сообщение — данные, которые поступают из одного сервиса, хранятся на узлах Kafka и читаются другими сервисами. Сообщение состоит из:
- Key — опциональный ключ, нужен для распределения сообщений по кластеру.
- Value — массив байт, бизнес-данные.
- Timestamp — текущее системное время, устанавливается отправителем или кластером во время обработки.
- Headers — пользовательские атрибуты key-value, которые прикрепляют к сообщению.
Продюсер — поставщик данных, который генерирует сообщения — например, служебные события, логи, метрики, события мониторинга.
Консьюмер — потребитель данных, который читает и использует события, пример — сервис сбора статистики.

Какие сложности решает распределенная система
Сообщения могут быть однотипными или разнородными, поскольку разным потребителям нужны разные данные. Один тип событий может быть нужен всем консьюмерам, а другие — только одному.
Без брокера продюсеры должны знать получателя и резервного консьюмера, если основной недоступен. К тому же, поставщикам данных придется самостоятельно регистрировать новых консьюмеров. С помощью брокера продюсеры просто отправляют информацию в единый узел.
Managed service для Apache Kafka
Сообщения хранятся на узлах-брокерах. Kafka — масштабируемый кластер со множеством взаимозаменяемых серверов, в которые добавляются новые брокеры, распределяющие задачи между собой.
ZooKeeper — инструмент-координатор, действует как общая служба конфигурации в системе. Работает как база для хранения метаданных о состоянии узлов кластера и расположении сообщений. ZooKeeper обеспечивает гибкую и надежную синхронизацию в распределенной системе, позволяя нескольким клиентам выполнять одновременно чтение и запись.
Kafka Controller — среди брокеров Zookeeper выбирает одного, который будет обеспечивать консистентность данных.
Topic — принцип деления потока данных, базовая и основная сущность Apache Kafka. В топик складывается стрим данных, единая очередь из входящих сообщений.
Partition — для ускорения чтения и записи топики делятся на партиции. Происходит параллелизация данных. Это конфигурируемый параметр, сообщения могут отправлять несколько продюсеров и принимать несколько консьюмеров.

Упорядочение событий происходит на уровне партиций. Принимающая сторона потребляет данные в порядке расположения в партиции. Пример: все события одного пользователя сервисы принимают упорядоченно, обработка сохраняет последовательность пути пользователя. Выстраивается конвейер данных, алгоритмы машинного обучения могут извлекать из сырой информации необходимую для бизнеса информацию.
Преимущества Apache Kafka
Брокер распределяет информацию в широковещательном режиме. Применяющийся в Apache Kafka подход нужен для масштабирования и репликации данных.
Горизонтальное масштабирование
Множество объединенных серверов гарантируют высокую доступность данных — выход из строя одного из узлов не нарушает целостность. Кластер состоит из обычных машин, а не суперкомпьютеров, их можно менять и дополнять. Система автоматически перебалансируется.
Чтобы события не потерялись, существуют механизмы репликации. Данные записываются на несколько машин, если что-то случается с сервером, он переключается на резервный. Кластер в режиме реального времени определяет, где находятся данные, и продолжает их использовать.
Офсеты
Если консьюмер падает в процессе получения данных, то, когда он запустится вновь и ему нужно будет вернутся к чтению этого сообщения, он воспользуется офсетом и продолжит с нужного места.
Взаимодействие через API
Брокеры решают проблему интеграции разных технических стеков и протоколов. Интеграция происходит просто: продюсерам и консьюмерам необходимо знать только API брокера. Они не контактируют между собой, с помощью чего достигается высокая интегрируемость с другими системами.
Принцип first in — first out
Принцип FIFO действует на консьюмеров. Чтение происходит в том же порядке, в котором пришла информация.
Где применяется Apache Kafka
Отказоустойчивая система используется в бизнесе, где необходимо собирать, хранить и обрабатывать большие неструктурированные данные. Примеры — платформы, где требуется интеграция данных из большого количества источников, сервисы стриминговой аналитики, mission-critical applications.
Big Data
Первоначально LinkedIn разработали «Кафку» для своих целей: обмена данными между службами, репликации баз данных, потоковой передачи информации о деятельности и операционных показателях приложений.
Для IBM Apache Kafka работает как средство обмена сообщениями между микросервисами. В аналитических системах американской корпорации Apache Kafka обрабатывает потоковые и событийные данные.
Uber, Twitter, Netflix и AirBnb с помощью хорошо развитых пайплайнов обработки данных передают миллиарды сообщений в день. «Кафка» решает проблемы перемещения Big data из одного источника в другой.
Издание The New York Times использует Apache Kafka для хранения и распространения опубликованного контента среди различных приложений и систем, которые делают его доступным для читателей в режиме реального времени.
Internet of Things
IoT-платформы используют архитектуру с большим количеством конечных устройств: контроллеров, датчиков, сенсоров и smart-гаджетов. ПО интернета вещей с помощью алгоритмов ML составляет графики профилактического ремонта оборудования, анализируя данные, поступающие с устройств.
ML-системы работают с онлайн-потоками, когда приборы, приложения и пользователи постоянно посылают данные, а сервисы обрабатывают их в реальном времени. Apache Kafka выступает центральным звеном в этом процессе.
Отрасли
Kafka используют организации практически в любой отрасли: разработка ПО, финансовые услуги, здравоохранение, государственное управление, транспорт, телеком, геймдев.
Сегодня Kafka пользуются тысячи компаний, более 60% входят в список Fortune 100. На официальном сайте представлен полный список корпораций и учреждений, которые используют брокера Apache.
Конкуренты
Чаще всего Kafka сравнивают с RabbitMQ. Обе системы — брокеры сообщений. Главное отличие в модели доставки: Kafka добавляет сообщение в журнал, и консьюмер сам забирает информацию из топика; брокер RabbitMQ самостоятельно отправляет сообщения получателям — помещает событие в очередь и отслеживает его статус.
«Кролик» удаляет событие после доставки, «Кафка» хранит до запланированной очистки журнала. Таким образом, брокер Apache используется как источник истории изменений.
Разработчики RabbitMQ создали системы управления потоком сообщений: мониторинг получения, маршрутизация и шаблоны доставки. Подобное гибкое управление подойдет для высокоскоростного обмена сообщениями между несколькими сервисами. Минус такого подхода в снижении производительности при высокой нагрузке.
Главный вывод — для сбора и агрегации событий из большого количества источников, логов и метрик больше подойдет Apache Kafka.
Заключение
Благодаря высокой пропускной способности и согласованности данных Apache Kafka обрабатывает огромные массивы данных в реальном времени. Системы горизонтального масштабирования и офсеты гарантируют надежность. Kafka — удачное решение для проекта с очень большими нагрузками на обработку данных. Установить это ПО можно на серверы Ubuntu, Windows, CentOS и других популярных операционных систем.
Часть 3. MPI — Как процессы общаются? Сообщения типа точка-точка
В этом цикле статей речь идет о параллельном программировании с использованием MPI.
- Часть 1. MPI — Введение и первая программа.
- Часть 2. MPI — Учимся следить за процессами.
- Часть 3. MPI — Как процессы общаются? Сообщения типа точка-точка.
В прошлой статье мы обсудили как распределяется работа между процессами, зачем нужно знать какой процесс выполняется на конкретном потоке и как фиксировать время выполнение работы программы. Что дальше?
Иногда нам требуется остановить один/все процессы чтобы они подождали какого-либо действия системы, пользователя, пересылать одни локальные данные с процесса в другой процесс, ведь, напомню, они работают с независимой памятью и изменение переменной в одном потоке не приведет к изменению переменной с тем же именем в другом процессе, более формально процессы работают в непересекающихся множествах адресов памяти, а изменение данных в одном из них никак не повлияет на другое множество(собственно из-за того, что они непересекающиеся).
Предисловие
Практически все программы с применением технологии MPI используют не только средства для порождения процессов и их завершения, но и одну из самых важных частей, заложенных в названии самой технологии (Message Passing Interface), конечно же явная посылка сообщений между процессами. Описание этих процедур начнем, пожалуй, с операций типа точка-точка.
Давайте представим себе ту самую картинку где много компьютеров соединены линиями, стандартная иллюстрация сети(локальной), только на месте узлов будут стоять отдельные процессы, процессоры, системы и т.п. Собственно отсюда становится понятно к чему клонит название «операции типа точка-точка», два процесса общаются друг с другом, пересылают какие-либо данные между собой.
Работает это так: один из процессов, назовем его P1, какого-то коммуникатора C должен указать явный номер процесса P2, который также должен быть под коммуникатором С, и с помощью одной из процедур передать ему данные D, но на самом деле не обязательно нужно знать номер процесса, но это мы обсудим далее.
Процедуры в свою очередь разделены на 2 типа: с блокировкой и без блокировки.
1. Процедуры обмена с блокировкой просто останавливают процесс который должен принять сообщение до момента выполнения какого-либо условия.
2. Процедуры без блокировки, как их ещё называют асинхронные, они возвращаются сразу после осуществления коммуникации, немедленно.
Оба типа процедур мы затронем подробно, а пока начнем с процедур с блокировкой.
Передаем и принимаем сообщения
Наконец приступим к практической части. Для передачи сообщений используется процедура MPI_Send. Эта процедура осуществляет передачу сообщения с блокировкой. Синтаксис у нее следующий:
int MPI_Send(void* buf, int count, MPI_Datatype datatype, int dest, int msgtag, MPI_Comm comm);
Что тут есть что:
buf — ссылка на адрес по которому лежат данные, которые мы пересылаем. В случае массивов ссылка на первый элемент.
count — количество элементов в этом массиве, если отправляем просто переменную, то пишем 1.
datatype — тут уже чутка посложнее, у MPI есть свои переопределенные типы данных которые существуют в С++. Их таблицу я приведу чуть дальше.
dest — номер процесса кому отправляем сообщения.
msgtag — ID сообщения (любое целое число)
comm — Коммуникатор в котором находится процесс которому мы отправляем сообщение.
А вот как называются основные стандартные типы данных С++ определенные в MPI_Datatype:
Название в MPI
Тип даных в С++
signed long int
signed long long int
MPI_UNSIGNED_*** (Вместо *** int и т.п.)
Аналогичным способом указываются и другие типы данных определенные в стандартной библиотеке С/С++ — MPI_[Через _ в верхнем регистре пишем тип так как он назван в С]. Еще один пример для закрепления понимания, есть тип беззнаковых 32 битных целых чисел, назван он uint32_t, чтобы получить этот тип данных переопределенным в MPI необходимо написать следующую конструкцию: MPI_UINT32_T. То есть все вполне логично и легко, верхний регистр, вместо пробелов знаки андерскора и в начале пишем MPI.
Блокировка защитит пересылаемые данные от изменений поэтому не стоит опасаться за корретность отправленных данных, после того как коммуникация завершится и выполнится какое либо условие, то процесс спокойно продолжит заниматься своими делами.
Теперь поговорим о приеме этих сообщений. Для этого в MPI определена процедура MPI_Recv. Она осуществляет, соответственно, блокирующий прием данных. Синтаксис выглядит вот так:
int MPI_Recv(void* buf, int count, MPI_Datatype datatype, int source, int tag, MPI_Comm comm, MPI_Status* status);
Что тут есть что:
buf — ссылка на адрес по которому будут сохранены передаваемые данные.
count — максимальное количество принимаемых элементов.
datatype — тип данных переопределенный в MPI(по аналогии с Send).
source — номер процесса который отправил сообщение.
tag — ID сообщения которое мы принимаем (любое целое число)
comm — Коммуникатор в котором находится процесс от которого получаем сообщение.
status — структура, определенная в MPI которая хранит информацию о пересылке и статус ее завершения.
Тут все практически идентично процедуре Send, только появился аргумент статуса пересылки. Зачем он нужен? Не всегда нужно явно указывать от какого процесса приходит сообщение, какой тег сообщения мы принимаем, чтобы избавиться от неопределенности MPI сохраняет информацию которая не указана в процессе-преемнике явно и мы можем к ней обратиться. Например чтобы узнать процесс который отправил сообщение и тэг этого сообщения:
MPI_Status status; MPI_Recv(&buffer, 1, MPI_Float, MPI_ANY_SOURCE, MPI_ANY_TAG, MPI_COMM_WORLD, &status) tag = status.MPI_SOURCE source = status.MPI_SOURCE
Здесь представлен кусочек возможного кода в процессе который принимает данные не более одного float элемента от любого процесса с любым тегом сообщения. Чтобы узнать какой процесс прислал это сообщение и с каким тэгом нужно собственно воспользоваться структурой MPI_Status.
Заметим появление констант MPI_ANY_SOURCE и MPI_ANY_TAG, они явно указывают, что можно принимать сообщения от любого процесса с любым тэгом.
Дабы закрепить знания об этих двух процедурах приведу пример как это выглядит на практике, данная программа определяет простые числа на заданном интервале:
#include #include #include "mpi.h" #define RETURN return 0 #define FIRST_THREAD 0 int* get_interval(int, int, int*); inline void print_simple_range(int, int); void wait(int); int main(int argc, char **argv) < // инициализируем необходимые переменные int thread, thread_size, processor_name_length; int* thread_range, interval; double cpu_time_start, cpu_time_fini; char* processor_name = new char[MPI_MAX_PROCESSOR_NAME * sizeof(char)]; MPI_Status status; interval = new int[2]; // Инициализируем работу MPI MPI_Init(&argc, &argv); // Получаем имя физического процессора MPI_Get_processor_name(processor_name, &processor_name_length); // Получаем номер конкретного процесса на котором запущена программа MPI_Comm_rank(MPI_COMM_WORLD, &thread); // Получаем количество запущенных процессов MPI_Comm_size(MPI_COMM_WORLD, &thread_size); // Если это первый процесс, то выполняем следующий участок кода if(thread == FIRST_THREAD) < // Выводим информацию о запуске printf("----- Programm information -----\n"); printf(">>> Processor: %s\n", processor_name); printf(">>> Num threads: %d\n", thread_size); printf(">>> Input the interval: "); // Просим пользователья ввести интервал на котором будут вычисления scanf("%d %d", &interval[0], &interval[1]); // Каждому процессу отправляем полученный интервал с тегом сообщения 0. for (int to_thread = 1; to_thread < thread_size; to_thread++) MPI_Send(&interval, 2, MPI_INT, to_thread, 0, MPI_COMM_WORLD); // Начинаем считать время выполнения cpu_time_start = MPI_Wtime(); >// Если процесс не первый, тогда ожидаем получения данных else MPI_Recv(&interval, 2, MPI_INT, MPI_ANY_SOURCE, MPI_ANY_TAG, MPI_COMM_WORLD, &status); // Все процессы запрашивают свой интервал range = get_interval(thread, thread_size, interval); // После чего отправляют полученный интервал в функцию которая производит вычисления print_simple_range(range[0], range[1]); // Последний процесс фиксирует время завершения, ожидает 1 секунду и выводит результат if(thread == thread_size - 1) < cpu_time_fini = MPI_Wtime(); wait(1); printf("CPU Time: %lf ms\n", (cpu_time_fini - cpu_time_start) * 1000); >MPI_Finalize(); RETURN; > int* get_interval(int proc, int size, int interval) < // Функция для рассчета интервала каждого процесса int* range = new int[2]; int interval_size = (interval[1] - interval[0]) / size; range[0] = interval[0] + interval_size * proc; range[1] = interval[0] + interval_size * (proc + 1); range[1] = range[1] == interval[1] - 1 ? interval[1] : range[1]; return range; >inline void print_simple_range(int ibeg, int iend) < // Прострейшая реализация определения простого числа bool res; for(int i = ibeg; i res = not res; if(res) printf("Simple value ---> %d\n", i); > > void wait(int seconds) < // Функция ожидающая в течение seconds секунд clock_t endwait; endwait = clock () + seconds * CLOCKS_PER_SEC ; while (clock() < endwait) <>; >
Поясню основные моменты и идею. Здесь передача сообщений задействована дабы остановить все процессы кроме первого до того момента пока пользователь не введет необходимые данные. Как только мы вводим требуемый интервал для вычисления, то процессы получают эту информацию от того потока, который занимался ее сбором и начинают свою работу независимо. В конце последний процесс фиксирует время вычислений, ожидает одну секунду(дабы остальные уж точно завершили за это время свою работу) и выводит результат расчета времени. Почему же именно последний?
Нужно учитывать тот факт, что интервал который приходит от пользователя разбивается на N равных частей, а значит последнему процессу достанется часть с самыми большими числами, вследствие вычисления займут наибольшее время именно в нем, а значит и программа в худшем случае отработает именно за это время.
Делаем прием данных более гибким
Не всегда мы имеем представление для конкретного процесса о том какой длины придут данные, для определения существует как раз выше упомянутая структура MPI_Status и некоторые процедуры помогающие эту информацию оттуда извлечь.
Первая процедура которую мы обсудим следующая:
int MPI_Get_count(MPI_Status* status, MPI_Datatype datatype, int* count);
По структуре status процедура определяет сколько данных типа datatype передано соответствующим сообщением и записывает результат по адресу count.
То есть буквально, если мы получаем сообщение от какого-либо процесса и не знаем сколько точно там передано данных, то можем вызвать процедуру MPI_Get_count и узнать какое количество ячеек памяти мы можем считать заполненными корректными данными(если конечно их отправляющий процесс корректно формирует).
Также есть еще одна процедура MPI_Get_elements. По синтаксису они отличаются лишь названиями, но назначение слегка разное. Если в сообщении передаются данные не базового типа, а типа который является производным от базовых(то есть составлен с помощью базовых типов), то нам вернет не количество этих данных, а именно количество данных базового типа. Однако в случае если передаются данные базового типа, то функции вернут одинаковые значения.
Также иногда случается так, что нам надо просто пропустить отправку сообщения, но саму процедуру из кода исключать не хочется, либо делать лишние условия в коде не рационально. В таких случаях процесс может отправить сообщение не существующему процессу, номер такого процесса определен константой MPI_PROC_NULL. В случае если мы передаем сообщение такому процессу, то процедура сразу завершается с кодом возврата SUCCESS.
Хорошо, мы можем принимать на вход какие либо данные и не знать сколько их точно поступает. В таком случае нужно рационально выделять какой-то объем памяти для их сохранения(буферизации). Возникает логичный вопрос о том какой объем выделять, в этом нам поможет процедура MPI_Probe. Она позволяет получить информацию о сообщении которое ожидает в очереди на прием не получая самого сообщения. Синтаксис ее выглядит следующим образом:
int MPI_Probe(int source, int tag, MPI_Comm comm, MPI_Status* status);
Тут мы также определяем от какого процесса получаемое сообщение, с каким тэгом, какой коммуникатор связывает эти процессы и передаем структуру которая запишет необходимую информацию.
Теперь на очень простом примере соединим эти процедуры вместе:
#include #include "mpi.h" using namespace std; void show_arr(int* arr, int size) < for(int i=0; i < size; i++) cout int main(int argc, char **argv) < int size, rank; MPI_Init(&argc, &argv); MPI_Comm_size(MPI_COMM_WORLD, &size); MPI_Comm_rank(MPI_COMM_WORLD, &rank); if(rank == 0) < int* arr = new int[size]; for(int i=0; i < size; i++) arr[i] = i; for(int i=1; i < size; i++) MPI_Send(arr, i, MPI_INT, i, 5, MPI_COMM_WORLD); >else < int count; MPI_Status status; MPI_Probe(MPI_ANY_SOURCE, MPI_ANY_TAG, MPI_COMM_WORLD, &status); MPI_Get_count(&status, MPI_INT, &count); int* buf = new int[count]; MPI_Recv(buf, count, MPI_INT, MPI_ANY_SOURCE, MPI_ANY_TAG, MPI_COMM_WORLD, &status); cout MPI_Finalize(); return 0; >
Что тут происходит?
В данной программе первый процесс создает массив размером равным количеству процессов и заполняет его номерами процессов по очереди. Потом соответствующему процессу он отправляет такое число элементов этого массива, какой номер у этого процесса. Напрмер: процесс 1 получит 1 элемент, процесс 2 получит 2 элемента этого массива и так далее.
Следующие же процессы должны принять это сообщение и для этого смотрят в очередь сообщений, видят это сообщение, собирают информацию в структуру status, после чего получают длину переданного сообщения и создают буфер для его сохранения, после чего просто выводят то, что получили. Для 5 запущенных процессов результат будет вот таким:
Process:1 || Count: 1 || Array: 0 Process:2 || Count: 2 || Array: 0 1 Process:4 || Count: 4 || Array: 0 1 2 3 Process:3 || Count: 3 || Array: 0 1 2
Собственно 4 результата потому что нулевой процесс занимается отправкой этих сообщений.
И еще несколько типов процедур посылки
Отметим, что вызов и возврат из процедуры MPI_Send далеко не гарантирует того, что сообщение покинуло данный процесс, было получено процессом, которому оно отправлено. Гарантия дается только на то, что мы после этой процедуры можем использовать передаваемые данные в своих целях не опасаясь того, что они как-то изменятся в отправленном сообщении.
Для того чтобы повысить степень определенности существуют еще несколько процедур которые никак не отличаются по синтаксису от MPI_Send, но отличаются по типу взаимодействия процессов и в способе отправки.
Первая из них это процедура MPI_Bsend. Такая процедура осуществляет передачу сообщения с буферизацией. Если процесс которому мы отправляем сообщение еще не запросил его получения, то информация будет записана в специальный буфер и процесс продолжит работу. Сообщение об ошибке возможно в случае, если места в буфере не хватит для сообщения, однако об этом размере может позаботиться и сам программист.
Вторая процедура это MPI_Ssend. Эта процедура синхронизирует потоки в процессе передачи сообщений. Возврат из этой процедуры произойдет ровно тогда, когда прием этого сообщения будет инициализирован процессом-получателем. То есть такая процедура заставляет процесс-отправитель ожидать приема сообщения, поэтому оба процесса участвующих в коммуникации продолжат работу после приема сообщения одновременно, а значит синхронизируются.
И еще одна процедура — MPI_Rsend. Она осуществляет передачу сообщения по готовности. Такой процедурой нужно пользоваться аккуратно, так как она требует чтобы процесс-получатель мыл уже готов принять это сообщение. Для того чтобы она выполнилась корректно необходимо заранее позаботиться о синхронизации процессов, либо явно знать, что этот процесс будет находиться в стадии ожидании получения сообщения к моменту вызова Rsend.
Резюме
Ну вот мы и ознакомились(а для опытных освежили в памяти) основные процедуры передачи сообщений типа точка-точка с блокировкой. В следующей статье я постараюсь показать на практике все изложенные ранее принципы и объяснить как написать программу которая будет выяснять одни из основополагающих характеристик техники используемой при параллелизации вычислений — латнетность и пропускная способность между процессами. А пока вот краткая сводка того что я здесь изложил:
Процедура/Константа/Структура
Назначение
Что такое SMS‑сообщение?

Аббревиатура SMS расшифровывается как Short Message Service (сервис коротких сообщений). Это служба текстовых сообщений, которая позволяет обмениваться короткими текстовыми сообщениями между мобильными устройствами. SMS-сообщения обычно имеют максимальную длину 160 символов и могут быть отправлены и получены в различных мобильных сетях. SMS широко используются для личного и делового общения, обеспечивая быстрый и удобный способ отправки кратких сообщений отдельным лицам или группам людей. Она стала неотъемлемой частью мобильной связи и поддерживается практически всеми мобильными устройствами.
Что означает SMS?
Аббревиатура SMS расшифровывается как Short Message Service (сервис коротких сообщений). Это служба текстовых сообщений, которая позволяет обмениваться короткими текстовыми сообщениями между мобильными устройствами.
Как работает SMS?
SMS работает на сигнальных каналах мобильных сетей, используя уже существующую инфраструктуру для голосовых вызовов. Ниже приведен упрощенный обзор того, как работает SMS.
- Отправитель инициирует сообщение: отправитель создает текстовое сообщение на своем мобильном устройстве и вводит номер телефона получателя.
- Сообщение, отправленное в SMSC: мобильное устройство отправителя отправляет сообщение в Центр обслуживания коротких сообщений (SMSC), который является централизованным сервером, отвечающим за обработку SMS-сообщений.
- Маршрутизация сообщений SMSC: SMSC проверяет телефонный номер получателя и определяет подходящую сеть для доставки сообщения.
- Доставка сообщения: затем SMSC отправляет сообщение в мобильную сеть получателя с помощью серии сигнальных сообщений.
- Сообщение, сохраненное в SMSC получателя: SMSC получателя принимает сообщение и временно сохраняет его до тех пор, пока устройство получателя не станет доступно для его получения.
- Уведомление на устройство получателя: как только устройство получателя будет доступно, SMSC получателя отправляет уведомление о доступности нового SMS.
- Получение сообщения: мобильное устройство получателя подключается к SMSC получателя для получения сообщения.
- Отображаемое сообщение: мобильное устройство получателя получает сообщение и отображает его.
- Дополнительное подтверждение доставки: мобильное устройство получателя может отправить подтверждение доставки обратно в SMSC отправителя, указывая, что сообщение было успешно получено.
Важно отметить, что SMS-сообщения обычно передаются по каналам управления и не используют те же голосовые каналы или каналы передачи данных, которые используются в других мобильных сервисах. Это позволяет SMS быть надежной и широко поддерживаемой формой связи даже в районах с ограниченным покрытием сети или во время перегрузки сети.
Как работает SMS-маркетинг?
SMS-маркетинг часто использует API SMS для автоматизации и оптимизации процесса. Вот обзор того, как работает SMS-маркетинг с помощью API SMS:
- Интеграция с API SMS: для реализации SMS-маркетинга компании интегрируют свои системы или приложения с API SMS. API служит интерфейсом, позволяющим программному обеспечению компании отправлять и получать SMS-сообщения программно.
- Составление списка подписчиков: компании собирают информацию о подписке от лиц, желающих получать маркетинговые SMS-сообщения. Эти данные можно собирать с помощью различных каналов, таких как веб-формы, мобильные приложения или регистрация в магазине. API позволяет легко интегрировать эти данные в список подписчиков.
- Создание и персонализация сообщений. Используя API SMS, компании создают персонализированные и целевые SMS-сообщения для взаимодействия со своими подписчиками. API позволяет динамически вставлять контент, позволяя настраивать сообщения на основе информации о подписчиках, такой как имена, история покупок или предпочтения.
- Запуск автоматических сообщений: API SMS позволяют автоматизировать отправку SMS-сообщений на основе заранее определенных триггеров или событий. Например, компании могут настроить автоматические приветственные сообщения для новых подписчиков или транзакционные сообщения для подтверждения заказов или обновлений касательно доставки.
- Отправка массовых SMS-кампаний: с помощью API SMS компании могут отправлять массовые SMS-кампании своему списку подписчиков. Они могут сегментировать свою аудиторию на основе демографических данных, предпочтений или поведения и отправлять целевые сообщения определенным сегментам, максимально повышая эффективность своих маркетинговых кампаний.
- Доставка и отслеживание сообщений: когда SMS отправляется через API SMS, сообщение передается на мобильные устройства абонентов через мобильную сеть. API управляет процессом доставки и предоставляет отчеты о доставке или обновления статуса, что позволяет компаниям отслеживать успех своих SMS-кампаний.
- Взаимодействие с абонентами и двусторонняя связь: API SMS позволяют компаниям упрощать двустороннюю связь со своими подписчиками. Подписчики могут отвечать на SMS-сообщения, что позволяет компаниям общаться в режиме реального времени, оказывать поддержку или собирать отзывы.
- Соответствие требованиям и нормативные требования. При использовании API SMS для SMS-маркетинга крайне важно соблюдать применимые правила, такие как получение надлежащего согласия и предоставление механизмов отказа от рассылки. API SMS часто включают функции, отвечающие требованиям соответствия, такие как управление отписками и обработка запросов на отказ от подписки.
Что такое API SMS?
API SMS, или прикладной программный интерфейс службы коротких сообщений, представляет собой механизм, позволяющий программным системам отправлять и получать SMS-сообщения программно. Он предоставляет набор определений и протоколов, обеспечивающих связь между различными программными компонентами. Подобно тому, как приложение погоды на вашем телефоне взаимодействует с программной системой метеорологического бюро с помощью API для отображения ежедневных обновлений погоды, API SMS позволяет приложениям взаимодействовать с SMS-сервисами и отправлять сообщения получателям.
Как работает API SMS?
API SMS использует архитектуру клиент-сервер. Приложение, которое инициирует запрос на отправку SMS, является клиентом, а сервер обрабатывает запрос и отправляет SMS предполагаемому получателю. API SMS выступает в качестве интерфейса между клиентским приложением и поставщиком услуг SMS.
Когда клиентское приложение хочет отправить SMS, оно отправляет запрос к API SMS с необходимой информацией, такой как номер телефона получателя и содержимое сообщения. API проверяет запрос, взаимодействует с поставщиком услуг SMS и доставляет сообщение получателю. Аналогичным образом, когда клиентское приложение хочет получать SMS-сообщения, оно может использовать API для получения входящих сообщений от поставщика услуг SMS.
Каковы преимущества использования API SMS?
Использование API SMS дает несколько преимуществ.
Автоматизация: API SMS позволяет автоматически отправлять и получать SMS-сообщения, что позволяет сэкономить время и силы по сравнению с ручными процессами.
Интеграция: API упрощают интеграцию функций SMS в существующие программные системы, позволяя компаниям включать SMS в свою коммуникационную стратегию.
Масштабируемость: с помощью API SMS компании могут легко масштабировать свои возможности SMS в соответствии с растущими коммуникационными потребностями, будь то отправка сообщений большой клиентской базе или обработка увеличенного объема сообщений.
Персонализация: API позволяют настраивать и персонализировать SMS-сообщения, интегрируя в них динамический контент, такой как имена клиентов или сведения о заказе.
Общение в реальном времени. Используя API SMS, компании могут общаться со своими клиентами в режиме реального времени, мгновенно доставляя важные уведомления, оповещения или рекламные сообщения.
Где найти поставщиков SMS-маркетинга?
API SMS доступны у различных поставщиков услуг SMS. Вы можете найти и изучить различные API SMS на торговых площадках или в каталогах API, посвященных демонстрации и предложению широкого спектра API. Некоторые популярные поставщики услуг SMS включают Amazon Pinpoint, Amazon SNS, Twilio и Sinch.
Не забудьте выбрать API SMS, соответствующие вашим конкретным потребностям, учитывая такие факторы, как цена, надежность, функции и качество документации.
Как использовать API SMS?
Чтобы использовать API SMS, выполните указанные ниже действия.
- Получите учетные данные API: зарегистрируйтесь у поставщика услуг SMS, который предлагает API SMS, и получите необходимые учетные данные API, такие как ключ API или токен доступа.
- Интегрируйте API: в зависимости от используемого языка программирования или платформы интегрируйте API SMS в код приложения. Обычно это связано с отправкой HTTP-запросов на адреса API, предоставленные поставщиком услуг SMS.
- Отправка SMS-сообщений: используйте методы или адреса API для отправки SMS-сообщений. Укажите необходимые параметры, такие как номер телефона получателя и содержимое сообщения, в запросе API.
- Обработка ответов: после отправки API SMS предоставит ответ с указанием статуса доставки сообщения. Обработайте эти ответы в своем приложении, чтобы обеспечить успешную доставку сообщений и устранить возможные ошибки.
- Получение SMS-сообщений (опционально): если ваш вариант использования связан с получением SMS-сообщений, API SMS может предоставить адреса или веб-хуки для получения входящих сообщений. Настройте приложение для прослушивания входящих сообщений и соответствующей их обработки.
Как AWS может удовлетворить ваши потребности в SMS?
Amazon Pinpoint – это многоканальный сервис связи и привлечения клиентов, предлагаемый AWS. Он предоставляет компаниям инструменты для эффективного привлечения клиентов по различным каналам, включая SMS, push-уведомления и голосовые сообщения. С помощью Amazon Pinpoint компании могут создавать целевые и персонализированные кампании, отслеживать взаимодействие с пользователями и анализировать эффективность кампаний для оптимизации стратегий привлечения клиентов. Amazon Pinpoint в первую очередь помогает компаниям отправлять SMS, push-уведомления или голосовые сообщения людям или конечным клиентам, что часто называют обменом сообщениями от приложения к пользователю (A2P). Amazon Pinpoint также позволяет компаниям создавать индивидуальные маршруты, чтобы клиенты получали нужное сообщение в нужное время.
Amazon Simple Notification Service (SNS) – это гибкий и масштабируемый сервис обмена сообщениями, предоставляемый AWS. Это позволяет разработчикам отправлять сообщения или уведомления большому количеству подписчиков или адресов, таких как мобильные устройства, адреса электронной почты или распределенные системы. Amazon SNS в первую очередь помогает компаниям отправлять сообщения A2A или из приложения в приложение с помощью автоматических уведомлений. Это позволяет приложениям взаимодействовать и обмениваться информацией друг с другом, обеспечивая беспрепятственную связь и интеграцию между различными программными компонентами.
Начните работу с SaaS на AWS, создав бесплатный аккаунт уже сегодня.
Apache Kafka: основы технологии
У Kafka есть множество способов применения, и у каждого способа есть свои особенности. В этой статье разберём, чем Kafka отличается от популярных систем обмена сообщениями; рассмотрим, как Kafka хранит данные и обеспечивает гарантию сохранности; поймём, как записываются и читаются данные.
Статья подготовлена на основе открытого занятия из видеокурса по Apache Kafka. Авторы — Анатолий Солдатов, Lead Engineer в Авито, и Александр Миронов, Infrastructure Engineer в Stripe. Базовые темы курса доступны на Youtube.
Kafka и классические сервисы очередей
Для первого погружения в технологию сравним Kafka и классические сервисы очередей, такие как RabbitMQ и Amazon SQS.
Системы очередей обычно состоят из трёх базовых компонентов:
1) сервер,
2) продюсеры, которые отправляют сообщения в некую именованную очередь, заранее сконфигурированную администратором на сервере,
3) консьюмеры, которые считывают те же самые сообщения по мере их появления.

Базовые компоненты классической системы очередей
В веб-приложениях очереди часто используются для отложенной обработки событий или в качестве временного буфера между другими сервисами, тем самым защищая их от всплесков нагрузки.
Консьюмеры получают данные с сервера, используя две разные модели запросов: pull или push.

pull-модель — консьюмеры сами отправляют запрос раз в n секунд на сервер для получения новой порции сообщений. При таком подходе клиенты могут эффективно контролировать собственную нагрузку. Кроме того, pull-модель позволяет группировать сообщения в батчи, таким образом достигая лучшей пропускной способности. К минусам модели можно отнести потенциальную разбалансированность нагрузки между разными консьюмерами, а также более высокую задержку обработки данных.
push-модель — сервер делает запрос к клиенту, посылая ему новую порцию данных. По такой модели, например, работает RabbitMQ. Она снижает задержку обработки сообщений и позволяет эффективно балансировать распределение сообщений по консьюмерам. Но для предотвращения перегрузки консьюмеров в случае с RabbitMQ клиентам приходится использовать функционал QS, выставляя лимиты.
Как правило, приложение пишет и читает из очереди с помощью нескольких инстансов продюсеров и консьюмеров. Это позволяет эффективно распределить нагрузку.

Типичный жизненный цикл сообщений в системах очередей:
- Продюсер отправляет сообщение на сервер.
- Консьюмер фетчит (от англ. fetch — принести) сообщение и его уникальный идентификатор сервера.
- Сервер помечает сообщение как in-flight. Сообщения в таком состоянии всё ещё хранятся на сервере, но временно не доставляются другим консьюмерам. Таймаут этого состояния контролируется специальной настройкой.
- Консьюмер обрабатывает сообщение, следуя бизнес-логике. Затем отправляет ack или nack-запрос обратно на сервер, используя уникальный идентификатор, полученный ранее — тем самым либо подтверждая успешную обработку сообщения, либо сигнализируя об ошибке.
- В случае успеха сообщение удаляется с сервера навсегда. В случае ошибки или таймаута состояния in-flight сообщение доставляется консьюмеру для повторной обработки.

Типичный жизненный цикл сообщений в системах очередей
С базовыми принципами работы очередей разобрались, теперь перейдём к Kafka. Рассмотрим её фундаментальные отличия.
Как и сервисы обработки очередей, Kafka условно состоит из трёх компонентов:
1) сервер (по-другому ещё называется брокер),
2) продюсеры — они отправляют сообщения брокеру,
3) консьюмеры — считывают эти сообщения, используя модель pull.

Базовые компоненты Kafka
Пожалуй, фундаментальное отличие Kafka от очередей состоит в том, как сообщения хранятся на брокере и как потребляются консьюмерами.
- Сообщения в Kafka не удаляются брокерами по мере их обработки консьюмерами — данные в Kafka могут храниться днями, неделями, годами.
- Благодаря этому одно и то же сообщение может быть обработано сколько угодно раз разными консьюмерами и в разных контекстах.
В этом кроется главная мощь и главное отличие Kafka от традиционных систем обмена сообщениями.
Теперь давайте посмотрим, как Kafka и системы очередей решают одну и ту же задачу. Начнём с системы очередей.
Представим, что есть некий сайт, на котором происходит регистрация пользователя. Для каждой регистрации мы должны:
1) отправить письмо пользователю,
2) пересчитать дневную статистику регистраций.
В случае с RabbitMQ или Amazon SQS функционал может помочь нам доставить сообщения всем сервисам одновременно. Но при необходимости подключения нового сервиса придётся конфигурировать новую очередь.

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

Kafka также позволяет тривиально подключать новые сервисы к стриму регистрации. Например, сервис архивирования всех регистраций в S3 для последующей обработки с помощью Spark или Redshift можно добавить без дополнительного конфигурирования сервера или создания дополнительных очередей.
Кроме того, раз Kafka не удаляет данные после обработки консьюмерами, эти данные могут обрабатываться заново, как бы отматывая время назад сколько угодно раз. Это оказывается невероятно полезно для восстановления после сбоев и, например, верификации кода новых консьюмеров. В случае с RabbitMQ пришлось бы записывать все данные заново, при этом, скорее всего, в отдельную очередь, чтобы не сломать уже имеющихся клиентов.
Структура данных
Наверняка возникает вопрос: «Раз сообщения не удаляются, то как тогда гарантировать, что консьюмер не будет читать одни и те же сообщения (например, при перезапуске)?».
Для ответа на этот вопрос разберёмся, какова внутренняя структура Kafka и как в ней хранятся сообщения.
Каждое сообщение (event или message) в Kafka состоит из ключа, значения, таймстампа и опционального набора метаданных (так называемых хедеров).

Сообщения в Kafka организованы и хранятся в именованных топиках (Topics), каждый топик состоит из одной и более партиций (Partition), распределённых между брокерами внутри одного кластера. Подобная распределённость важна для горизонтального масштабирования кластера, так как она позволяет клиентам писать и читать сообщения с нескольких брокеров одновременно.
Когда новое сообщение добавляется в топик, на самом деле оно записывается в одну из партиций этого топика. Сообщения с одинаковыми ключами всегда записываются в одну и ту же партицию, тем самым гарантируя очередность или порядок записи и чтения.
Для гарантии сохранности данных каждая партиция в Kafka может быть реплицирована n раз, где n — replication factor. Таким образом гарантируется наличие нескольких копий сообщения, хранящихся на разных брокерах.

У каждой партиции есть «лидер» (Leader) — брокер, который работает с клиентами. Именно лидер работает с продюсерами и в общем случае отдаёт сообщения консьюмерам. К лидеру осуществляют запросы фолловеры (Follower) — брокеры, которые хранят реплику всех данных партиций. Сообщения всегда отправляются лидеру и, в общем случае, читаются с лидера.
Чтобы понять, кто является лидером партиции, перед записью и чтением клиенты делают запрос метаданных от брокера. Причём они могут подключаться к любому брокеру в кластере.

Основная структура данных в Kafka — это распределённый, реплицируемый лог. Каждая партиция — это и есть тот самый реплицируемый лог, который хранится на диске. Каждое новое сообщение, отправленное продюсером в партицию, сохраняется в «голову» этого лога и получает свой уникальный, монотонно возрастающий offset (64-битное число, которое назначается самим брокером).
Как мы уже выяснили, сообщения не удаляются из лога после передачи консьюмерам и могут быть вычитаны сколько угодно раз.
Время гарантированного хранения данных на брокере можно контролировать с помощью специальных настроек. Длительность хранения сообщений при этом не влияет на общую производительность системы. Поэтому совершенно нормально хранить сообщения в Kafka днями, неделями, месяцами или даже годами.

Consumer Groups
Теперь давайте перейдём к консьюмерам и рассмотрим их принципы работы в Kafka. Каждый консьюмер Kafka обычно является частью какой-нибудь консьюмер-группы.
Каждая группа имеет уникальное название и регистрируется брокерами в кластере Kafka. Данные из одного и того же топика могут считываться множеством консьюмер-групп одновременно. Когда несколько консьюмеров читают данные из Kafka и являются членами одной и той же группы, то каждый из них получает сообщения из разных партиций топика, таким образом распределяя нагрузку.
Вернёмся к нашему примеру с топиком сервиса регистрации и представим, что у сервиса отправки писем есть своя собственная консьюмер-группа с одним консьюмером c1 внутри. Значит, этот консьюмер будет получать сообщения из всех партиций топика.

Если мы добавим ещё одного консьюмера в группу, то партиции автоматически распределятся между ними, и c1 теперь будет читать сообщения из первой и второй партиции, а c2 — из третьей. Добавив ещё одного консьюмера (c3), мы добьёмся идеального распределения нагрузки, и каждый из консьюмеров в этой группе будет читать данные из одной партиции.

А вот если мы добавим в группу ещё одного консьюмера (c4), то он не будет задействован в обработке сообщений вообще.
Важно понять: внутри одной консьюмер-группы партиции назначаются консьюмерам уникально, чтобы избежать повторной обработки.
Если консьюмеры не справляются с текущим объёмом данных, то следует добавить новую партицию в топик. Только после этого консьюмер c4 начнёт свою работу.
Механизм партиционирования является нашим основным инструментом масштабирования Kafka. Группы являются инструментом отказоустойчивости.
Кстати, как вы думаете, что будет, если один из консьюмеров в группе упадёт? Совершенно верно: партиции автоматически распределятся между оставшимися консьюмерами в этой группе.
Добавлять партиции в Kafka можно на лету, без перезапуска клиентов или брокеров. Клиенты автоматически обнаружат новую партицию благодаря встроенному механизму обновления метаданных. Однако, нужно помнить две важные вещи:
- Гарантия очерёдности данных — если вы пишете сообщения с ключами и хешируете номер партиции для сообщений, исходя из общего числа, то при добавлении новой партиции вы можете просто сломать порядок этой записи.
- Партиции невозможно удалить после их создания, можно удалить только весь топик целиком.
И ещё неочевидный момент: если вы добавляете новую партицию на проде, то есть в тот момент, когда в топик пишут сообщения продюсеры, то важно помнить про настройку auto.offset.reset=earliest в консьюмере, иначе у вас есть шанс потерять или просто не обработать кусок данных, записавшихся в новую партицию до того, как консьюмеры обновили метаданные по топику и начали читать данные из этой партиции.
Помимо этого, механизм групп позволяет иметь несколько несвязанных между собой приложений, обрабатывающих сообщения.

Как мы обсуждали ранее, можно добавить новую группу консьюмеров к тому же самому топику, например, для обработки и статистики регистраций. Эти две группы будут читать одни и те же сообщения из топика тех самых ивентов регистраций — в своём темпе, со своей внутренней логикой.
А теперь, зная внутреннее устройство консьюмеров в Kafka, давайте вернёмся к изначальному вопросу: «Каким образом мы можем обозначить сообщения в партиции, как обработанные?».
Для этого Kafka предоставляет механизм консьюмер-офсетов. Как мы помним, каждое сообщение партиции имеет свой собственный, уникальный, монотонно возрастающий офсет. Именно этот офсет и используется консьюмерами для сохранения партиций.
Консьюмер делает специальный запрос к брокеру, так называемый offset-commit с указанием своей группы, идентификатора топик-партиции и, собственно, офсета, который должен быть отмечен как обработанный. Брокер сохраняет эту информацию в своём собственном специальном топике. При рестарте консьюмер запрашивает у сервера последний закоммиченный офсет для нужной топик-партиции, и просто продолжает чтение сообщений с этой позиции.
В примере консьюмер в группе email-service-group, читающий партицию p1 в топике registrations, успешно обработал три сообщения с офсетами 0, 1 и 2. Для сохранения позиций консьюмер делает запрос к брокеру, коммитя офсет 3. В случае рестарта консьюмер запросит свою последнюю закоммиченную позицию у брокера и получит в ответе 3. После чего начнёт читать данные с этого офсета.

Консьюмеры вольны коммитить совершенно любой офсет (валидный, который действительно существует в этой топик-партиции) и могут начинать читать данные с любого офсета, двигаясь вперёд и назад во времени, пропуская участки лога или обрабатывая их заново.
Ключевой для понимания факт: в момент времени может быть только один закоммиченный офсет для топик-партиции в консьюмер-группе. Иными словами, мы не можем закоммитить несколько офсетов для одной и той же топик-партиции, эмулируя каким-то образом выборочный acknowledgment (как это делалось в системах очередей).
Представим, что обработка сообщения с офсетом 1 завершилась с ошибкой. Однако мы продолжили выполнение нашей программы в консьюмере и запроцессили сообщение с офсетом 2 успешно. В таком случае перед нами будет стоять выбор: какой офсет закоммитить — 1 или 3. В настоящей системе мы бы рекомендовали закоммитить офсет 3, добавив при этом функционал, отправляющий ошибочное сообщение в отдельный топик для повторной обработки (ручной или автоматической). Подобные алгоритмы называются Dead letter queue.
Разумеется, консьюмеры, находящиеся в разных группах, могут иметь совершенно разные закоммиченные офсеты для одной и той же топик-партиции.
Apache ZooKeeper
В заключение нужно упомянуть об ещё одном важном компоненте кластера Kafka — Apache ZooKeeper.
ZooKeeper выполняет роль консистентного хранилища метаданных и распределённого сервиса логов. Именно он способен сказать, живы ли ваши брокеры, какой из брокеров является контроллером (то есть брокером, отвечающим за выбор лидеров партиций), и в каком состоянии находятся лидеры партиций и их реплики.
В случае падения брокера именно в ZooKeeper контроллером будет записана информация о новых лидерах партиций. Причём с версии 1.1.0 это будет сделано асинхронно, и это важно с точки зрения скорости восстановления кластера. Самый простой способ превратить данные в тыкву — потеря информации в ZooKeeper. Тогда понять, что и откуда нужно читать, будет очень сложно.
В настоящее время ведутся активные работы по избавлению Kafka от зависимости в виде ZooKeeper, но пока он всё ещё с нами (если интересно, посмотрите на Kafka improvement proposal 500, там подробно расписан план избавления от ZooKeeper).
Важно помнить, что ZooKeeper по факту является ещё одной распределённой системой хранения данных, за которой необходимо следить, поддерживать и обновлять по мере необходимости.
Традиционно ZooKeeper раскатывается отдельно от брокеров Kafka, чтобы разделить границы возможных отказов. Помните, что падение ZooKeeper — это практически падение всего кластера Kafka. К счастью, нагрузка на ZooKeeper при нормальной работе кластера минимальна. Клиенты Kafka никогда не коннектятся к ZooKeeper напрямую.