Есть ли готовые решения для очереди с десятком тысяч заранее не заданных партиций?
Проблема: нужно синхронизировать мои данные со сторонними сервисами.
Дано:
- Мои данные ClientId, CustomerId, CustomerEmail, CustomerPoints
- Четыре (в перспективе до 20) сервиса с которыми нужно синхронизировать мои данные
- Примерно 15 000 клиентов (ClientId), у клиентов от 1000 до 500 000 customer-ов
- У каждого сервиса есть rate-limit. Доступ к каждому сервису для каждого клиента осуществляется по приватному ключу, т.е. rate-limit индивидуальный для связки ClientId-ServiceName
Текущее решение:
- Есть kafka со 100 партициями, партиция это ClientId % 100 = (0 - 99)
- При обновлении данных customer-а в kafka topic пишется сообщение о том чей он клиент (ClientId) и какой у него id
- Проверяется с какими сервисами у клиента настроена синхронизация
- Запускается параллельно синхронизация со всеми активными сервисами
- При неудаче синхронизация повторяется до тех пор пока не выполнится (rate-limit и отказы сервисов учитываются именно сдесь)
Какие сейчас есть проблемы:
- При подключении нового сервиса клиентом, мне необходимо синхронизировать всех его customer-ов с этим сервисом, а их у него может быть до 500 000. Партиция забьётся и те клиенты которым не повезло оказаться с ним в одной партиции будут очень долго синхронизироваться учитывая все rate-limit-ы.
- При достижении rate-limit-а у определённого клиента на каком либо из сервисов, страдает синхронизация для всех остальных сервисов и для всех клиентов сообщения от которых тоже находятся в этой партиции.
- Если сторонний сервис выходит из строя на 20 минут, то каждая партиция в которой появится сообщение для клиента данные которого необходимо синхронизировать с этим сервисом "зависнет" на 20 минут.
Какое решение для очереди сообщений мне нужно реализовать (найти готовое):
- Неограниченное количество партиций
- Партиция должна быть строкой или составным значением (ClientId, ServiceName)
- Каждая партиция должна обрабатываться конкретным consumer-ом
- Когда в партиции заканчиваются данные она пропадает и consumer переключается на обработку другой
- Возможность повторной обработки сообщений в случае ошибок
- Проверка есть ли в очереди сообщение по его ключу и значению (это не обязательно, но желательно что-бы избежать дубликатов)
- Проверка размера необработанных сообщений в определённой партиции (это не обязательно, но желательно)
Сразу отброшу kafka т.к. нет возможности заранее узнать все партиции и держать огромное их количество.
Рассматривал Apache Pulsar, но там количество партиций задаётся глобально в сетингах что тоже не подходит, или я невнимательно изучал этот вопрос, если так то поправьте.
RabbitMQ совсем как мне кажется про другое.
В голову лезут мысли про свой собственный велосипед поверх:
- redis streams
- redis hset+list
- mongodb (но тут точно возникнут проблемы с выполнением на нескольких серверах)
- sqs + в памяти разбивать по очередям для каждого ClientId+ServiceName (плохо параллелится и возможны проблемы по памяти)
Может у кого-то были схожие задачи и есть готовые реализации подобных очередей?
Дополнительно:
Что значит
Когда в партиции заканчиваются данные она пропадает и consumer переключается на обработку другой
?
Как узнать что в стриминговой системе вдруг закончились данные?
Может быть ты имеешь в виду диспетчеризацию внутри процесс-консюмера?
Например я - процесс потребитель событий. И с одной стороны я могу подключаться к 10 топикам Kafka.
Сливать все 10 в один внутренний для обработки. Те события которые удалось обработать - коммитить
в кафку (да там есть режим фиксации транзакций). А те которые еще не готовы к обработке
я буду вращать по кругу во внутреннем буфере. Какое-то время. 5-10 минут. Выбери сам.
Потом сказать Кафке rollback. Дескыть пока не судьба. Положу обратно на полочку.
Kafka бы мне идеально подошла если бы я мог создать порядка 100 000 партиций в топике, а в перспективе мне могут понадобиться ещё несколько таких топиков, но вроде у kafka есть глобальный лимит на количество партиций.
Если бы у меня было 30 клиентов и 3 сервиса я бы создал 90 партиций и указал в параметре partitionsConsumedConcurrently например 10 , то сообщения из каждой партиции обрабатывались бы параллельно и с примерно одинаковой скоростью не "мешая" другим партициям. Если бы я в одну партицию я записал миллион сообщений, это бы не привело к задержкам обработки сообщений из других партиций.
Но к сожалению я не могу создать 100 000 партиций в kafka, или могу?
Я нашёл ещё два варианта которые могли бы решить проблему:
1. Apache Pulsar, но пока что тоже не уверен что он мне подойдёт.
2. AWS SQS FIFO с указанием MessageGroupId - это точно решает мою проблему, но это будет стоить дороже
создашь 100 тыщ партиций то тебе надо некоторое количество дисковых хранилищ и серверов
сразу. Беря во внимание что обычно бизнес нагрузка эволюционирует плавно то скорее всего
тебе эти мощности сразу будут не нужны. Но будучи созданными они будут зря тратить деньги
за uptime.
Ты можешь для начала создать 16 партиций а свои топики отобразить на на партиции по ключу (Key).
Кафка поддерживает ключ в каждом месседже поэтому создай себе правильную хеш функцию
которая отображает ключи на партиции. Это в кафке заложено. Вот. Если мощности не хватит - тогда
сделай 32 партиции. Потом 64 и так далее. Вот это и будет правильный подход.
Почему один ключ должен "мешать" другому я не понимаю. Если такое было то у тебя должен
быть реальный кейс такой ситуации из практики. Но я в этом сомненваюсь.
Было бы проще переварить одну картинку (диаграмму) вместо тыщи слов.
При отправке сообщения в очередь помимо самого сообщения отправляется еще и группа. При получении сообщений гарантируется что все пришедшие за один запрос сообщения будут из одной группы и их после получения можно обработать последовательно.
При чём если сообщения полученые первым запросом еще не обработались, то при повторном получении сообщений гарантируется что сообщения придут из другой группы.
Но есть одно неприятное ограничение. При получении сообщений проверяются только первые 20 000 и если среди них нет ни одного из не занятой обработкой группы, то ответ придет пустой.
Тут про ограничение https://docs.aws.amazon.com/AWSSimpleQueueService/...
For FIFO queues, there can be a maximum of 20,000 in flight messages (received from a queue by a consumer, but not yet deleted from the queue). If you reach this quota, Amazon SQS returns no error messages. A FIFO queue looks through the first 20k messages to determine available message groups.
И это неприятно, т.к. необходимо запустить синхронизацию всем клиентам, а сейчас только для 10 запустил и уже появилась проблема связанная с тем что приходят сообщения только из одной группы
Из этих 2.7 миллионов примерно полтора миллиона от одного клиента из-за этого обрабатавается одна группа, т.к сообщения других клиентов видимо лежат дальше чем первые 20 000
Опишите проблему, и специалист поможет с настройкой, исправлением ошибки или доработкой сайта. Подберём понятный план работ без лишней переписки.
Пока нет других ответов. Будьте первым, кто поможет автору.
Ответить на вопрос

