Форум программистов, компьютерный форум, киберфорум
Javaican
Войти
Регистрация
Восстановить пароль
Блоги Сообщество Поиск  

Kafka или Pulsar: Что лучше для потоковой обработки в Java

Запись от Javaican размещена 14.03.2025 в 12:33
Показов 1503 Комментарии 0
Метки apache, java, kafka, pulsar, steaming

Нажмите на изображение для увеличения
Название: d366db0f-3224-446d-9857-809b95b339c2.jpg
Просмотров: 174
Размер:	125.9 Кб
ID:	10390
Среди множества решений для потоковой обработки данных Apache Kafka долгое время удерживала лидирующие позиции, став де-факто стандартом в индустрии. Однако в последние годы всё больше внимания привлекает Apache Pulsar — относительно новый игрок, предлагающий альтернативный подход к построению распределенных очередей сообщений. Оба инструмента широко используются в Java-проектах, предоставляя API для разработки потоковых приложений. При этом они имеют существенные различия в архитектуре, модели хранения и особенностях эксплуатации, что делает выбор между ними нетривиальной задачей.

Архитектура и основные принципы



Apache Kafka и Apache Pulsar, несмотря на схожесть выполняемых функций, имеют принципиально разную архитектуру, что накладывает отпечаток на их характеристики и области применения.

Архитектурная модель Kafka



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

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

Java
1
2
3
4
5
6
7
8
9
10
// Пример конфигурации продюсера в Kafka
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092");
props.put("acks", "all");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
 
Producer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("my-topic", "key", "value"));
producer.close();
Важной особенностью Kafka является использование системы журналирования на диске для хранения сообщений. Kafka оптимизирует операции ввода-вывода, используя последовательные операции записи и активное кэширование на уровне операционной системы. Сообщения хранятся на диске определённый период (настраиваемый через политики хранения), что позволяет клиентам "отматывать" чтение к более ранним сообщениям.

Архитектурная модель Pulsar



В отличие от Kafka, Pulsar использует многоуровневую архитектуру с чётким разделением обработки и хранения:
Брокеры — обрабатывают операции публикации и подписки, но не хранят данные постоянно.
BookKeeper — распределенное хранилище с журналированием, где фактически сохраняются сообщения.
ZooKeeper — используется для координации работы кластера.

Такое разделение обязанностей даёт Pulsar ряд архитектурных преимуществ. Брокеры становятся легко масштабируемыми и относительно простыми компонентами, поскольку вся сложность хранения данных переносится на BookKeeper.

Базовая единица в Pulsar — это топик, который может иметь несколько партиций, называемых в экосистеме Pulsar "разделами" (segments). Топики организованы в пространства имен (namespaces), а те, в свою очередь, в арендаторов (tenants), создавая трёхуровневую иерархию.

Java
1
2
3
4
5
6
7
8
9
10
11
12
// Пример работы с Pulsar в Java
PulsarClient client = PulsarClient.builder()
    .serviceUrl("pulsar://localhost:6650")
    .build();
 
Producer<byte[]> producer = client.newProducer()
    .topic("my-topic")
    .create();
 
producer.send("Hello World".getBytes());
producer.close();
client.close();
Одна из ключевых инноваций Pulsar — модель хранения сообщений. BookKeeper использует распределенную систему журналов (ledgers), где каждый журнал состоит из фрагментов (entries). Каждый фрагмент реплицируется на несколько узлов BookKeeper, обеспечивая высокую доступность и надежность данных.

Управление состоянием и отказоустойчивость



Архитектурные различия между Kafka и Pulsar отражаются в их подходах к управлению состоянием и механизмах отказоустойчивости. В Kafka отказоустойчивость обеспечивается репликацией партиций между брокерами. Каждая партиция имеет один брокер-лидер и несколько брокеров-последователей (followers). При отказе брокера-лидера один из последователей автоматически повышается до статуса лидера. Однако, при полном отказе брокера, все партиции, для которых он был лидером, становятся временно недоступными до завершения процесса переизбрания нового лидера. Это создаёт короткое окно уязвимости. Кроме того, перебалансировка кластера при изменении топологии может вызвать значительную нагрузку на сеть и диски.

Pulsar, благодаря разделению управления и хранения, демонстрирует инной подход к отказоустойчивости. Брокеры являются легковесными и не имеют состояния (stateless), что позволяет быстро восстанавливаться при сбоях. Данные хранятся в BookKeeper, который обеспечивает высокую надежность через репликацию журналов.

Java
1
2
3
4
5
6
// Пример настройки надежного потребителя в Pulsar
Consumer<byte[]> consumer = client.newConsumer()
    .topic("persistent://my-tenant/my-namespace/my-topic")
    .subscriptionName("my-subscription")
    .subscriptionType(SubscriptionType.Failover) // Отказоустойчивая подписка
    .subscribe();
