Как использовать Kafka для создания приложений, основанных на сильных событиях
Понимание Apache Kafka и его роли в архитектуре, основанной на событиях
Apache Kafka - это распределенная платформа потокового воспроизведения событий, способная обрабатывать триллионы событий в день. Изначально разработанная в LinkedIn, Kafka стала основой современных архитектур, управляемых событиями, позволяя приложениям публиковать, хранить, обрабатывать и реагировать на потоки данных в режиме реального времени. Его способность сочетать высокую пропускную способность, отказоустойчивость и горизонтальную масштабируемость делает его идеальным выбором для создания надежных, производственных систем, управляемых событиями. Независимо от того, синхронизируете ли вы микросервисы, питаете аналитику в реальном времени или строите конвейер данных между устаревшими системами и современными приложениями, Kafka обеспечивает прочную, устойчивую основу, необходимую для этих требовательных рабочих нагрузок.
Что отличает Kafka от традиционных очередей сообщений, так это его основной дизайн в виде распределенного журнала фиксации. Вместо удаления сообщений после потребления Kafka сохраняет их на настраиваемый период (или навсегда), позволяя нескольким потребителям воспроизводить или перерабатывать события. Это разделение производителей и потребителей означает, что каждая сторона может масштабироваться независимо, а сбои в одной части системы не каскадируются. Для приложений, управляемых событиями, этот архитектурный выбор напрямую переводится в надежность: вы можете добавлять новых потребителей, не нарушая существующих, и вы можете оправиться от сбоев, просто перечитав из известного смещения.
Основные компоненты Kafka: более глубокое погружение
Чтобы создать надежные приложения, основанные на событиях, с помощью Kafka, вы должны сначала понять его основные строительные блоки. Каждый компонент играет решающую роль в производительности и надежности платформы:
- Топики — это логические каналы, по которым публикуются записи.Тема может иметь любое количество разделов, а стратегия разделения определяет, как данные распределяются между брокерами.
- Разделы являются единицей параллелизма и упорядочения.В рамках раздела записи строго упорядочены за счет смещения. Производители могут выбрать ключ раздела (например, идентификатор пользователя), чтобы гарантировать, что все события для одного и того же ключа переходят в один и тот же раздел, сохраняя порядок для этого объекта.
- Производители публикуют записи по темам. Они могут настроить подтверждения (аки) для балансировки скорости по сравнению с долговечностью:
- — без подтверждения, быстрее, но с риском потери данных.
- — лидер признает, хороший баланс.
- — все реплики в синхронизации признают, самую сильную долговечность.
- Потребители читают записи с разделов. Они принадлежат к группе потребителей, что позволяет балансировать нагрузку: каждый раздел назначается ровно одному потребителю в группе. Если потребитель терпит неудачу, разделы перебалансируются на остальные члены, гарантируя, что никакие данные не останутся необработанными.
- Брокеры — это серверы Kafka, которые хранят данные и обслуживают запросы клиентов. Кластер Kafka обычно состоит из нескольких брокеров. Каждый раздел реплицируется через настраиваемое число брокеров (фактор репликации) для обеспечения отказоустойчивости. Набор in-sync реплики (ISR) гарантирует, что для лидерства рассматриваются только полностью догоняемые реплики.
Понимание того, как эти компоненты взаимодействуют, имеет решающее значение для разработки развертывания Kafka, которое соответствует требованиям вашего приложения к пропускной способности, задержке, долговечности и согласованности.
Создание Kafka для потокового вещания событий
Настройка разработки с одним брокером хороша для обучения, но надежное приложение, ориентированное на события, требует производственной конфигурации. Вот ключевые шаги и соображения:
Кластерный размер и конфигурация брокера
Начните с по крайней мере трех брокеров, чтобы обеспечить кворум для выборов лидера и разрешить техническое обслуживание без простоев. Настройте коэффициент репликации до 3 для критических тем. Настройте до 2, чтобы гарантировать, что по крайней мере две реплики признают записи при использовании . Настройте политику хранения журнала на основе ваших потребностей в хранении данных. Например, (7 дней) является общим для многих потоковых рабочих нагрузок.
Тема Дизайн и стратегия разделения
Количество разделов определяет максимальный параллелизм как для производителей, так и для потребителей. Хорошее эмпирическое правило состоит в том, чтобы начинать с 10-50 разделов по теме в зависимости от ожидаемой пропускной способности. Каждый раздел по существу является файлом, поэтому слишком много разделов могут привести к накладным расходам на обработку файлов и увеличению нагрузки на Zookeeper. Рассмотрите возможность использования руководящих принципов размера раздела FLT:0 для вашей конкретной рабочей нагрузки. Используйте значимые ключи раздела (например, идентификатор заказа, идентификатор клиента) для сохранения порядка в потоке событий объекта.
Интеграция с реестром сменных схем
Для поддержания совместимости данных по мере развития ваших схем событий интегрируйте Реестр сменных схем. Этот сервис хранит определения Avro, Protobuf или JSON Schema и обеспечивает соблюдение правил совместимости (обратно, вперед, полный). Производители и потребители ссылаются на идентификатор схемы, а не на встраивание полных схем, уменьшая накладные расходы сети. Например, производитель может отправить сообщение с кодом Protobuf вместе с идентификатором схемы, и потребитель использует реестр схем для его декодирования. Это важно для надежных, долгоживущих систем, управляемых событиями, где несколько команд владеют различными частями трубопровода.
Внедрение производителей и потребителей с передовой практикой
Kafka предлагает богатые клиентские библиотеки для Java, Python, Go, .NET и многих других языков. Следующие примеры используют Java, но шаблоны применяются повсеместно.
Создать надежного производителя
Надежный производитель должен обрабатывать повторы, идемпотентность и транзакционную семантику:
- Включить идемпотентность, установив . Это предотвращает дублирование записей в случае повторных записей, обеспечивая точно один раз семантику для однораздельных записей.
- , , , , , , , , [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]], [[FLT:]],
- Используйте асинхронные отправки с обратной связью для изящной обработки сбоев: регистрируйте ошибку, оповещение или маршрут к теме с мертвой буквой.
- Выберите разделитель, равномерно распределяющий нагрузку. По умолчанию липкий разделитель повышает эффективность пакетирования.
Пример фрагмента (псевдокод):
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.ByteArraySerializer");
props.put("enable.idempotence", true);
props.put("acks", "all");
props.put("retries", Integer.MAX_VALUE);
KafkaProducer<String, byte[]> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("orders", orderKey, orderBytes), (metadata, exception) -> {
if (exception != null) {
// handle exception – log, alert, send to DLT
}
});
Создание устойчивого потребителя
Потребители должны изящно справляться с перебалансировкой, управлять смещениями и idempotently обрабатывать:
- Установите и вручную совершите смещение после обработки партии. Это предотвращает потерю данных, если потребитель терпит крах перед совершением.
- Используйте , чтобы контролировать размер партии и избегать обработки слишком большого количества записей перед совершением.
- Внедрить перебалансировщик для хранения смещений до отзыва раздела и стремиться к сохранению смещений при назначении.
- Сделайте обработку идемпотентной, чтобы дубликаты от переработки не вызывали побочных эффектов. Например, размножение по идентификатору события или использование базы данных.
Для высокой пропускной способности рассмотрите возможность использования петли опроса , которая обрабатывает записи параллельно с использованием пула потоков, но гарантирует, что офсетные фиксации происходят только после обработки всех записей в партии. Потребительская документация Apache Kafka обеспечивает глубокое погружение в эти механики.
Обработка событий с помощью Kafka Streams и KSQL
Помимо простого производства/потребления, Kafka предоставляет первоклассные возможности обработки потоков.
Кафка Стримс
Kafka Streams - это клиентская библиотека для создания государственных потоковых приложений. Она работает как стандартное приложение (без отдельного кластера) и использует собственные темы Kafka для государственных магазинов и журналов изменений. Ключевые функции включают в себя:
- Ровно-однократная семантика для государственных операций (соединений, агрегации).
- Родная поддержка окон (скакивание, прыжки, окна сеанса).
- Процессор API и DSL (например, ).
Например, вы можете вычислить общую сумму заказов на клиента, создав KTable из темы заказа и используя оператора . Kafka Streams автоматически обрабатывает государственный магазин и журнал изменений, делая ваше приложение автоматически устойчивым к сбоям - если узел падает, состояние перестраивается из темы журнала изменений.
KSQL (Kafka SQL)
KSQL — это движок потокового SQL для Kafka. Он позволяет запускать SQL-подобные запросы на потоковых данных без написания кода Java. Используйте его для специального анализа, прототипирования или простого ETL. Например:
CREATE STREAM orders WITH (KAFKA_TOPIC='orders', VALUE_FORMAT='JSON');
CREATE TABLE high_value_orders AS
SELECT customer_id, COUNT(*) AS order_count, SUM(amount) AS total
FROM orders WINDOW TUMBLING (SIZE 1 HOUR)
WHERE amount > 1000
GROUP BY customer_id;
KSQL особенно полезен для групп по обработке данных, которые хотят быстро создавать преобразования, основанные на событиях.
Лучшие практики для создания надежных производственных систем
Устойчивое приложение, основанное на событиях, выходит за рамки простого написания статей для производителей и потребителей. Это требует целостного подхода к проектированию, операциям и мониторингу.
Обработка ошибок и очереди из мертвых писем
Даже при наличии надежных потребителей некоторые записи будут необработанными (например, неполноценный JSON, временные отключения в нисходящем потоке). Внедрите шаблон, в котором потребитель улавливает исключения, регистрирует исходную запись и публикует ее на тему с мертвой буквой (например, ]. Отдельный процесс может позже воспроизвести эти записи после расследования. Это гарантирует, что основной поток никогда не блокируется ядовитыми таблетками.
Гарантия точной семантики один раз
Для приложений, где дубликаты неприемлемы (например, финансовые транзакции), используйте семантику Kafka (EOS) для производителей и потребителей. На стороне производителя, как упоминалось, не гарантирует дубликатов в течение сессии. На стороне потребителя используйте транзакционный API для записи как выходных записей, так и смещений атомарно. Альтернативно, реализуйте идемпотентных потребителей с помощью таблицы дедупликации во внешней базе данных.
Мониторинг и наблюдаемость
Kafka предоставляет множество метрик через JMX. Мониторинг ключевых метрик:
- Подреплицированные разделы: Указывает на проблему с репликацией.
- Отставание потребителей: Разница между последним офсетом и совершенным офсетом потребителя. Высокий отставание означает, что потребители отстают.
- Запросить время ожидания: Время производить или потреблять.
Используйте такие инструменты, как Prometheus с экспортером Kafka JMX, для сбора метрик и настройки приборных панелей в Grafana. Кроме того, включите встроенный анализатор журналов Kafka (например, ) для отладки.
Лучшие практики безопасности
Защитите свои данные в пути и в покое:
- Аутентификация: Используйте SASL/SCRAM или SASL/SSL для аутентификации клиента.
- Авторизация: Определите ACL, чтобы контролировать, какие пользователи могут читать/писать на темы.
- Шифрование: Включить TLS/SSL для общения между клиентом и брокером и брокером.
- Сетевые политики: Используйте брандмауэры и VPC для ограничения доступа к брокерам.
Конфлюентная документация безопасности для всеобъемлющего руководства.
Масштабирование и настройка
По мере роста объема событий вам может потребоваться настроить количество разделов, увеличить коэффициент репликации или добавить брокеров. План емкости путем мониторинга использования диска, сетевого ввода / вывода и процессора. Используйте инструмент Kafka для перебалансировки данных по новым брокерам. Для сценариев с высокой пропускной способностью настройте размеры партий (, ]) для производителей и размеры приносят для потребителей. Настройки буферной памяти и гнезда также требуют внимания.
Реальные случаи использования и шаблоны
Чтобы проиллюстрировать, как эти концепции объединяются, рассмотрим типичную платформу электронной коммерции, которая использует Kafka в качестве центральной нервной системы.
- Служба заказов публикует события «OrderPlaced» на тему .
- Служба инвентаризации потребляет эти события для резервирования запасов, а затем публикует «Запас инвентаря» или «Запас вне магазина».
- Платежный сервис потребляет «Инвентарный резерв» событий и обрабатывает платежи, публикуя «Завершенный платеж».
- Служба уведомлений использует «PaymentCompleted» и отправляет подтверждения по электронной почте / SMS.
- Аналитический сервис использует все события заказа для создания панели инструментов в реальном времени.
- Приложение Kafka Streams присоединяется к потокам событий для обнаружения мошеннических схем (например, слишком много заказов с одного и того же IP за короткое время).
В этой архитектуре каждая услуга масштабируется независимо. Если Служба уведомлений не работает для обслуживания, события остаются в Kafka и обрабатываются позже. Если Платежная служба не справляется после совершения, событие PaymentCompleted обеспечивает идемпотентное восстановление. Использование реестра схем гарантирует, что при добавлении Службой заказов нового поля (например, «код дисконта») службы нисходящего потока не будут немедленно нарушены.
Другой распространенный шаблон - это шаблон Event Sourcing, где основным источником истины является сам поток событий. Журнал Kafka служит хранилищем событий. Государственные службы восстанавливают свое состояние, воспроизводя события с самого начала (или с момента снимка). Этот шаблон обеспечивает полный контрольный след и возможность задним числом исправлять ошибки, воспроизводя исправленные события.
Сравнение с другими технологиями, основанными на событиях
Хотя Кафка и мощная, но это не единственное решение. Понимание того, когда использовать ее против альтернатив, поможет вам сделать правильный архитектурный выбор:
- RabbitMQ отличается низкой задержкой, обменом сообщениями между точками со сложной маршрутизацией (обменами, привязками). Он легче для небольших развертываний, но не имеет гарантий прочности Kafka и возможности воспроизведения. Используйте RabbitMQ, когда вам нужна гарантированная доставка одному потребителю с низкими накладными расходами.
- Amazon Kinesis — управляемый потоковый сервис, похожий на Kafka, но он устраняет эксплуатационные накладные расходы. Однако он может иметь более высокую стоимость в масштабе и меньшую гибкость в настройке. Kafka предлагает больше возможностей управления и локального развертывания.
- Apache Pulsar обеспечивает многоуровневое хранение и многопользовательское использование нативных систем, но имеет меньшее сообщество и меньше инструментов экосистемы. зрелость Kafka, массовое сообщество и обширные клиентские библиотеки часто делают его более безопасным выбором для крупномасштабных событийных систем.
В конечном счете, Kafka лучше всего подходит для приложений, которые требуют упорядоченных, прочных, воспроизводимых потоков событий с высокой пропускной способностью и низкой задержкой, особенно при интеграции нескольких микросервисов или создании озера данных.
Заключение
Создание надежных приложений, основанных на событиях, с Apache Kafka требует не только понимания его API - он требует тщательного понимания его архитектуры, тщательной настройки для производства и соблюдения лучших практик для обработки ошибок, мониторинга и безопасности. Используя основные компоненты Kafka (темы, разделы, производители, потребители, брокеры) и расширенные возможности, такие как Kafka Streams и реестр схем, вы можете создавать системы, которые устойчивы при сбоях, масштабируемы до высоких нагрузок и обслуживаемы с течением времени.
Начните с тщательного моделирования ваших событий, тщательного проектирования ваших тем с учетом будущего роста и всегда планируйте неожиданные: сетевые разделы, сбои брокеров и изменения схемы. С Kafka вы получаете возможность разъединять службы, включать поток данных в реальном времени и создавать приложения, которые не только выживают, но и процветают в условиях сложности. Для дальнейшего чтения изучите документацию Apache Kafka и Библиотека ресурсов для глубоких руководств и эталонных архитектур.