Проект распределенной очереди с задержкой V1~2
Введение
задний план
Когда мы работаем над системой, во многих случаях мы имеем дело с задачами в реальном времени, обрабатывая запросы по мере их поступления, а затем немедленно предоставляя обратную связь пользователям. Но иногда встречаются задачи не в реальном времени, например, сделать важные объявления в определенный момент времени. Или нужно X минут/Y часов после того, как пользователь что-то сделал, например:
«PM: Нам нужно напомнить этому пользователю, чтобы отправить ему вознаграждение через 10 минут после начала звонка»
Для конкретных действий, таких как уведомления, купоны и т. д. Как правило, решения, с которыми я сталкивался, будут поддерживать бэкенд в относительно небольшом сервисе, но с увеличением количества таких бэкэндов и серверов этот метод в значительной степени связан с собственным бизнесом, поэтому в настоящее время требуется задержка службы очереди.
Глоссарий
очередь theme_list: каждый входящий запрос задержки должен ссылаться на kafka для другой темы задержки, и логически разделить очередь из каждого дела и обрабатывать его отдельно;
очередь тема_информация: каждая тема очереди существует в новой очереди, и каждый раз информация о теме сканируется для определения количества новых сопрограмм службы управления созданием и уничтожением темы;
смещение: ход текущего потребления;
new_offset: прогресс нового потребления, готово изменить смещение;
topic_offset_lock: Распределенная блокировка.
2. Цели дизайна
Список функций
1. Интерфейс для добавления информации о задержке основан на http-вызовах.
2. Благодаря функции очереди хранения он может сохранять данные об использовании очереди за последние 3 дня.
3. Обеспечьте функцию потребления
4. Задержка уведомления
Представление
Предполагаемый объем вызовов интерфейса: 3500 одноклассовых задач в секунду, 1300 одноклассовых задач за несколько секунд
Результаты испытаний под давлением:
Простое испытание под давлением
wrk записывает qps: 259,3 с, записывает 9000 записей, один поток, без параллелизма
Производительность/точность триггера: 1000 в секунду, без расширения на тестовой машине. При 3000 в секунду иногда бывают задержки на 1-2 секунды. Влияет на память и процессор.
3. Дизайн системы
Интерактивный процесс
Временная диаграмма
Этот дизайн основан на вызовах интерфейса http. Когда сообщение добавляется в существующую очередь темы, сообщение будет добавлено в конец соответствующей очереди темы для хранения. При добавлении в соответствующую очередь темы, которая не существует, сначала будет создана новая очередь темы.Когда срабатывает триггер или распределенная блокировка, экземпляр, который захватывает блокировку, сначала получает смещение соответствующей очереди, устанавливает новое смещение, затем блокировка может быть снята для других экземпляров для борьбы так как в голове очереди извлекается определенное количество элементов, затем получается сегмент смещения, экземпляр отправляется в хранилище за подробной информацией и обрабатывает ее в сопрограмме, а основная сопрограмма ждет следующего триггера. Затем добавьте сопрограммы для мониторинга триггеров.
Модульное деление
1. Модуль хранения очереди
1. Задерживаемый модуль delay.base в основном отвечает за получение запросов на запись, запись информации об очереди в хранилище, не отвечает за внутреннюю логику и вызов модуля хранилища.
2. Бэкенд-модуль. Модуль delay.backend под задержкой отвечает за синхронизированное сканирование соответствующей очереди тем, вызов модуля хранения, в основном отвечающего за доступ к модулю хранения чтения, и вызов модуля обратного вызова.
1. Отсканируйте тему, чтобы добавить программу
2. Сканировать информацию о потреблении тематического_списка
3. Просканируйте список тем и закройте программу, если она не используется в течение определенного периода времени.
3. Модуль обратного звонка в основном отвечает за отправку поступивших в момент данных и уведомление соответствующей службы
3. Модуль хранения
1. Модуль распределенной блокировки, система развертывается на нескольких машинах, чтобы обеспечить уникальность каждого потребления, и блокируется сегмент смещения каждого потребления темы.Смещение сегмента new_offset является эксклюзивным для одной машины.
2. Список управления темами, управление количеством тем для управления количеством сопрограмм.
3. тема_список, очередь сообщений
4. Topic_info, объекту сообщения, может потребоваться передать некоторую информацию в обратном вызове для унифицированной обработки.
4. Модуль генерации уникальных номеров.
5. Дизайн кэша
В настоящее время используется режим полного кеша
ключевой дизайн:
ключ списка управления темами: XX:DELAY_TOPIC_LIST type:list
Ключ темы_списка: XX:DELAY_SIMPLE_TOPIC_TASK-%s (по ключу темы) тип:zset
Ключ темы_информации: XX:DELAY_REALL_TOPIC_TASK-%s (по ключу темы) тип: хэш
Ключ темы_смещения: XX:DELAY_TOPIC_OFFSET-%s (по ключу темы) тип:строка
topic_lock ключ: xx: delay_topic_reload_lock-% s (по теме ключ) Тип: строка
6. Дизайн интерфейса
delay.task.addv1 (добавление очереди задержки v1)
Пример запроса
|
вернуться к примеру
|
Метод обратного вызова Pull возвращает (v2 больше не поддерживает)
Пример запроса
|
вернуться к примеру
|
delay.task.addv2 (добавление очереди задержки v2)
Пример запроса
|
Пример
|
вернуться к примеру
|
Seven, дизайн MQ (v2 больше не поддерживается)
Возвращает о методе потребления kafka:
|
8. Другие конструкции
Уникальный дизайн номера
Вызовите модуль хранилища и используйте логику автоинкремента Redis для создания уникального номера. Конкретная логика выглядит следующим образом:
func (c *CacheManager) OperGenTaskid() (uint64, error) {
now := time.Now().Unix()
key := c.getDelayTaskIdKey()
reply, err := c.DelayRds.Do("INCR", key)
if err != nil {
log.Errorf("genTaskid INCR key:%s, error:%s", key, err)
return 0, err
}
version := reply.(int64)
if version == 1 {
//默认认为1秒能创建100个任务
c.DelayRds.Expire(key, time.Duration(100)*time.Second)
}
incrNum := version % 10000
taskId := (uint64(now)*10000 + uint64(incrNum))
log.Debugf("genTaskid INCR key:%s, taskId:%d", key, taskId)
return taskId, nil
}
Распределенная конструкция замка
func (c *CacheManager) SetDelayTopicLock(ctx context.Context, topic string) (bool, error) {
key := c.getDelayTopicReloadLockKey(topic)
reply, err := c.DelayRds.Do("SET", key, "lock", "NX", "EX", 2)
if err != nil {
log.Errorf("SetDelayTopicLock SETNX key:%s, cal:%v, error:%s", key, "lock", err)
return false, err
}
if reply == nil {
return false, nil
}
log.Debugf("SetDelayTopicLock SETNXEX topic:%s lock:%d", topic, false)
return true, nil
}
9. Вопросы дизайна
прочность
стратегия автоматического выключателя:
В этой версии дизайна есть много недостатков.Когда redis недоступен, большое количество невыполненных запросов будет оказывать давление на машины или экземпляры, что приведет к недоступности других сервисов, поэтому принимается стратегия понижения (стратегия понижения также недостаточна ); добавить его при запросе redis Retry, когда количество повторных попыток больше, чем количество тревог, будет записана атомарная операция atomic.StoreInt32(&stopFlag,1), где stopFlag — глобальная переменная, после atomic.LoadInt32(&stopFlag ), значение stopFlag равно 1 временно Не запрашивать redis, одновременно записывать текущее время, добавить таймер, фьюз разделен на три уровня, вкл, выкл, полуоткрыт, по истечении таймера stopFlag= 2, второй таймер будет отсчитывать время полуоткрытого состояния, и есть вероятность повторного доступа, когда количество успешных попыток достигнет порога stopFlag=0, в противном случае stopFlag=1 продолжит отсчет времени
недостаточный
1. Время звонка
Обычно golang записывает временную задачу, выполняемую циклом, тремя способами:
1. время. Метод сна:
for {
time.Sleep(time.Second)
fmt.Println("test")
}
2. Функция time.Tick:
t1:=time.Tick(3*time.Second)
for {
select {
case <-t1:
fmt.Println("test")
}
}
3. Среди задач синхронизации тиков вы также можете использовать функцию time.Ticker, чтобы сначала получить структуру тикера, а затем заблокировать информацию мониторинга.Таким образом, вы можете вручную остановить задачу синхронизации и сократить потери памяти при остановке задачи.
t:=time.NewTicker(time.Second)
for {
select {
case <-t.C:
fmt.Println("test")
t.Stop()
}
}
В начале я думал, что сон обрабатывается отдельно и останавливал сопрограмму напрямую, поэтому первая версия тоже использовала сон, но после сбора данных я обнаружил, что в этих методах были созданы таймеры, а также добавлена сопрограмма обработки тайминга задачи. Фактически, таймеры, сгенерированные этими двумя функциями, помещаются в одну и ту же кучу таймеров (колесо времени Golang), и все они ожидают обработки в сопрограмме обработки задачи синхронизации. Структуры таймера, используемые функциями Tick, Sleep и time.After, будут обрабатываться единообразно в одной и той же сопрограмме, поэтому кажется, что нет никакой разницы между использованием Tick и Sleep. На самом деле разница есть.Эта статья не для обсуждения плюсов и минусов time.sleep и time.tick для временного выполнения задач в golang, о чем будет рассказано в последующих статьях. Более гибко использовать сопрограмму блокировки канала для выполнения задач по времени.Вы можете установить время ожидания и метод выполнения по умолчанию в сочетании с выбором, и вы можете настроить таймер на автоматическое закрытие.Поэтому рекомендуется использовать time.Tick to выполнять задания на время.
2. Проблема с модулем памяти
В настоящее время это полный кеш и никакая БД не задействована.Во-первых, проблемой является высокая доступность redis (codis), а также проблематично принять решение «ничего не делать» после прерывателя цепи.Поэтому, для перспектив на будущее первое:
1. Структура данных одной машины использует несколько колес времени. Чтобы уменьшить расстояние между данными, процесс загрузки данных асинхронно загружается в машину, чтобы уменьшить потери времени, вызванные сетевым вводом-выводом. В то же время это также снижает зависимость от Redis.
2. Внедрите ZooKeeper или добавьте резервную копию кластера, лидер. Убедитесь, что по крайней мере две машины в кластере загружают данные темы, а лидер может координировать потребление для обеспечения высокой доступности.