При отказе брокера в Pulsar клиенты автоматически переключаются на другой брокер без потери данных. Это возможно потому, что состояние подписок и данные сообщений находятся не в брокере, а в BookKeeper. Такое разделение обязанностей в архитектуре Pulsar обеспечивает большую гибкость и устойчивость при различных сценариях сбоев, но ценой большей сложности развертывания и обслуживания кластера.

Различия в реализации сегментирования и хранения



Помимо общей архитектуры, Kafka и Pulsar существенно различаются в способах организации сегментирования данных и механизмах хранения сообщений. В Kafka сообщения организованы в партициях в виде аппендикс-лога — последовательной структуры, где новые данные всегда добавляются в конец. Каждая партиция делится на сегменты фиксированного размера для оптимизации доступа и управления политиками хранения. Записи идентифицируются по смещениям (offsets), которые монотонно растут внутри партиции.

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

Java
1
2
3
4
5
6
7
8
9
10
// Пример управления смещениями (offset) в Kafka
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("my-topic"));
 
// Ручное управление смещениями
TopicPartition partition = new TopicPartition("my-topic", 0);
consumer.seek(partition, 1234L); // Перемещение к конкретному смещению
 
// Получение текущей позиции
long position = consumer.position(partition);
Pulsar, в свою очередь, использует совершенно иной механизм сегментирования. Физическое хранение сообщений реализовано через BookKeeper с его концепцией распределенных журналов (ledgers). Топик в Pulsar представлен набором сегментов, каждый из которых хранится как отдельный журнал в BookKeeper. BookKeeper распределяет фрагменты журнала между несколькими физическими узлами (bookie), обеспечивая параллелизм и отказоустойчивость. Каждый фрагмент реплицируется на несколько узлов согласно настраиваемому коэффициенту репликации.

Существенное архитектурное отличие Pulsar — двухуровневая адресация сообщений. В Kafka сообщения идентифицируются по смещению (offset) в партиции, тогда как Pulsar использует комбинацию идентификатора журнала (ledger ID) и номера записи (entry ID). Эта система адресации позволяет Pulsar поддерживать более гибкие сценарии хранения и доступа к сообщениям.

Java
1
2
3
4
5
6
7
8
9
10
11
// Пример продвинутого управления курсорами в Pulsar
Consumer<byte[]> consumer = client.newConsumer()
    .topic("my-topic")
    .subscriptionName("my-subscription")
    .startMessageIdInclusive() // Включая стартовое сообщение
    .subscriptionInitialPosition(SubscriptionInitialPosition.Earliest) // Начать с самого раннего доступного сообщения
    .subscribe();
 
// Ручное подтверждение сообщений
Message<byte[]> msg = consumer.receive();
consumer.acknowledge(msg); // Подтверждение одного сообщения
Ещё одно существенное различие касается способа управления политиками хранения. Kafka предлагает относительно простую модель, где сообщения удаляются на основе их возраста или совокупного размера лога. Pulsar идёт дальше и поддерживает многоуровневое хранение (tiered storage), позволяя автоматически перемещать старые сегменты во внешние хранилища, такие как Amazon S3, Google Cloud Storage или HDFS. Этот механизм многоуровневого хранения даёт Pulsar преимущество при работе с долгоживущими данными — можно хранить практически неограниченное количество исторических данных при относительно низких затратах.

Механизмы партиционирования и маршрутизации сообщений



Оба инструмента используют партиционирование для обеспечения параллелизма, но делают это несколько по-разному. В Kafka партиционирование является фундаментальным аспектом дизайна. Каждый топик обязательно разделен на партиции, и именно число партиций определяет максимальный параллелизм обработки. Распределение сообщений по партициям обычно производится на основе ключа сообщения (при его наличии).

Java
1
2
3
4
5
// Настройка партиционирования по ключу в Kafka
producer.send(new ProducerRecord<>("my-topic", "specific-key", "message-value"));
 
// Использование кастомного партиционера
props.put("partitioner.class", "com.example.CustomPartitioner");
Pulsar, хотя и поддерживает аналогичное партиционирование, предлагает дополнительные возможности:
1. Неразделенные топики: Топики могут быть как партиционированными, так и неразделенными. В последнем случае топик обрабатывается одним брокером, но всё равно пользуется преимуществами распределенного хранения BookKeeper.
2. Ключевые разделяемые топики: Pulsar поддерживает концепцию ключевого разделения (key-shared), где сообщения маршрутизируются на потребителей на основе хеша ключа. Это позволяет нескольким потребителям работать с одним топиком, гарантируя при этом, что сообщения с одинаковым ключом всегда попадут к одному потребителю.

Java
1
2
3
4
5
6
// Пример использования Key-Shared режима в Pulsar
Consumer<byte[]> consumer = client.newConsumer()
    .topic("my-topic")
    .subscriptionName("key-shared-subscription")
    .subscriptionType(SubscriptionType.Key_Shared) // Режим распределения по ключам
    .subscribe();
