
Коннектор Redis для redb.Route: все структуры данных как шаги маршрута, Pub/Sub, Streams с группами потребителей, список как очередь и Claim Check на Redis.
В интеграционном проекте Redis редко бывает чем‑то одним. Это кэш ответов внешних систем, счётчики и блокировки, шина событий через Pub/Sub, журнал сообщений через Streams, очередь задач на списке. Обычно каждая из этих ролей живёт в своём куске кода со своим подключением, своими повторами и своим логированием.
В redb.Route все они собраны в одном коннекторе redb.Route.Redis. Любая операция Redis становится шагом маршрута, а подписка, стрим или список становятся его входом. Подключение, телеметрия и плавная остановка при этом общие с остальными коннекторами. Под капотом StackExchange.Redis 2.11.8.
Потребители Redis описаны здесь начиная с версии 4.1.0: подписка по шаблону, стрим без группы со своей позицией, повторная обработка зависших записей стрима и надёжная очередь на списке появились в ней.
Подключение пакета и адрес
dotnet add package redb.Route.Redis
services.AddRedbRoute(route => { route.Services.AddRedbRouteRedis(); route.AddRouteBuilder<OrderRoutes>(); });
Адрес эндпоинта устроен так: redis:ОПЕРАЦИЯ:ресурс?параметры. Первый сегмент пути это операция Redis (регистр не важен), всё после первого двоеточия это ключ, канал или имя стрима. Двоеточия внутри ключа сохраняются, поэтому redis:SET:session:42?ttl=300 записывает ключ session:42 со сроком жизни 300 секунд.
То же самое собирает fluent‑построитель, и дальше в статье он используется везде:
.To(Redis.Set("session:42").Ttl(300)) // redis:SET:session:42?ttl=300
Для ключей, Pub/Sub, стримов и списков в Redis есть готовые фабрики (Set, Get, Incr, Publish, Subscribe, XAdd, XRead, LPush и остальные). Любую другую операцию задаёт Redis.Command(операция, ключ): например, Redis.Command("HSET", "customer:42").Field("email").
Все структуры данных как шаги маршрута
Продюсер поддерживает операции над всеми основными структурами Redis:
Структура |
Операции |
|---|---|
Строки и ключи |
|
Списки |
|
Хэши |
|
Множества |
|
Sorted sets |
|
Гео |
|
HyperLogLog |
|
Битовые карты |
|
Сообщения |
|
Любая команда |
|
Откуда операция берёт данные, определено одинаково для всех:
значение или элемент это тело сообщения. В списках, хэшах, множествах, sorted sets и стримах тело
byte[]пишется как есть, без перевода в строку;скалярные параметры задаются в адресе:
fieldдля хэшей,score,minScoreиmaxScoreдля sorted sets,startиstopдля диапазонов,longitude,latitude,member1,member2иgeoUnitдля гео,offsetиbitдля битовых карт;наборы значений передаются заголовками: словарь полей для
HMSET(redbRedis.HashFields), имена полей дляHMGET(redbRedis.FieldNames), исходные ключи дляPFMERGE(redbRedis.SourceKeys), центр и радиус дляGEORADIUS, аргументы дляCOMMAND.
Ключ, поле, канал и имя стрима могут содержать шаблоны ${...}: они вычисляются на каждом сообщении.
Результат операции становится телом для следующего шага. GET отдаёт строку или null, INCR новое значение счётчика, LRANGE и SMEMBERS массив строк, HGETALL словарь, GEORADIUS имена найденных точек:
From("direct://order-created") .To(Redis.Incr("stats:orders:today")) .Log("заказов сегодня: ${body}");
SETNX записывает значение, только если ключа ещё нет, и возвращает "OK" или null. Вместе с Ttl(...) это готовая основа для блокировки на время обработки или для отсечения повторов.
Pub/Sub: события всем, кто сейчас слушает
From("direct://publish-event") .To(Redis.Publish("orders.events")); From(Redis.Subscribe("orders.events")) .To("direct://notify");
После публикации число получателей лежит в теле и в заголовке redbRedis.Publish.Recipients. У подписчика тело это текст сообщения, а в заголовках канал (redbRedis.Channel), тип (PubSub) и время получения.
Подписка по шаблону делается через PSUBSCRIBE: From(Redis.PSubscribe("orders.*")) получает сообщения всех каналов, подходящих под шаблон, а в заголовке redbRedis.Channel лежит канал, в который сообщение пришло. Параметр usePattern превращает в подписку по шаблону и обычный SUBSCRIBE.
Pub/Sub в Redis доставляет сообщение только тем, кто подключён в момент публикации, и нигде его не хранит. Это правильный выбор для оповещений и сброса кэшей. Когда нужна доставка с подтверждением, берите Streams.
Streams: журнал сообщений с группами потребителей
Запись в стрим:
From("direct://order-created") .To(Redis.XAdd("orders").StreamMaxLength(100_000));
Тело попадает в поле data вместе с полем timestamp в миллисекундах. Если нужен свой набор полей, передайте словарь в заголовке redbRedis.StreamFields. StreamMaxLength обрезает стрим по длине, по умолчанию приблизительно (MAXLEN ~); точную обрезку включает StreamApproximate(false). Идентификатор новой записи возвращается в теле и в заголовке redbRedis.Stream.MessageId.
Чтение группой потребителей:
From(Redis.XRead("orders") .ConsumerGroup("billing") .ConsumerName("node-1") .StreamReadCount(50) .StreamClaimMinIdle(60_000)) .To("direct://bill");
Что происходит под капотом:
группа создаётся сама при старте, вместе со стримом, если его ещё нет. По умолчанию она начинает с конца стрима и получает записи, добавленные после её создания;
StreamStartPosition("0")начинает с первой записи. Если группа уже существует, консюмер подключается к ней;записи читаются пачками по
streamReadCount. Когда новых записей нет, консюмер ждётstreamBlockTimeMs(по умолчанию 1000 мс) и опрашивает снова;тело это словарь полей записи, каждое поле дополнительно лежит в заголовке
redbRedis.Stream.<поле>, а идентификатор записи вredbRedis.MessageId;имя потребителя по умолчанию равно имени машины.
Несколько узлов с одной группой и разными именами потребителей делят стрим между собой: каждую запись получает один из них. Так один стрим обрабатывается кластером без дублей.
Подтверждение и повторная обработка
Консюмер группы подтверждает запись (XACK) после того, как вся единица работы завершилась успешно, включая транзакцию маршрута. Запись, обработка которой упала, не подтверждается и остаётся в списке ожидающих группы. Так получается доставка «хотя бы один раз».
Вернуть такие записи в работу помогает StreamClaimMinIdle(мс). На каждом опросе консюмер забирает командой XAUTOCLAIM ожидающие записи группы, простоявшие не меньше заданного времени, и обрабатывает их заново: и свои упавшие, и оставленные потребителем с упавшего узла. Время простоя стоит брать больше самого долгого маршрута, иначе запись заберут, пока её ещё обрабатывают.
Для потоков, где потеря записи допустима, есть StreamNoAck(): чтение идёт с NOACK, запись считается доставленной в момент чтения и в список ожидающих не попадает (не более одного раза). С StreamClaimMinIdle этот режим не сочетается, эндпоинт с обеими опциями не создаётся.
Стрим без группы
Без группы каждый потребитель читает стрим целиком и сам ведёт позицию: после каждой прочитанной записи она сдвигается дальше. Так один журнал раздаётся нескольким независимым получателям или перечитывается с начала:
From(Redis.XRead("orders").StreamStartPosition("0")) .To("direct://rebuild-projection");
Без StreamStartPosition потребитель начинает с записей, добавленных после его запуска, "0" читает с первой записи, а идентификатор записи продолжает после неё. Позиция > означает «ещё не выданное группе» и имеет смысл только для группы: без группы консюмер с ней не создаётся. Упавшую запись без группы ничто не удерживает, и чтение идёт дальше.
Список как очередь задач
From("direct://enqueue-job") .To(Redis.LPush("jobs")); From(Redis.Command("BRPOP", "jobs") .ProcessingList("jobs:processing") .PollDelay(500)) .To("direct://run-job");
LPUSH в голову и чтение с хвоста дают очередь в порядке поступления. Консюмер забирает элементы по одному опросом, а когда список пуст, ждёт pollDelayMs. BLPOP и BRPOP здесь имена операций, а не блокирующие команды: блокирующее чтение заняло бы общее соединение, через которое работает весь процесс.
С ProcessingList очередь надёжная. Элемент переносится командой LMOVE в список «в работе» и удаляется оттуда после успешного маршрута. Упавший элемент атомарно возвращается в очередь, туда, откуда его взяли, и обрабатывается снова, а то, что осталось в списке «в работе» после прошлого запуска, при старте возвращается в очередь. На один список «в работе» приходится один консюмер. LMOVE есть в Redis начиная с версии 6.2.
Без ProcessingList элемент снимается со списка до обработки и при сбое маршрута теряется. Такой режим подходит для задач, которые не жалко потерять, например для прогрева кэша.
Подключение и секреты
Самый короткий путь это параметры в адресе: connectionString (по умолчанию localhost:6379), database и password. Пароль помечен как секрет и в логах скрыт.
Для продакшена удобнее именованная фабрика подключений в реестре контекста, тогда секреты не попадают в адрес маршрута:
context.AddToRegistry("prod", new RedisConnectionFactory { ConnectionString = "redis-1:6379,redis-2:6379", User = "route", Password = secrets.RedisPassword, Ssl = true, SslProtocols = "Tls12, Tls13", }); From("direct://rates") .To(Redis.Set("rates:usd").ConnectionFactory("prod").Ttl(300));
Фабрика знает ACL‑пользователя Redis 6+, TLS с выбором протоколов, Sentinel (ServiceName), таймауты подключения и операций, keep‑alive, политику переподключения (экспоненциальную или линейную) и префикс каналов. Ошибки конфигурации не маскируются: имя фабрики, которого нет в реестре, и опечатка в SslProtocols дают явную ошибку при подключении вместо тихого перехода на параметры по умолчанию. То же с параметрами адреса: опцию, которой эндпоинт не знает, или значение не того типа он не принимает и называет такой параметр по имени, а не теряет его молча.
Каждый эндпоинт держит одно подключение StackExchange.Redis и переподключается сам. Потеря и восстановление связи пишутся в лог.
Claim Check на Redis
Шаблон Claim Check из каталога Enterprise Integration Patterns убирает тяжёлое тело из сообщения на время маршрута и возвращает его, когда оно снова понадобится. Для него есть хранилище на Redis, RedisClaimCheckRepository:
var claims = new RedisClaimCheckRepository( new RedisConnectionFactory { ConnectionString = "redis.internal:6379" }, defaultTtl: TimeSpan.FromHours(1)); From("direct://inbound-invoice") .ClaimCheck(claims, ClaimCheckOperation.Push) // тело в Redis, в сообщении ключ .To("direct://route-by-headers") .ClaimCheck(claims, ClaimCheckOperation.Pop) // тело вернулось, запись удалена .To("direct://archive");
Хранилище пишет данные как есть, в двоичном виде, и ставит на запись родной срок жизни Redis. Ключи получают префикс redb:claimcheck:, его можно заменить. Чтение с удалением атомарно: оно идёт одним Lua‑скриптом, поэтому работает и на версиях Redis до 6.2, где ещё нет GETDEL.
Наблюдаемость и остановка
Каждый вызов продюсера открывает span OpenTelemetry с именем redis <ОПЕРАЦИЯ>, тегом db.system=redis, ресурсом и операцией. Вызовы Redis видны в той же трассе, что и остальные шаги маршрута.
Консюмеры останавливаются плавно, как и в других коннекторах redb.Route: перестают брать новые сообщения и дорабатывают те, что уже в обработке.
Как поставить
dotnet add package redb.Route dotnet add package redb.Route.Redis
Исходники: github.com/redbase‑app/redb‑route. Пакет на NuGet.
Если было полезно, ⭐ на GitHub поможет другим это найти.
Другие мои статьи — redb.ru/articles, ещё — на Хабре.