Подготовка окружения
Компонента kovalevdmv/1CRabbitMQ, обработка КлиентRMQ. Брокер - как в части 1.
Код примеров собран в расширении 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 - другие.
На чём сосредоточена часть

В частях 1–2 очередь после подтверждения (ПодтвердитьСообщение) забывает сообщение. Stream-очередь ближе к журналу: записи остаются в пределах политики хранения (retention - объём / возраст, например max-length-bytes / max-age). Получатель читает поток и может пройти его снова.
Эта часть - минимальный контур, как «Hello World»:
- объявить durable stream-очередь;
- канал с предзагрузкой (prefetch);
- отправить одно сообщение;
- прочитать его без смещения
x-stream-offset = "first".
Подробности о смещении в - часть 9.
Объявление stream-очереди и канал
Stream-очередь должна быть durable. Тип - аргумент x-queue-type:
КлиентRMQ = Обработки.КлиентRMQ.Создать();
URI = ОМ_РМКУ_Настройки.СтрокаПодключения();
Подключение = КлиентRMQ.ПодключитьсяКСерверу(URI);
// Prefetch обязателен: иначе брокер откажет в consumer на stream
КоличествоПредзагруженных = 10;
Канал = КлиентRMQ.СоздатьКанал(Подключение, Ложь, КоличествоПредзагруженных);
ИмяОчереди = "hello_python_stream"; // или хелпер расширения
АргументыОчереди = Новый Массив;
АргументыОчереди.Добавить(КлиентRMQ.НовыйАргумент("x-queue-type", "stream", "строка"));
// при необходимости retention:
// АргументыОчереди.Добавить(КлиентRMQ.НовыйАргумент("max-length-bytes", 5368709120, "число")); // байты; 5368709120 W76; 5 GiB
// АргументыОчереди.Добавить(КлиентRMQ.НовыйАргумент("max-age", "7D", "строка")); // 7 дней; единицы: Y M D h m s
Очередь = КлиентRMQ.ОбъявитьОчередь(
Канал, ИмяОчереди, "", "",
Истина, // durable
Ложь, Ложь,
АргументыОчереди);
В расширении то же вынесено в ОМ_РМКУ_Настройки.ОбъявитьОчередьStream(...).
Retention. Пока журнал не обрезан политикой хранения, старые записи остаются доступны для повторного чтения. Границы задают аргументами объявления:
max-length-bytes- максимальный размер журнала в байтах (число; например5368709120W76; 5 GiB);max-age- максимальный возраст записи - строка «число + единица»:Yгод,Mмесяц,Dдень,hчас,mминута,sсекунда (например7D- неделя).
Для Hello World достаточно x-queue-type = stream без этих лимитов.
Зачем prefetch. Stream - это журнал: после ack запись остаётся в истории, её можно читать снова. Лимит «сколько сообщений висит без ack» для stream обязателен - иначе брокер откажет в basic.consume. В КлиентRMQ его задают третьим параметром СоздатьКанал (КоличествоПредзагруженных > 0).
Объявление идемпотентно при тех же параметрах. Объявляйте очередь и у отправителя, и у получателя - неизвестно, кто стартует первым (как в части 1).
Отправка

Одно сообщение в stream - через обменник по умолчанию (пустая ТочкаОбмена), routing_key = имя очереди:
ТекстСообщения = "Hello, World!";
Ответ = КлиентRMQ.ОпубликоватьСообщение(
Канал, ТекстСообщения, ИмяОчереди, "", Истина); // persistent по желанию
Если КлиентRMQ.ЭтоОшибка(Ответ) Тогда
ВызватьИсключение Ответ.Текст;
КонецЕсли;
ОМ_РМКУ_Настройки.ЗафиксироватьОтправку(ИмяОчереди, ТекстСообщения, "Урок8_Отправитель");
Получение

Для stream при СоздатьПолучателя нужен аргумент x-stream-offset - откуда начать. В Hello World достаточно "first" (с самого раннего ещё хранимого сообщения):
АргументыПолучателя = Новый Массив;
АргументыПолучателя.Добавить(
КлиентRMQ.НовыйАргумент("x-stream-offset", "first"));
ТаймаутСек = 10;
Получатель = КлиентRMQ.СоздатьПолучателя(
Канал, ИмяОчереди, "stream_hello", Ложь,,,, ТаймаутСек, АргументыПолучателя);
Если КлиентRMQ.СледующееСообщение(Получатель) Тогда
ДанныеСообщения = КлиентRMQ.ДанныеСообщения();
Если Не КлиентRMQ.ЭтоОшибка(ДанныеСообщения) Тогда
Текст = Строка(ДанныеСообщения.Данные);
// ДанныеСообщения.Тег - в демо компоненты W76; offset (подробнее в части 9)
Ид = ОМ_РМКУ_Настройки.ЗафиксироватьПолучение(
ИмяОчереди, Текст, "Урок8_Получатель", ДанныеСообщения.Тег);
КлиентRMQ.ПодтвердитьСообщение(Канал, ДанныеСообщения.Тег);
ОМ_РМКУ_Настройки.ЗафиксироватьПодтверждение(ИмяОчереди, Ид);
КонецЕсли;
КонецЕсли;
КлиентRMQ.ОтменитьПолучателя(Получатель);
КлиентRMQ.ЗакрытьКанал(Канал);
КлиентRMQ.ОтключитьсяОтСервера(Подключение);
Подтверждение на stream не «стирает» запись из журнала - оно двигает окно доставки для этого consumer. Повторное чтение с "first" снова увидит то же сообщение (пока retention его не срезала).
Другие значения offset ("next", число, timestamp) - в части 9.
Запуск решения
YAxUnit - тесты набора «Урок 8» Урок8_StreamHello
Консоль кода (серверный контекст) - те же сценарии:
// smoke: durable stream U94; publish U94; читать с x-stream-offset = "first" U94; ack
ОМ_ТестRMQ_Lessons.Урок8_StreamHello();
Результат работы: в регистре rmq_ОчередьОбменаRMQ - исходящая и входящая запись. В ЖР / окне сообщений после Урок8_StreamHello примерно:
… | Старт | очередь: hello_python_stream | данные: Урок8_StreamHello
… | Отправка в очередь | очередь: hello_python_stream | данные: Hello, World!
… | Получение из очереди | очередь: hello_python_stream | данные: Hello, World!
… | Окончание | очередь: hello_python_stream | данные: Урок8_StreamHello
Канал с prefetch > 0; в конце cleanup: ОтменитьПолучателя U94; ЗакрытьКанал U94; ОтключитьсяОтСервера.
Результат
Подняли Hello World на stream через КлиентRMQ: тип очереди stream, обязательный prefetch, одна публикация и чтение с "first". Дальше - часть 9: смещение, пачка сообщений, курсор на стороне 1С.
Ссылки на остальные части
- Часть 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
Вступайте в нашу телеграмм-группу Инфостарт