Подготовка окружения
Компонента kovalevdmv/1CRabbitMQ, обработка КлиентRMQ. Брокер - как в части 1 и части 8.
Код примеров собран в расширении RMQ_lessons (malikov-pro/1CRabbitMQ) для запуска через YAxUnit.
Те же сценарии можно вызывать из консоли кода в серверном контексте - например bsl_console и аналоги.
Важно для 1С. Официальный Python-туториал ходит по native stream protocol (порт 5552, клиент rstream). КлиентRMQ работает по AMQP 0-9-1 (порт 5672): stream - это тип очереди x-queue-type = stream, а не отдельный сокет. Идея сценария та же (журнал + чтение со смещения), транспорт и API - другие.
В оригинале consumer останавливается по телу marker:…. В 1С выход из Пока СледующееСообщение - idle-таймаут у СоздатьПолучателя, как в части 2.
На чём сосредоточена часть

В части 8 - Hello World на stream. Здесь - как листать одну и ту же ленту с разных точек (как в оригинале Offset Tracking):
- продюсер публикует 100 сообщений;
- читатель с
"first"- вся пачка; в сообщении - начальное и конечное смещение; - читатель с числом 42 - с середины до конца пачки; снова first/last offset;
- читатель с
"next", затем ещё 100 сообщений - только новая волна; снова first/last offset.
Объявление очереди и канал с prefetch - как в части 8 (ОбъявитьОчередьStream / prefetch > 0).
Смещение: откуда начать читать
Stream - лента событий. Каждое сообщение лежит на своём месте - смещение (offset): абсолютный номер в журнале брокера (0, 1, 2…). Это не «счётчик вашего читателя», а координата в общей ленте.

