Руководство RabbitMQ - Часть 9. Streams - отслеживание смещения

06.08.26

Интеграция - Внешние источники данных

Адаптация сценария RabbitMQ Streams под платформу 1С и КлиентRMQ (очередь stream, x-stream-offset). Источник https://www.rabbitmq.com/tutorials/tutorial-two-python-stream

Подготовка окружения

Компонента 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):

  1. продюсер публикует 100 сообщений;
  2. читатель с "first" - вся пачка; в сообщении - начальное и конечное смещение;
  3. читатель с числом 42 - с середины до конца пачки; снова first/last offset;
  4. читатель с "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: 0hello: 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" - с фиксацией начального и конечного смещения.

 

Ссылки на остальные части

Благодарю за внимание.

Создано совместно с Cursor Grok 4.5

Вступайте в нашу телеграмм-группу Инфостарт

RabbitMQ КлиентRMQ Streams stream offset смещение отслеживание смещения x-stream-offset prefetch retention журнал AMQP 1CRabbitMQ интеграция

Вы можете заказать платную адаптацию этой статьи под ваши задачи на «Бирже заказов».

  • 0% комиссии — оплата напрямую исполнителю;
  • Исполнители любого масштаба — от отдельных специалистов до команд под проект;
  • Прямой обмен контактами между заказчиком и исполнителем;
  • Безопасная сделка — при необходимости;
  • Рейтинги, кейсы и прозрачная система откликов.

См. также

Внешние источники данных Программист Бизнес-аналитик Пользователь 1С:Предприятие 8 1C:Бухгалтерия Узбекистан Беларусь Кыргызстан Молдова Россия Казахстан Платные (руб)

Готовое решение для автоматической выгрузки данных из 1С 8.3 в базу данных ClickHouse, PostgreSQL или Microsoft SQL для работы с данными 1С в BI-системах. «Экстрактор данных 1С в BI» работает со всеми типовыми и нестандартными конфигурациями 1С 8.3 и упрощает работу бизнес-аналитиков. Благодаря этому решению, специалистам не требуется быть программистами, чтобы легко получать данные из 1С в вашей BI-системе.

35000 руб.

15.11.2022    32442    50    49    

49

Внешние источники данных Кадровый учет Файловый обмен (TXT, XML, DBF), FTP Перенос данных 1C Программист 1С:Предприятие 8 1С:Зарплата и кадры государственного учреждения 3 Государственные, бюджетные структуры Россия Бухгалтерский учет Бюджетный учет Платные (руб)

Обработка позволяет перенести кадровую информацию и данные по заработной плате, фактическим удержаниям, НДФЛ, вычетам, страховым взносам из базы Парус 10 учреждений (далее Парус) в конфигурацию 1С:Зарплата и кадры государственного учреждения ред. 3 (далее 1С) и начать с ней работать с любого месяца года.

85400 руб.

05.10.2022    14066    16    8    

17

Розничная торговля Внешние источники данных Файловый обмен (TXT, XML, DBF), FTP Системный администратор Программист 1С:Предприятие 8 1С:Бухгалтерия 3.0 Фармацевтика, аптеки Россия Бухгалтерский учет Платные (руб)

Внешняя обработка загрузки данных из файла-выгрузки, сформированного в программе F3 TAIL версии 3.4 (и выше) или еФарма версии 2.1, в базу конфигурации 1С: Бухгалтерия предприятия 8, ред. 3.0 (Базовая, ПРОФ, КОРП, ФРЕШ (тонкий клиент)).

17080 руб.

19.12.2016    54771    126    107    

86

Производство готовой продукции (работ, услуг) Внешние источники данных 1С:Предприятие 8 1С:Управление нашей фирмой 1.6 Лесное и деревообрабатывающее хозяйство Россия Управленческий учет Платные (руб)

Обработка предназначена для загрузки файлов, выгруженных из системы Базис-мебельщик, в справочник 1С "Спецификации" для последующих процессов учета и диспетчирования полуфабрикатов и изделий.

10370 руб.

24.06.2021    26214    64    55    

47

Внешние источники данных Программист Бизнес-аналитик 1С:Предприятие 8 1С:Управление производственным предприятием 1С:Бухгалтерия 3.0 1С:Управление торговлей 11 1С:Комплексная автоматизация 2.х 1С:Зарплата и Управление Персоналом 3.x 1С:Управление нашей фирмой 3.0 1С:Розница 3.0 Платные (руб)

Обработка для выгрузки данных из подготовленных СКД в фоновом режиме в базу ClickHouseDB, PostgreSQL, MySQL, в шину данных с поддержкой REST API (CSV, JSON. SQL), в локальные файлы (CSV, JSON, XLS, XLSX) или в Google Sheets. Это дополнительная подключаемая обработка.

18000 руб.

21.08.2024    9813    25    4    

22

Оптовая торговля Розничная торговля Внешние источники данных Прайсы 1С:Предприятие 8 1С:ERP Управление предприятием 2 1С:Управление торговлей 11 1С:Комплексная автоматизация 2.х Розничная и сетевая торговля (FMCG) Оптовая торговля, дистрибуция, логистика Управленческий учет Платные (руб)

Хотите, чтобы остатки и цены товаров в вашей базе всегда были актуальными без лишних усилий? Теперь это возможно - автоматизируйте процесс загрузки и обновления данных о номенклатуре от ваших поставщиков или конкурентов. Как это работает? Вы сами настраиваете правила и расписание для каждого поставщика, чтобы обновление информации из произвольных форматов прайс-листов происходило автоматически.

15250 руб.

15.05.2024    4822    8    1    

9
Для отправки сообщения требуется регистрация/авторизация