Руководство RabbitMQ - Часть 8. Streams - Hello World

06.08.26

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

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

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

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

  1. объявить durable stream-очередь;
  2. канал с предзагрузкой (prefetch);
  3. отправить одно сообщение;
  4. прочитать его без смещения 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 - максимальный размер журнала в байтах (число; например 5368709120 W76; 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С.

 

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

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

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

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

RabbitMQ КлиентRMQ Streams stream Hello World очередь AMQP 1CRabbitMQ интеграция rstream журнал retention prefetch YAxUnit

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

  • 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
Для отправки сообщения требуется регистрация/авторизация