Такой механизм позволяет реализовать эффективную параллельную обработку, сохраняя при этом порядок сообщений с одинаковым ключом — что особенно важно для сценариев, где требуется строгое упорядочение специфических потоков событий. Гибкость моделей подписки в Pulsar также заслуживает внимания. В отличие от группировки потребителей в Kafka, Pulsar предлагает четыре различных типа подписки:
1. Exclusive (эксклюзивная): только один потребитель может читать топик.
2. Failover (с отказоустойчивостью): несколько потребителей в режиме активный/резервный.
3. Shared (разделяемая): сообщения распределяются между несколькими потребителями в режиме раунд-робин.
4. Key_Shared (с ключевым распределением): сообщения маршрутизируются по ключу.

Эта гибкость позволяет адаптировать поведение брокера под конкретные сценарии использования без необходимости вносить изменения в код приложения.

Что лучше изучить Java Backend разработчику, Angular или React?
Прежде всего я знаю, что это холиварный вопрос, и о нем много информации в интернете. Мне нужно знать, что лучше в моей конкретной ситуации. Я...

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

Какой язык лучше изучать для разработки сайтов Java или PHP?
Скажите, какой язык лучше изучать для разработки сайтов и какой больше востребованный, Java или PHP?

Spring Kafka. Ошибка Connection refused при подключении к брокеру Kafka
Пишу Kafka Broker и Consumer, чтобы ловить сообщения от приложения. При попытке достать сообщения из Consumer вылетает ошибка ...


Производительность и масштабируемость



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

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



Kafka традиционно славится своей высокой пропускной способностью. Её архитектура, ориентированная на последовательные операции ввода-вывода и эффективное использование кэша операционной системы, позволяет достигать впечатляющих результатов при потоковой записи и чтении больших объёмов данных. В тестах производительности Kafka обычно демонстрирует превосходные результаты при сценариях с высокой пропускной способностью и относительно простой моделью потребления, особенно когда число партиций заранее хорошо спланировано. Группа исследователей из Confluent провела серию тестов, показавших, что кластер из пяти брокеров Kafka способен обрабатывать до 4 миллионов сообщений в секунду при оптимальной конфигурации.

Однако, Pulsar с его многоуровневой архитектурой демонстрирует другие характеристики. Разделение обработки и хранения позволяет Pulsar поддерживать более стабильную задержку даже при высоких нагрузках. Это особенно заметно в сценариях, где есть множество топиков с разной интенсивностью трафика.

Java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
// Пример настройки Kafka Producer для максимальной пропускной способности
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092");
props.put("acks", "1"); // Компромисс между производительностью и надёжностью
props.put("batch.size", 131072); // Увеличенный размер пакета
props.put("linger.ms", 5); // Небольшая задержка для накопления сообщений
props.put("buffer.memory", 67108864); // 64MB буфер
props.put("compression.type", "lz4"); // Сжатие для экономии пропускной способности сети
 
// Аналогичные настройки для Pulsar
Producer<byte[]> producer = client.newProducer()
    .topic("my-topic")
    .batchingMaxPublishDelay(5, TimeUnit.MILLISECONDS) // Задержка пакетирования
    .batchingMaxMessages(1000) // Максимальное количество сообщений в пакете
    .compressionType(CompressionType.LZ4) // Сжатие
    .sendTimeout(0, TimeUnit.SECONDS) // Без тайм-аута
    .blockIfQueueFull(true) // Блокировка при заполнении очереди
    .create();
Интересно отметить результаты тестирования, проведенного независимой группой GrabTaxi. В их исследовании, при малых и средних нагрузках (до 100К сообщений в секунду) обе системы показывают сопоставимую производительность. Однако при росте числа топиков и партиций до нескольких тысяч, Pulsar демонстрирует более стабильную задержку, в то время как Kafka начинает испытывать проблемы, связанные с управлением метаданными.

Что касается латентности, здесь картина сложнее. При небольшом количестве топиков Kafka часто выигрывает по минимальной задержке. Однако при увеличении числа топиков и потребителей задержка в Kafka может расти быстрее, чем в Pulsar, благодаря архитектурному разделению в последнем.

Возможности горизонтального масштабирования



Масштабируемость — еще один ключевой аспект, по которому эти платформы имеют существенные различия. Модель масштабирования Kafka тесно связана с партиционированием. Увеличение числа партиций — основной способ повышения параллелизма. Однако с ростом числа партиций растут и накладные расходы на управление ими. В частности:
  • Каждая партиция требует файловых дескрипторов и памяти на брокере.
  • С ростом числа партиций увеличивается время восстановления после сбоев и перебалансировки.
  • Большое количество партиций создаёт дополнительную нагрузку на ZooKeeper (или KRaft в новых версиях).

Java
1
2
3
4
5
6
// Программное создание топика с множеством партиций в Kafka
AdminClient adminClient = AdminClient.create(props);
NewTopic newTopic = new NewTopic("high-throughput-topic", 
                                 100, // Много партиций для параллелизма
                                 (short)3); // Фактор репликации
