Redis в маршрутах .NET: кэш, Pub/Sub, Streams с группами и Claim Check одним коннектором

redb.Route

В интеграционном проекте 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:

Структура Операции
Строки и ключи SET, GET, DEL, EXISTS, EXPIRE, INCR, DECR, SETNX
Списки LPUSH, RPUSH, LPOP, RPOP, LLEN, LRANGE
Хэши HSET, HGET, HMSET, HMGET, HGETALL, HDEL, HLEN
Множества SADD, SREM, SMEMBERS, SCARD, SISMEMBER
Sorted sets ZADD, ZREM, ZRANGE, ZCARD, ZSCORE, ZRANGEBYSCORE
Гео GEOADD, GEODIST, GEORADIUS
HyperLogLog PFADD, PFCOUNT, PFMERGE
Битовые карты SETBIT, GETBIT, BITCOUNT
Сообщения PUBLISH, XADD
Любая команда COMMAND с именем команды в CustomCommand(...) и аргументами в заголовке

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

  • значение или элемент это тело сообщения. В списках, хэшах, множествах, 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, ещё — на Хабре.