Аргумент x-stream-offset при СоздатьПолучателя:
| Значение | Смысл |
|---|---|
"first" |
с самого раннего ещё хранимого |
"next" |
только то, что появится после подписки |
"last" |
с последнего «куска» записей |
| число | точный offset |
дата + тип "timestamp" |
первое сообщение не раньше этой отметки |
Числовой offset сообщения - в заголовке x-stream-offset (ДанныеСообщения.Заголовки). Поле Тег - это delivery tag для ПодтвердитьСообщение, не смещение в ленте.
Отправка
Продюсер объявляет stream (если ещё нет) и публикует 100 сообщений hello: 0 … hello: 99:
ИмяОчереди = ОМ_РМКУ_Настройки.ИмяОчередиЧасть9(); // stream_offset_tracking
Очередь = ОМ_РМКУ_Настройки.ОбъявитьОчередьStream(КлиентRMQ, Канал, ИмяОчереди);
ЧислоСообщений = 100;
Для Номер = 0 По ЧислоСообщений - 1 Цикл
КлиентRMQ.ОпубликоватьСообщение(
Канал, "hello: " + Номер, ИмяОчереди, "", Истина);
КонецЦикла;
Publisher confirms (часть 7) - по желанию, если нужно дождаться приёма пачки на брокере.
Читатель с first
Подписка с начала ленты. В цикле запоминаем offset первого и последнего полученного сообщения; после idle-таймаута выводим их:
АргументыПолучателя = Новый Массив;
АргументыПолучателя.Добавить(
КлиентRMQ.НовыйАргумент("x-stream-offset", "first"));
ТаймаутОжиданияСек = 1; // idle: выход, когда пачка разобрана
Получатель = КлиентRMQ.СоздатьПолучателя(
Канал, ИмяОчереди, "stream_first", Ложь,,,,
ТаймаутОжиданияСек, АргументыПолучателя);
НачальноеСмещение = Неопределено;
КонечноеСмещение = Неопределено;
Пока КлиентRMQ.СледующееСообщение(Получатель) Цикл
ДанныеСообщения = КлиентRMQ.ДанныеСообщения();
Если КлиентRMQ.ЭтоОшибка(ДанныеСообщения) Тогда
Продолжить;
КонецЕсли;
ТекущееСмещение = Неопределено;
Для Каждого Заголовок Из ДанныеСообщения.Заголовки Цикл
Если Заголовок.Ключ = "x-stream-offset" Тогда
ТекущееСмещение = Число(Заголовок.Значение);
Прервать;
КонецЕсли;
КонецЦикла;
Если НачальноеСмещение = Неопределено Тогда
НачальноеСмещение = ТекущееСмещение;
КонецЕсли;
КонечноеСмещение = ТекущееСмещение;
КлиентRMQ.ПодтвердитьСообщение(Канал, ДанныеСообщения.Тег); // Тег = delivery tag
КонецЦикла;
КлиентRMQ.ОтменитьПолучателя(Получатель);
ОМ_РМКУ_Настройки.СообщитьКонтрольнуюТочку(
"Читатель first",
СтрШаблон("first_offset=%1 last_offset=%2", НачальноеСмещение, КонечноеСмещение),
ИмяОчереди);
// ожидаемо: first_offset=0 last_offset=99
Таймаут - ожидание пустой выдачи, не «время на обработку» (см. часть 2).
Читатель с offset 42
Тот же stream можно перечитать с произвольной точки. Подставьте число вместо "first":
АргументыПолучателя = Новый Массив;
АргументыПолучателя.Добавить(
КлиентRMQ.НовыйАргумент("x-stream-offset", 42)); // тип "число" подставится сам
Цикл тот же: фиксируете first_offset и last_offset. Ожидаемо для пачки 0…99: first_offset=42, last_offset=99.
Читатель с next и вторая пачка
"next" - только сообщения, появившиеся после подписки. В оригинале receiver ждёт во втором терминале, пока sender шлёт новую волну. В одном сеансе 1С порядок такой: сначала создать получателя с "next", затем опубликовать ещё 100 сообщений, затем цикл чтения.
АргументыПолучателя = Новый Массив;
АргументыПолучателя.Добавить(
КлиентRMQ.НовыйАргумент("x-stream-offset", "next"));
Получатель = КлиентRMQ.СоздатьПолучателя(
Канал, ИмяОчереди, "stream_next", Ложь,,,,
ТаймаутОжиданияСек, АргументыПолучателя);
// вторая волна - уже после подписки
Для Номер = 0 По ЧислоСообщений - 1 Цикл
КлиентRMQ.ОпубликоватьСообщение(
Канал, "hello: " + Номер, ИмяОчереди, "", Истина);
КонецЦикла;
НачальноеСмещение = Неопределено;
КонечноеСмещение = Неопределено;
Пока КлиентRMQ.СледующееСообщение(Получатель) Цикл
// … тот же учёт x-stream-offset U94; Начальное / Конечное, ack по Тег …
КонецЦикла;
КлиентRMQ.ОтменитьПолучателя(Получатель);
ОМ_РМКУ_Настройки.СообщитьКонтрольнуюТочку(
"Читатель next",
СтрШаблон("first_offset=%1 last_offset=%2", НачальноеСмещение, КонечноеСмещение),
ИмяОчереди);
// ожидаемо: first_offset=100 last_offset=199
Итого «обзор» stream: с начала, с любого offset, только новых сообщений.
Server-Side Offset Tracking
Server-side offset tracking - механизм RabbitMQ Streams, при котором брокер сам хранит прогресс consumer’а в ленте. У приложения есть стабильное имя (tracking reference); клиент периодически вызывает store_offset, а при следующем запуске - query_offset и продолжает с сохранённой точки, не обрабатывая уже прочитанные сообщения заново. Подробнее - в блоге RabbitMQ и оригинале. Эта возможность есть только в native stream protocol.
Запуск решения
YAxUnit - тесты набора «Урок 9» Урок9_StreamOffset
Консоль кода (серверный контекст) - те же сценарии:
// 100 U94; first (0..99) U94; 42 (42..99) U94; next + ещё 100 (100..199)
ОМ_ТестRMQ_Lessons.Урок9_StreamOffset();
Результат работы: в регистре rmq_ОчередьОбменаRMQ - исходящие по волнам и входящие с first_offset / last_offset по каждому читателю. В ЖР / окне сообщений после Урок9_StreamOffset примерно:
… | Старт | очередь: stream_offset_tracking | данные: Урок9_StreamOffset
… | Отправка в очередь | очередь: stream_offset_tracking | данные: Publishing 100 messages
… | Читатель first | очередь: stream_offset_tracking | данные: first_offset=0 last_offset=99
… | Читатель offset 42 | очередь: stream_offset_tracking | данные: first_offset=42 last_offset=99
… | Отправка в очередь | очередь: stream_offset_tracking | данные: Publishing 100 messages
… | Читатель next | очередь: stream_offset_tracking | данные: first_offset=100 last_offset=199
… | Окончание | очередь: stream_offset_tracking | данные: Урок9_StreamOffset
Канал с prefetch > 0; выход из циклов чтения - по idle-таймауту; в конце cleanup: ОтменитьПолучателя → ЗакрытьКанал → ОтключитьсяОтСервера.
Результат
Научились листать stream через x-stream-offset: "first", число, "next" - с фиксацией начального и конечного смещения.
Ссылки на остальные части
- Часть 1. Hello World
- Часть 2. Work Queues - рабочие очереди
- Часть 3. Publish/Subscribe - одно сообщение многим (fanout)
- Часть 4. Routing - маршрутизация по ключу (direct)
- Часть 5. Topics - маршрутизация по шаблону
- Часть 6. RPC - удалённый вызов процедур
- Часть 7. Publisher Confirms - надёжная публикация
- Часть 8. Streams - Hello World
- Часть 9. Streams - отслеживание смещения (эта статья)
Благодарю за внимание.
Создано совместно с Cursor Grok 4.5
Вступайте в нашу телеграмм-группу Инфостарт