adminClient.createTopics(Collections.singleton(newTopic));
На практике это означает, что в кластере Kafka обычно рекомендуется ограничивать общее число партиций несколькими десятками тысяч. Превышение этого порога может привести к снижению стабильности системы и усложнению администрирования.
Pulsar, благодаря разделению обработки и хранения, имеет иную модель масштабирования. Брокеры Pulsar не хранят данные локально, а значит, добавление брокеров не требует перемещения данных. Кроме того, BookKeeper позволяет независимо масштабировать уровень хранения.

Java
1
2
3
4
5
6
7
8
9
10
11
// Пример создания партиционированного топика в Pulsar
admin.topics().createPartitionedTopic("my-tenant/my-namespace/my-partitioned-topic", 32);
 
// Настройка многоуровневого хранения для масштабируемости
admin.topicPolicies().setOffloadPolicies("my-tenant/my-namespace/my-topic", 
    OffloadPoliciesImpl.create("aws-s3", "bucket-name", "path/prefix", 
                               100, // Размер блока в МБ
                               Long.MAX_VALUE, // Порог для выгрузки
                               1, // Минимальное время до выгрузки в минутах
                               0, // Время инициализации кластера
                               null, null, null)); // Другие опциональные настройки
Эта архитектурная особенность даёт Pulsar преимущество при динамическом масштабировании: система может подстраиваться под изменяющиеся нагрузки более гибко, без необходимости сложной перебалансировки данных между узлами.

Метрики производительности при высоких нагрузках



Исследователи из Yahoo! (где изначально разрабатывался Pulsar) провели сравнительное тестирование Kafka и Pulsar при экстремальных нагрузках. Результаты показывают интересную картину:
1. При постоянной нагрузке с фиксированным числом топиков (до 1000) и стабильным потоком данных Kafka демонстрирует более высокую пропускную способность, особенно когда большая часть данных читается из кэша.
2. При динамических нагрузках, когда число активных топиков и потребителей меняется, Pulsar показывает лучшую стабильность задержки и более предсказуемое поведение.
3. При большом количестве неактивных топиков (десятки тысяч) Pulsar демонстрирует значительное преимущество, поскольку неактивные топики практически не потребляют ресурсы брокера.

Практическое тестирование, проведённое командой Splunk, выявило еще одну интересную особенность: Kafka демонстрирует более линейный рост производительности при добавлении новых брокеров, но только до определенного предела (обычно около 20-30 брокеров). После этого каждый дополнительный брокер даёт все меньший прирост, в основном из-за растущих накладных расходов на координацию. Pulsar может масштабироваться до сотен брокеров с более стабильным линейным ростом производительности благодаря разделению функций хранения и обработки.
Важно отметить, что обе системы поддерживают кэширование на стороне клиента для повышения производительности чтения часто запрашиваемых данных:

Java
1
2
3
4
5
6
7
8
9
10
// Настройка кэширования на стороне потребителя в Kafka
props.put("fetch.min.bytes", 1024);
props.put("fetch.max.wait.ms", 500);
 
// Настройка кэширования на стороне потребителя в Pulsar
Consumer<byte[]> consumer = client.newConsumer()
    .topic("my-topic")
    .subscriptionName("my-subscription")
    .receiverQueueSize(1000) // Размер буфера сообщений
    .subscribe();

Работа с гетерогенными кластерами



Kafka традиционно лучше работает в однородных кластерах, где все брокеры имеют сходные характеристики. При наличии узлов с разной производительностью могут возникать проблемы с дисбалансом нагрузки. В новых версиях Kafka (2.4+) появилась поддержка "неравномерной" балансировки через rack-awareness и кастомные стратегии распределения, но это всё ещё развивающаяся функциональность.

Java
1
2
// Настройка rack-awareness в Kafka для работы с гетерогенными кластерами
props.put("client.rack", "us-east-1a"); // Указание стойки для клиента
Pulsar, с другой стороны, изначально проектировался с учётом гетерогенности. Его многоуровневая архитектура позволяет размещать брокеры и BookKeeper на разных типах оборудования, оптимизируя использование ресурсов. Например, брокеры могут работать на машинах с большим объёмом ОЗУ, а BookKeeper — на узлах с быстрыми дисками.

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

Интеграция с Java-экосистемой



Java остаётся основной платформой для разработки корпоративных приложений, использующих системы обработки потоковых данных. И Kafka, и Pulsar предоставляют богатые Java API, однако их подходы к интеграции, программным моделям и взаимодействию с популярными фреймворками имеют важные отличия.

Java API и клиентские библиотеки



Kafka предлагает два уровня API для Java-разработчиков:
1. Низкоуровневый Producer/Consumer API — предоставляет полный контроль над всеми аспектами публикации и потребления сообщений.
2. Высокоуровневый Streams API — позволяет строить топологии потоковой обработки с использованием функционального стиля программирования.

