Материал подготовлен командой Simple-Server для администраторов VPS и выделенных серверов. Команды и пути проверяйте на тестовой машине перед production.
Способы и инструменты реализации очереди
У Redis есть несколько встроенных инструментов (функций), которые позволяют реализовать механику очереди сообщений и их последовательную обработку. У каждого способа есть свои плюсы и минусы:
-
«Pub/Sub » (Публикация/Подписка). Один сервис (или несколько сервисов) публикует в отдельную очередь свое сообщение, после чего его могут обработать только те получатели, которые подписаны на эту очередь. В том случае, если на очередь никто не подписан, сообщение будет утеряно.
-
«List » (Список). Очередь типа FIFO (первым пришел — первым ушел). После отправки сервисом сообщения его получит только один подписчик.
«Pub/Sub » (Публикация/Подписка). Один сервис (или несколько сервисов) публикует в отдельную очередь свое сообщение, после чего его могут обработать только те получатели, которые подписаны на эту очередь. В том случае, если на очередь никто не подписан, сообщение будет утеряно.
«List » (Список). Очередь типа FIFO (первым пришел — первым ушел). После отправки сервисом сообщения его получит только один подписчик.
- «Stream » (Поток). Тоже самое, что и Pub/Sub, но с гарантией доставки сообщения. То есть при отсутствии чтения сообщение остается в очереди и ждет обработки.**
**
«Stream » (Поток). Тоже самое, что и Pub/Sub, но с гарантией доставки сообщения. То есть при отсутствии чтения сообщение остается в очереди и ждет обработки.**
**
Создание базы данных Redis в Simple-Server
В этом руководстве будет использоваться облачная база данных Simple-Server.
Соответственно, конфигурация и управления Redis будет выполняться через графический интерфейс панели управления Simple-Server. Поэтому сперва нужно авторизоваться в ней.
После того, как откроется главная страница панели управления, необходимо найти в левом меню пункт «Базы данных » и кликнуть по нему.
Если ранее вы еще не создавали какие-либо базы данных на своем аккаунте, то на открывшейся странице будет кнопка «Создать », на которую нужно кликнуть.
После этого откроется страница конфигурации базы данных.
Самый главный параметр — тип базы данных. В нашем случае выбирается «Redis » версии 6. Остальные настройки можно выставить по своему усмотрению.
На кнопке «Заказать » будет выводится стоимость месячной аренды базы данных. Для тестового проекта можно выбрать самую минимальную конфигурацию, чтобы снизить стоимость.
После завершения конфигурации жмем на кнопку «Заказать ».
Откроется страница управления созданной базой данных. Некоторое время сервер будет неактивен — он будет запускаться и конфигурироваться.
Рассмотрим процесс реализации очереди в Redis пошагово.
В этом руководстве использовался облачный сервер Simple-Server под управлением операционной системы Ubuntu 22.04.
Перед началом конфигурации сервера под приложение Python сперва рекомендуется обновиться систему:
sudo apt update
sudo apt upgradeСперва следует убедиться, установлен ли в системе Python. Это можно сделать, запросив его версию:
В консоли появится примерно такой вывод:
Обратите внимание, что в этом руководстве используется Python версии 3.10.12.
В случае, если в системе отсутствует Python, его следует установить через пакетный менеджер APT:
sudo apt install -y python3Флаг -y необходим для автоматического ответа «yes» на все возникающие во время установки вопросы.
Установка виртуальной среды Python
Для работы приложения потребуется виртуальная среда Python. Поэтому устанавливаем ее:
sudo apt install python3-venv -yСоздание рабочего каталога
Создадим отдельную директория для проекта на Python:
Активация виртуальной среды
Перед началом написания кода нужно создать виртуальную среду Python в рабочем каталоге:
Теперь можно проверить состояние каталога:
Если все прошло успешно, внутри появится папка виртуальной среды:
Далее выполняется активация:
source ./venv/bin/activateУстановка пакетного менеджера Pip
Для работы с Redis потребуется специальный Python-модуль, через который будет выполняться подключение к серверу Redis и непосредственно управление базой данных.
Чтобы получить актуальную версию модуль, сперва необходимо установить пакетный менеджер Pip:
sudo apt install python3-pip -yДля проверки корректности установки можно запросить версию Pip:
В консоли должно появится примерно такое сообщение:
pip 22.0.2 from /usr/lib/python3/dist-packages/pip (python 3.10)В этом руководстве использовался Pip версии 22.0.2.**
**
Теперь можно установить Python-модуль для работы с Redis:
Впоследствии он будет импортироваться в коде приложения Python.
Написание приложения Python
Рассмотрим базовые способы создания очереди на основе Pub/Sub, List и Stream.
Очередь на основе Pub/Sub
В рабочем каталоге создадим файл обработчика, в котором будет читаться очередь сообщений:
import redis
import time
connection = redis.Redis(
host="IP", # указываем IP-адрес сервера Redis
password="PASSWORD", # указываем root-пароль сервера Redis
port=6379, # стандартный порт для подключения к Redis без SSL
db=0,
decode_responses=True # ответы сервера Redis будут автоматически декодироваться в читаемый вид
)
queue = connection.pubsub() # создаем очередь типа Pub/Sub
queue.subscribe("channelFirst", "channelSecond") # подписываемся на указанные каналы
# бесконечный цикл обработки очереди сообщений
while True:
time.sleep(0.01)
msg = queue.get_message() # извлекаем сообщение
if msg: # проверяем сообщение на пустоту
if not isinstance(msg["data"], int): # проверяем, какой тип информации хранится в переменной data (msg имеет тип словаря)
print(msg["data"]) # выводим сообщение в консольСперва выполняется подключение к удаленному серверу Redis, после чего создается очередь типа Pub/Sub.
Обратите внимание, что при подключении к Redis нужно указать адрес удаленного хоста и root-пароль сервера. В этом примере используется подключение без использования SSL, поэтому указывается порт 6379.
Очередь подписывается на два канала — channelFirst и channelSecond.
Внутри бесконечного цикла с некоторым интервалом выполняется проверка наличия новых сообщений. Если в очереди есть сообщение — оно выводится в консоль.
Далее создадим файл отправителя, который будет продуцировать сообщения в очередь:
Его содержимое должно быть таким:
import redis
# аналогичное подключение к удаленному серверу Redis
connection = redis.Redis(
host="IP",
password="PASSWORD",
port=6379,
db=0,
decode_responses=True
)
connection.publish('channelFirst', 'Данное сообщение было отправлено в первый канал') # отправляем сообщение в первый канал
connection.publish('channelSecond', 'Данное сообщение было отправлено во второй канал') # отправляем сообщение во второй каналСперва выполняется подключение к удаленному серверу Redis, аналогичное файлу subscriberPS.py. Далее по открытому соединению отправляется два сообщения — одно в первый канал, другое во второй.
Теперь можно последовательно запустить написанные скрипты и убедиться в работоспособности примера.
Сначала запускаем обработчик сообщений в уже открытом терминале консоли:
Далее открываем второй консольный терминал и выполняем активацию среды:
source ./venv/bin/activateПосле чего запускаем отправителя сообщений:
Соответственно в первом терминале появится следующий вывод:
Данное сообщение было отправлено в первый каналДанное сообщение было отправлено во второй канал
Теперь мы реализуем похожую очередь, но с использованием сущности List.
Сперва создадим файл обработчика:
sudo nano consumerList.pyИ напишем в него следующий код:
import redis
import random
import time
connection = redis.Redis(
host="IP",
password="PASSWORD",
port=6379,
db=0,
decode_responses=True
)
len = connection.llen("listQueue") # получаем размер листа очереди сообщений
# читаем сообщения из листа до тех пор, пока размер листа не станет равен нулю
while connection.llen("listQueue") != 0:
msg = connection.rpop("listQueue") # читаем сообщение, которое является типом данных словаря
if msg:
print(msg) # выводим сообщение в консольОбратите внимание, что в этом примере мы извлекаем и удаляем сообщение «справа», а не «слева». То есть вместо функции lpop используется rpop.
Теперь реализуем файл отправителя:
sudo nano producerList.py
import redis
import random
connection = redis.Redis(
host="IP",
password="PASSWORD",
port=6379,
db=0,
decode_responses=True
)
# отправляем сразу 3 сообщения
for i in range(0,3):
connection.lpush("listQueue", "Сообщение №" + str(random.randint(0, 100))) # добавляем в лист сообщение с уникальным номеромВажно отметить, что сообщения добавляются в лист справа налево.
Например, было отправлено 3 сообщения:
-
Сообщение №1
-
Сообщение №2
-
Сообщение №3
После этого лист очереди будет выглядеть так:
[ Сообщение №3, Сообщение №2, Сообщение №1 ]Соответственно, если в коде обработчика сообщений используется функция rpop, то сообщения будут обрабатываться в порядке их отправки. Если же lpop — в порядке, обратном их отправке.
Тоже самое справедливо и для отправки сообщений с помощью функции rpush вместо функции lpush.
Запустим скрипт отправителя для заполнения листа очереди сообщений:
После этого выполним обработку сообщений:
В консоли должен появится примерно такой вывод — отличаться будут только номера сообщений:
Сообщение №94
Сообщение №96
Сообщение №24Еще один полезный инструмент для реализации очереди — «Потоки».
Для управления потоками существует несколько базовых команд:
XADD. Добавляет новую запись в поток.
XADD. Добавляет новую запись в поток.
XREAD. Читает одну или несколько записей, начиная с заданной позиции и продвигаясь вперед во времени.
XREAD. Читает одну или несколько записей, начиная с заданной позиции и продвигаясь вперед во времени.
XRANGE. Возвращает диапазон записей между двумя предоставленными идентификаторами записей.
XRANGE. Возвращает диапазон записей между двумя предоставленными идентификаторами записей.
XLEN. Возвращает длину потока.
XLEN. Возвращает длину потока.
Создадим файл отправителя сообщений:
sudo nano producerStream.py
import redis
import random
connection = redis.Redis(
host="IP",
password="PASSWORD",
port=6379,
db=0,
decode_responses=True
)
# отправляем сразу 3 сообщения
for i in range(0,3):
connection.xadd("queueStream", { "data":"Сообщение №" + str(random.randint(0, 100))}) # добавляем в очередь сообщение с уникальным номером (тип словаря)
print("Длина очереди: " + str(connection.xlen("queueStream"))) # выводим в консоль размер очереди сообщенийВ этом примере мы отправляем сразу 3 сообщения с уникальным номером в поток через цикл for. После отправки в терминал консоли выводится размер потока.
Теперь реализуем функционал обработчика сообщений:
sudo nano consumerStream.pyКод внутри файла следующий:
import redis
import random
connection = redis.Redis(
host="IP",
password="PASSWORD",
port=6379,
db=0,
decode_responses=True
)
len = connection.xlen("queueStream") # узнаем длину потока
if len > 0:
messages = connection.xread(count=len, streams={"queueStream":0}) # получаем весь список сообщений в потоке
# перебираем список сообщений
for msg in messages:
print(msg) # выводим сообщение в консольСперва мы целиком извлекаем очередь сообщений из потока, после чего последовательно обрабатываем ее в цикле for.
В аналогичном порядке запускаем написанные скрипты. Начинаем с отправителя сообщений:
В консоли должен появиться следующий вывод:
После чего выполняем обработку очереди сообщений:
Соответственно, вывод в консоли будет примерно таким:
['queueStream', [('1711712995031-0', {'data': 'Сообщение №74'}), ('1711712995033-0', {'data': 'Сообщение №54'})]]По этому выводу можно заметить, что у каждого сообщения есть свой уникальный идентификатор, автоматически присвоенный Redis.
Однако у этого примера есть один недостаток — каждый раз мы читаем один и тот же поток целиком.
Давайте усложним код в файле consumerStream.py, чтобы каждый новый запуск скрипта обработчика читал только новые сообщения потока:
import redis
import random
connection = redis.Redis(
host="IP",
password="PASSWORD",
port=6379,
db=0,
decode_responses=True
)
# создаем переменную Redis для хранения ID последнего сообщения (если эта переменная еще не существует)
if connection.get("last") == None:
connection.set("last", 0)
len = connection.xlen("queueStream") # узнаем длину потока
if len > 0:
messages = connection.xread(count=len, block=1000, streams={"queueStream":connection.get("last")}) # в момент чтения передаем ID последнего сообщения в качестве аргумента (либо 0)
print(connection.get("last")) # выводим в консоль ID последнего сообщения (либо 0)
# перебираем список новых сообщений
for msg in messages:
print(msg) # выводим сообщение в консоль
connection.set("last", msg[-1][-1][0]) # устанавливаем в качестве значения переменной ID последнего прочитанного сообщенияТеперь каждый новый запрос к Redis будет выводить в консоль только свежие сообщения.
Разверните Redis в облаке в пару кликов
В этом руководстве было продемонстрировано несколько базовых способов создания очереди в Redis: Pub/Sub, List и Stream.
Показанные примеры — лишь минимальные реализации, выполняющие логику очереди сообщений. Реальный проект потребует усложнения этой логики, чтобы удовлетворять критериям разработчика и решать поставленные задачи.
Например, функционал очереди сообщений может быть обернут в классы и объекты, либо выполнен в виде отдельной внутренней библиотеки проекта.
Поэтому в каждом конкретном проекте потребуется дальнейшее уникальное развитие этой реализации, чтобы решать задачи
Более подробно ознакомится с командами Redis, которые предназначены для работы с разными инструментами очереди сообщений, можно в специальных разделах официальной документации Redis:
Нужен сервер для практики? Арендуйте VPS/VDS в России — root-доступ, NVMe, DDoS-защита и поддержка 24/7.