Для обработки очереди с десятками тысяч заранее не заданных партиций, можно воспользоваться различными готовыми решениями. Одним из наиболее эффективных способов решения данной проблемы является использование структуры данных "Priority Queue" или "Приоритетная очередь".
Приоритетная очередь позволяет хранить элементы с определенным приоритетом, где элемент с наивысшим приоритетом будет извлекаться первым. Это идеальное решение для обработки очереди с различными приоритетами и неограниченным числом партиций.
Для реализации приоритетной очереди в PHP можно воспользоваться стандартной библиотекой SPL (Standard PHP Library), которая предоставляет класс SplPriorityQueue. Этот класс представляет собой приоритетную очередь, в которой элементы хранятся в порядке их приоритета.
Пример использования SplPriorityQueue для обработки очереди с десятками тысяч партиций:
$queue = new SplPriorityQueue(); // Добавление элементов в очередь с указанием их приоритета $queue->insert('Partition 1', 1); $queue->insert('Partition 2', 2); $queue->insert('Partition 3', 3); // Добавление всех необходимых партиций // Извлечение элементов из очереди в порядке их приоритета $queue->top(); // Возвращает элемент с наивысшим приоритетом (Partition 3) $queue->extract(); // Извлекает и удаляет элемент с наивысшим приоритетом (Partition 3)
Таким образом, использование приоритетной очереди позволит эффективно обрабатывать очередь с десятками тысяч заранее не заданных партиций, сохраняя порядок их приоритета и обеспечивая быстрый доступ к элементам.