Базовый пример использования Java API для Kafka демонстрирует его относительную простоту:

Java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
// Создание потребителя в Kafka
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "test-group");
props.put("key.deserializer", StringDeserializer.class.getName());
props.put("value.deserializer", StringDeserializer.class.getName());
 
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("test-topic"));
 
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
    for (ConsumerRecord<String, String> record : records) {
        // Обработка полученных данных
        System.out.printf("offset = %d, key = %s, value = %s%n", 
                         record.offset(), record.key(), record.value());
    }
}
Pulsar также предоставляет два уровня API:
1. Базовый Producer/Consumer API — основной способ взаимодействия с топиками.
2. Pulsar Functions — легковесная вычислительная модель для создания обработчиков данных.

Пример использования Java API в Pulsar:

Java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
// Создание потребителя в Pulsar
PulsarClient client = PulsarClient.builder()
    .serviceUrl("pulsar://localhost:6650")
    .build();
 
Consumer<String> consumer = client.newConsumer(Schema.STRING)
    .topic("test-topic")
    .subscriptionName("test-subscription")
    .subscriptionType(SubscriptionType.Shared)
    .subscribe();
 
while (true) {
    Message<String> msg = consumer.receive();
    try {
        // Обработка сообщения
        System.out.printf("Got message: %s%n", msg.getValue());
        consumer.acknowledge(msg);
    } catch (Exception e) {
        consumer.negativeAcknowledge(msg);
    }
}
Одно из главных отличий — модель подтверждения получения сообщений. В Kafka потребитель фиксирует свою позицию через offset commits, что происходит автоматически или по явному вызову. В Pulsar используется явная модель acknowledgements, где каждое сообщение должно быть подтверждено, что даёт больше гибкости, но требует явного управления.
Kafka Streams API предлагает богатые возможности для потоковой обработки:

Java
1
2
3
4
5
6
7
8
9
10
// Пример преобразования потока с помощью Kafka Streams
StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> inputStream = builder.stream("input-topic");
KStream<String, String> processedStream = inputStream
    .filter((key, value) -> value.contains("important"))
    .mapValues(value -> value.toUpperCase());
processedStream.to("output-topic");
 
KafkaStreams streams = new KafkaStreams(builder.build(), streamsConfig);
streams.start();
В Pulsar аналогичная функциональность предоставляется через Pulsar Functions и Pulsar IO:

Java
1
2
3
4
5
6
7
// Пример функции преобразования в Pulsar
public class UpperCaseFunction implements Function<String, String> {
    @Override
    public String process(String input, Context context) {
        return input.contains("important") ? input.toUpperCase() : null;
    }
}

Интеграция с популярными Java-фреймворками



Spring Boot — один из самых популярных фреймворков для Java-приложений, и обе системы имеют свои интеграционные модули. Spring для Kafka предоставляет богатый набор абстракций:

Java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// Пример обработчика сообщений в Spring для Kafka
@Service
public class MessageProcessor {
    
    @KafkaListener(topics = "test-topic", groupId = "test-group")
    public void processMessage(String message) {
        // Обработка сообщения
        System.out.println("Получено сообщение: " + message);
    }
    
    @Autowired
    private KafkaTemplate<String, String> kafkaTemplate;
    
    public void sendMessage(String topic, String message) {
        kafkaTemplate.send(topic, message);
    }
}
Для Pulsar существует аналогичный проект Spring for Apache Pulsar:

Java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
// Пример обработчика сообщений в Spring для Pulsar
@Service
public class PulsarMessageProcessor {
    
    @PulsarListener(subscriptionName = "test-subscription", topics = "test-topic")
    public void processMessage(String message) {
        // Обработка сообщения
        System.out.println("Получено сообщение: " + message);
    }
    
    @Autowired
    private PulsarTemplate<String> pulsarTemplate;
    
    public void sendMessage(String topic, String message) {
        pulsarTemplate.send(topic, message);
    }
}
Однако стоит отметить, что Spring для Kafka более зрел и широко используется в продакшн-системах, тогда как Spring для Pulsar — относительно новый проект с меньшей экосистемой.

Особенности работы с Reactive Streams API



Reactive Streams — важная спецификация для работы с асинхронными потоками данных в Java. Kafka и Pulsar предоставляют интеграцию с этой моделью, но делают это по-разному. Kafka поддерживает Reactive Streams через проект Reactor Kafka:

Java
1
2
3
4
5
6
7
8
9
10
11
12
// Пример использования Reactor Kafka
KafkaSender<String, String> sender = KafkaSender.create(senderOptions);
 
Flux<SenderRecord<String, String, String>> recordFlux = 
    Flux.range(1, 100)
        .map(i -> SenderRecord.create(
            new ProducerRecord<>("test-topic", "key-" + i, "value-" + i),
            "correlation-" + i));
 
