redb.Route.Redis
redb.Route.Redis

Коннектор 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 есть готовые фабрики (SetGetIncrPublishSubscribeXAddXReadLPush и остальные). Любую другую операцию задаёт Redis.Command(операция, ключ): например, Redis.Command("HSET", "customer:42").Field("email").

Все структуры данных как шаги маршрута

Продюсер поддерживает операции над всеми основными структурами Redis:

Структура

Операции

Строки и ключи

SETGETDELEXISTSEXPIREINCRDECRSETNX

Списки

LPUSHRPUSHLPOPRPOPLLENLRANGE

Хэши

HSETHGETHMSETHMGETHGETALLHDELHLEN

Множества

SADDSREMSMEMBERSSCARDSISMEMBER

Sorted sets

ZADDZREMZRANGEZCARDZSCOREZRANGEBYSCORE

Гео

GEOADDGEODISTGEORADIUS

HyperLogLog

PFADDPFCOUNTPFMERGE

Битовые карты

SETBITGETBITBITCOUNT

Сообщения

PUBLISHXADD

Любая команда

COMMAND с именем команды в CustomCommand(...) и аргументами в заголовке

Откуда операция берёт данные, определено одинаково для всех:

  • значение или элемент это тело сообщения. В списках, хэшах, множествах, sorted sets и стримах тело byte[] пишется как есть, без перевода в строку;

  • скалярные параметры задаются в адресе: field для хэшей, scoreminScore и maxScore для sorted sets, start и stop для диапазонов, longitudelatitudemember1member2 и geoUnit для гео, offset и bit для битовых карт;

  • наборы значений передаются заголовками: словарь полей для HMSET (redbRedis.HashFields), имена полей для HMGET (redbRedis.FieldNames), исходные ключи для PFMERGE (redbRedis.SourceKeys), центр и радиус для GEORADIUS, аргументы для COMMAND.

Ключ, поле, канал и имя стрима могут содержать шаблоны ${...}: они вычисляются на каждом сообщении.

Результат операции становится телом для следующего шага. GET отдаёт строку или nullINCR новое значение счётчика, 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) и время получения.

Подписка по шаблону делается через PSUBSCRIBEFrom(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.StreamFieldsStreamMaxLength обрезает стрим по длине, по умолчанию приблизительно (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 в голову и чтение с хвоста дают очередь в порядке поступления. Консюмер забирает элементы по одному опросом, а когда список пуст, ждёт pollDelayMsBLPOP и 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, ещё — на Хабре.

Комментарии (0)