sender.send(recordFlux)
      .doOnNext(r -> System.out.printf("Message %s sent successfully\n", r.correlationMetadata()))
      .subscribe();
Pulsar имеет непосредственную поддержку Reactive Streams через реактивный API:

Java
1
2
3
4
5
6
7
8
9
10
11
// Пример использования реактивного API Pulsar
Flux<String> messageFlux = Flux.range(1, 100)
    .map(i -> "value-" + i);
 
pulsarReactiveClient.newProducer(Schema.STRING)
    .topic("test-topic")
    .producerName("reactive-producer")
    .create()
    .sendMessages(messageFlux)
    .subscribe(messageId -> 
        System.out.printf("Message sent successfully with ID: %s\n", messageId));
Здесь стоит отметить, что модель сообщений Pulsar с её поддержкой различных типов подписки хорошо сочетается с реактивным программированием. Особенно это проявляется при работе с непрерывными потоками данных и сценариями с множеством производителей и потребителей.

Эффективная сериализация в Java-приложениях



Выбор формата сериализации важен для производительности и гибкости системы. Оба инструмента поддерживают различные схемы сериализации, но имеют отличия в подходах. Kafka полагается на внешний Schema Registry (обычно от Confluent) для управления схемами:

Java
1
2
3
4
5
6
7
8
9
10
// Пример использования Avro с Schema Registry в Kafka
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", StringSerializer.class.getName());
props.put("value.serializer", KafkaAvroSerializer.class.getName());
props.put("schema.registry.url", "http://localhost:8081");
 
Producer<String, User> producer = new KafkaProducer<>(props);
User user = new User("John Doe", 30, "john@example.com");
producer.send(new ProducerRecord<>("users", user.getEmail(), user));
Pulsar предлагает встроенную поддержку схем через Schema Registry, интегрированный прямо в брокеры:

Java
1
2
3
4
5
6
7
8
9
// Пример использования Avro в Pulsar
Schema<User> schema = Schema.AVRO(User.class);
 
Producer<User> producer = client.newProducer(schema)
    .topic("users")
    .create();
    
User user = new User("John Doe", 30, "john@example.com");
producer.send(user);
Встроенная поддержка схем в Pulsar упрощает развертывание и администрирование системы, поскольку не требует отдельного сервиса Schema Registry. Однако решение Confluent для Kafka более зрелое и предоставляет расширенные возможности управления эволюцией схем и совместимостью. Оба инструмента поддерживают популярные форматы сериализации: JSON, Avro, Protobuf и др. При выборе формата следует учитывать требования к производительности, размеру данных и совместимости между различными версиями схем.

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

Практические сценарии применения



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

Случаи, когда предпочтительнее Kafka



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

Java
1
2
3
4
5
6
7
// Пример конфигурации для высоконагруженной системы на Kafka
Properties props = new Properties();
props.put("bootstrap.servers", "broker1:9092,broker2:9092,broker3:9092");
props.put("linger.ms", 5); // Накопление сообщений для пакетной отправки
props.put("batch.size", 16384); // Оптимальный размер батча
props.put("compression.type", "snappy"); // Компромис между скоростью и сжатием
// ...остальные настройки...
Системы, где критична последовательная обработка партиций.
Модель партиций Kafka гарантирует, что партиция обрабатывается ровно одним потребителем в группе, что упрощает разработку приложений, требующих строгой последовательности обработки.

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

Системы с развитой экосистемой коннекторов.
Если вы планируете интегрироваться с множеством различных систем данных, экосистема Kafka Connect с сотнями готовых коннекторов может сэкономить значительное время разработки.

Java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
// Пример использования Kafka Streams для реализации бизнес-логики
StreamsBuilder builder = new StreamsBuilder();
KStream<String, Order> orders = builder.stream("incoming-orders",
                                           Consumed.with(Serdes.String(), orderSerde));
 
KStream<String, Order> validatedOrders = orders
    .filter((key, order) -> order.getAmount() > 0)
    .mapValues(order -> {
        // Обработка заказа
        order.setStatus("VALIDATED");
        return order;
    });
 
validatedOrders.to("validated-orders");
Зрелые проекты с существующими командами DevOps.
Если у вашей организации уже есть опыт работы с Kafka, переход на Pulsar потребует значительных инвестиций в обучение и адаптацию инфраструктуры.

Ситуации, где Pulsar имеет преимущество



Системы с большим количеством топиков.
Если вы работаете с тысячами или десятками тысяч топиков, многие из которых имеют нерегулярный трафик, Pulsar с его архитектурой разделения хранения и обработки будет работать стабильнее. Типичный пример — мультитенантные SaaS-платформы, где каждый клиент может иметь свои потоки данных.

Java
1
2
3
4
5
6
7
8
9
10
11
// Пример работы с топиками в мультитенантной среде Pulsar
// Создание топика для конкретного клиента
admin.topics().createPartitionedTopic(
    "persistent://tenant-" + clientId + "/namespace/events", 
    4  // Количество партиций
);
 
// Продюсер с изоляцией по клиенту
Producer<byte[]> producer = client.newProducer()
    .topic("persistent://tenant-" + clientId + "/namespace/events")
    .create();
Долговременное хранение сообщений.
Благодаря многоуровневому хранению (tiered storage), Pulsar позволяет эффективно хранить терабайты исторических данных, автоматически перемещая старые сегменты в дешёвое облачное хранилище. Это особенно ценно для систем, требующих аудита, нормативного соответствия или ретроспективного анализа (например, платформы кибербезопасности или финансовой аналитики).

Гибридные и мультиоблачные развёртывания.
Расзделение вычислений и хранения в Pulsar упрощает создание распределенных конфигураций, где компоненты системы могут работать в разных дата-центрах или облачных провайдерах.

Системы с разнообразными моделями подписки.
Если ваши сценарии требуют разных паттернов потребления (эксклюзивные потребители, active-standby, распределение по ключам), встроенная поддержка этих моделей в Pulsar сэкономит время на разработке.

Java
1
2
3
4
5
6
7
// Использование key-shared подписки для гарантированной обработки 
// сообщений с одним ключом одним потребителем
Consumer<String> consumer = client.newConsumer(Schema.STRING)
    .topic("orders")
    .subscriptionName("order-processor")
    .subscriptionType(SubscriptionType.Key_Shared)
    .subscribe();
Приложения с требованиями к задержке на 99-м перцентиле.
В системах, где критична стабильная задержка (например, торговые платформы или системы онлайн-принятия решений), Pulsar обычно демонстрирует лучшие результаты P99 благодаря своей архитектуре, меньше подверженной проблемам "тени сборки мусора" или I/O-джиттера.

Проекты с ожидаемым быстрым ростом.
Если вы предвидите резкое увеличение объёмов данных или числа тем/подписчиков, архитектура Pulsar обеспечивает более линейное масштабирование, что может быть решающим фактором для быстрорастущих стартапов или проектов с непредсказуемыми нагрузками.

Гибридные подходы



Интересно отметить, что некоторые организации успешно используют оба инструмента вместе, выбирая наиболее подходящий для конкретной задачи. Например:
  • Kafka для основных высоконагруженных потоков данных.
  • Pulsar для сценариев с множеством топиков и долгим хранением.

Существуют даже проекты, позволяющие объединить их преимущества. Например, можно использовать протокол Kafka с клиентами Kafka для взаимодействия с Pulsar:

Java
1
2
3
4
5
6
7
8
9
// Использование протокола Kafka для взаимодействия с Pulsar 
Properties props = new Properties();
props.put("bootstrap.servers", "pulsar-kafka-proxy:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
 
// Стандартный Kafka Producer API работает с Pulsar через прокси
Producer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("my-topic", "key", "value"));
Этот подход позволяет постепенно мигрировать существующие системы с Kafka на Pulsar или создавать гибридные решения, использующие сильные стороны обеих платформ. В конечном счёте, выбор между Kafka и Pulsar должен основываться на тщательном анализе конкретных потребностей проекта, имеющихся ресурсов и долгосрочных планов развития. Оба инструмента продолжают активно развиваться, добавляя новые функции и улучшая производительность, что делает выбор между ними всё более сложным, но и более интересным.

Критерии выбора между Kafka и Pulsar



Выбор между Apache Kafka и Apache Pulsar не сводится к простому определению "победителя". Обе платформы имеют сильные стороны и области применения, где они превосходят друг друга.

При принятии решения стоит учитывать несколько факторов:
  • Архитектурные требования проекта. Если вы строите систему с ярко выраженным разделением между хранением и вычислениями, многоуровневая архитектура Pulsar может оказаться естественным выбором. Если же вы предпочитаете более монолитный подход с простотой развёртывания — Kafka может быть предпочтительнее.
  • Масштаб и динамика топиков. При работе с десятками тысяч топиков с неравномерной нагрузкой Pulsar демонстрирует лучшую стабильность. Если же у вас относительно небольшое число высоконагруженных потоков — Kafka может предложить более высокую пропускную способность.
  • Зрелость экосистемы и внутренняя экспертиза. У Kafka более развитая экосистема инструментов и большее сообщество, что упрощает поиск решений типовых задач и найм специалистов. Pulsar, хотя и развивается быстро, всё ещё отстаёт в этом аспекте.
  • Стратегическое видение. Если вы смотрите в перспективу 5-10 лет, следует принять во внимание направления развития обеих платформ. Kafka движется в сторону упрощения самоуправления через KRaft, в то время как Pulsar развивается в направлении более глубокой интеграции с облачными средами.
  • Экономические соображения. При ограниченном бюджете Kafka может быть более доступным решением, требующим меньше ресурсов для первоначального развёртывания. Pulsar, с его дополнительным уровнем BookKeeper, требует больше серверов для минимальной конфигурации.

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

В конечном счёте, и Kafka, и Pulsar — серьезные платформы для потоковой обработки данных в Java-экосистеме. Выбор между ними должен опираться на конкретные требования вашего бизнеса и технические аспекты проекта, а не на общие рекомендации или текущие тренды индустрии.

Java & Apache Kafka
Всем доброго времени суток! С кафкой раньше не сталкивался. Задача такая: генератор генерит сообщение, в котором сериализуется объект с полями...

Проблемы с java kafka и zookeeper на windows 10
Здраствуйте. Я сейчас пытаюсь настроить zookeeper и kafka по https://habr.com/ru/post/496182/ вот что я сделал. в файл zoo в...

С чего лучше начать учить Java? С книг или сайтов, или лекций?
Всем привет! Нужна ваша помощь. Помогите пожалуйста новичку в изучении Java! Скажите пожалуйста и (если не сложно) киньте ссылки на книги и...

Что легче/перспективнее для изучения Java или Abap?
Здравствуйте! Для многих вопрос может показаться смешным. Я понимаю, что всё требует большого количества времени для изучения и возможно сам...

Что лучше: SDL2 или GLFW для обработки событий в OpenGl?
Здравствуйте! Хочу побаловаться с OpenGL и теперь выбираю, как обрабатывать события, при помощи какой библиотеки: SDL2 или GLFW. Какая из них...

C++, Java или Python. Что лучше для кроссплатформенного десктопа?
C++ vs Java vs Python ? P.S Для декстопа P.S.S Я новичок.

Что лучше: Java или C#?
Что лучше: Java или C#?

Что лучше: C# или Java?
Что лучше: C# или Java?

Что лучше C# или Java ?
Привет,народ!Я программирую на C# около года,и вот думаю почему не перейти на Java ? Ведь она кроссплафтормена,и вообще open-source. C# стал...

Что лучше, учить команды CMD или BASH или PowerShell или все они важны или лучше язык программирования?
В заголовке имел в виду, что если изучаю распространенный язык программирования, например Python, то команды из этих сред командных оболочек можно не...

Что лучше начать изучать, java или javascript?
Здравствуйте, я новичок в программирований. В школе изучали PascalABC и pascalABC.net. Создавали проекты по этим языкам программирования. Так теперь...

QT что лучше использовать для построения и обработки графиков
Как лучше выводить графики по имеющемуся массиву данных. В последствии планируется сделать аппроксимацию и т.д (если это будет нужно)

Метки apache, java, kafka, pulsar, steaming
Размещено в Без категории
Надоела реклама? Зарегистрируйтесь и она исчезнет полностью.
Всего комментариев 0
Комментарии
 
Новые блоги и статьи
Был там один разговор по поводу свободы в материальном мире.
kumehtar 19.08.2026
Суть: рассматривается живое существо, оказавшееся внутри довольно странной системы (этого мира) и пытающееся обустроить в ней свой кусок пространства. Жизнь действительно предъявляет каждому. . .
Когда логика программы не спасает от человеческих ошибок
Maks 18.08.2026
В последнее время всё чаще и чаще сталкиваюсь с таким явлением, как абсолютная невнимательность (или глупость) пользователей. Проявляется это чаще всего на работе в коллективе. Допустим, человек с. . .
Лето уходит
kumehtar 17.08.2026
Мысли в слух
kumehtar 17.08.2026
Забавно, насколько сейчас стала доступна информация. Например о магии, духовном развитии, медитациях, и других подобных направлениях, ранее зачастую тайных, передаваемых от учителя к ученику. Хотя. . .
Перемещение строк из ТЧ в другой документ с учетом текущего пробега
Maks 17.08.2026
Реализация из решения ниже выполнена на примере нетипового документа "Автозапчасти", с ТЧ "Шины". За основу взят алгоритм отсюда: https:/ / www. cyberforum. ru/ blogs/ 359708/ 10838. html Задача: . . .
Саморегулирующийся социальный контракт для сервера cross-section.
Hrethgir 14.08.2026
С кодом конечно таких глубоких размышлений пока не было, впрочем я уже привык к алгоритмизации. Суть предмета записи: снова в диалоге с нейросетью (я взял пока себе ник для учётки админа - Rector). . . .
Часы электронные
Uhbif79 12.08.2026
Выкладываю программу часов. Программа позволяет: 1. Использовать системное время и дату, 2. Есть возможность вводить время и дату вручную. 3. Реализованы 2 будильника: начало и конец рабочего дня. . . .
Часы с будильником на основе класса QLCDNumber
Uhbif79 12.08.2026
Всем добрый день, выкладываю программу часов с будильником на основе класса QLCDNumber. Здесь я пробовал самостоятельно создавал классы, впервые столкнулся с видимостью переменной одного класса из. . .
КиберФорум - форум программистов, компьютерный форум, программирование
Powered by vBulletin
Copyright ©2000 - 2026, CyberForum.ru