Механизм
Поставщика информации принято называть
Информация, посылаемая как сообщение, характеризуется темой.
(рис 10.1) Схема взаимодействия брокеров, издателей и подписчиковОбщий подход работы механизма
Сценарий
Сценарий
Сценарий
Механизм
Для работы с механизмом
Регион/СекторРынка/Компания
Конкретные темы выглядят следующим образом:
NewYork/InformationTechnology/IBM
Этот синтаксис разрешает в строках иметь символы:
* - означающий ноль или любое количество
произвольных символов;
? – означающий один произвольный символ.
Данный синтаксис позволит в дальнейшем определить, например, такие темы для подписки на рынке акций как:
* получение всех цен на акции
со всех рынков мира;
London/* получение всех цен на акции
с Лондонской биржи;
NewYork/Banks/* получение цен акций всех
нью-йоркских банков;
*/*/IBM получение всех цен на акции
компании ИБМ на всех рынках мира.
Символ ? в строках, например, 'ABC%*D' означает 'ABC*D'.
В числе основных функций для работы с механизмом
RegPub |
Регистрация издателя (Register Publisher) |
RegSub |
Регистрация |
Publish |
Публикация |
ReqUpdate |
Запрос от издателя к |
DeletePub |
Удаление публикации (Delete Publication) |
DeregPub |
Отказ от регистрации издателя (Deregister Publisher) |
DeregSub |
Отказ от регистрации |
Более подробно эти функции будут рассмотрены на примере ниже, а описание их форматов можно найти в документации , .
Для осуществления работы издателя и
typedef struct tagMQRFH {
MQCHAR4 StrucId;
/* Идентификатор структуры */
MQLONG Version;
/* Номер версии структуры */
MQLONG StrucLength;
/* Общая длина MQRFH */
MQLONG Encoding;
/* Data encoding */
MQLONG CodedCharSetId;
/* Идентификатор множества перекодировки */
MQCHAR8 Format;
/* Имя формата */
MQLONG Flags;
/* Флаги */
} MQRFH;
Этот заголовок также включает строку NameValueString, в которой приложение публикации или подписки помещает команды, которые должен выполнить StrucLength в заголовке определяет длину структуры заголовка, включительно с переменной длины NameValueString в конце структуры. Поля Encoding, CodedCharSetId и Format описывают структуру данных в заголовке публикации.
Задача создания собственных приложений Издателя и
(рис 10.2) Модель Издатель/ПодписчикПриложение Издателя должно публиковать символьные строки, задаваемые пользователем на определенные темы. Следовательно, такая программа с именем publisher может иметь следующий формат запуска:
publisher Topic Command QMgrName PubQueue
где Topic – тема публикации, символьная строка длиной не более 256 байт; Command – команды Издателя для RegPub, Publish, ReqUpdate, DeletePub, DeregPub ), PubQueue - очередь издателя. Текст публикации вводится в командном окне с запущенной программой.
Приложение subscriber может иметь формат запуска:
subscriber Topic QMgrName SubQueue
где Topic – тема SubQueue - очередь Topic отображается в командном окне с запущенной программой subscriber.
Для работы приложений требуется создание очередей: одна у Издателя и одна у Publisher_queue и Subscriber_queue. Publisher_queue. Subscriber_queue. Создание этих очередей осуществляется простой командой :
runmqsc QMgrName define qlocal(имя очереди) end
Следует рассмотреть потоки сообщений и то, как сообщения по определенной теме находят своего
SYSTEM.BROKER.CONTROL.QUEUE, содержащие TestTopic. SYSTEM.BROKER.*.TestTopic и оно попадает в потоковую очередь SYSTEM.BROKER.DEFAULT.STREAM
Сообщения
Также как и сообщения
Основные шаги, из которых складывается создание Приложения издателя. Таких шагов для программы publisher будет 7.
Шаг 1. Для того чтобы положить сообщение в очередь потоков для помещения сообщений.
Шаг 2. Для публикации необходимо сформировать MQRFH структуру (см. руководство по MQRFH должны быть изменены по сравнению со значениями по умолчанию перед публикацией или в тот момент, когда мы определим значения этих полей.
Шаг 3. Сразу за структурой MQRFH должна следовать символьная строка NameValueString. Указатель pNameValueString определяет начальную позицию NameValueString в теле сообщения.
Шаг 4. Содержимое строки NameValueString должно включать всю необходимую информацию о публикации и всю необходимую последовательность команд для формирования непрерывной символьной строки нашей публикации. Содержимое строки NameValueString позволяет
Шаг 5. Поле StrucLength должно содержать длину MQRFH структуры и сопровождающей его строки NameValueString. Длина MQRFH фиксирована и длина NameValueString – переменная. Значение StrucLength позволяет не применять разделитель конца строки в NameValueString, хотя и он может быть использован, если это необходимо приложению- MQRFH и NameValueString выравниваются на границе 16 байт автоматически.
Шаг 6. Данные о публикации (в данном случае строка данных для публикации) помещаются сразу после MQFRH и NameValueString структуры, указатель pUserData должен быть установлен на начальную позицию строки введенных данных о публикации.
Шаг 7. Осуществляется процедура тестирования работы программы. Программа компилируется, устраняются ошибки компиляции. Она вызывается, например, командой:
publisher TestTopic Publish QM_broker Publisher_queue
В командном окне программы publisher вводится сообщение по теме: Hello
На этом этапе еще нет приложения
Publisher_queue и исправить это; есть ли успешные сообщения от SYSTEM.DEFAULT.LOCAL.STREAM, используя MQ explorer или команду amqsbcg и убедиться, что в формате сообщения о публикации нет ошибок.Основные шаги, из которых складывается создание Приложения подписчика.
Шаг 1. Прежде всего, командой необходимо открыть очередь
Шаг 2. На этом шаге необходимо зарегистрировать интересы в виде темы и послать SendBrokerCommand, одним из аргументов которой является командная строка для помещения в NameValueString командного сообщения. Эта функция SendBrokerCommand должна быть аналогичной по коду, как и для приложения Издателя.
Шаг 3. Теперь, когда . Первое, что необходимо сделать при получении сообщения, это проверить, что оно в нужном формате MQRFH.
Шаг 4. После распознавания формата MQRFH сообщения, можно выделить порцию наиболее интересного сообщения. В данном случае нет необходимости смотреть на строку NameValueString так как интересуют данные, которые следуют за этой строкой, они задаются указателем pUserData и могут быть напечатаны на экране или выведены в файл.
Шаг 5. На этом последнем шаге необходимо отказаться от регистрации и, как это делалось раньше на шаге 2, вызвать функцию SendBrokerCommand, добавив соответствующую команду для отказа от регистрации. Теперь, так же как и раньше, надо скомпилировать код и затем исправить ошибки компиляции. После этого программа готова для работы, ее можно выполнить:
subscriber TestTopic QM_broker Subscriber_queue
Этот запуск должен отобразить информацию об успешной регистрации у publisher на тему TestTopic с сообщением, например, PublisherReady. В результате, в программе- PublisherReady.
Если этого не произойдет, то необходимо найти ошибку.
Получив успешно одну публикацию от одного издателя, можно расширить эксперимент и попробовать запустить много программ-издателей и программ-
Теперь после знакомства с технологией
После этого можно стартовать
strmqbrk -m QMgrName
Для отображения состояния dspmqbrk:
dspmqbrk -m QMgrName
В ответ появится следующее сообщение:
WebSphere MQ message broker for queue manager QMgrName running
Теперь
На каждом менеджере может быть стартован только один
В менеджере есть необходимые
runmqsc QMgrName display qlocal(SYSTEM.BROKER.*) end
Следует сразу отметить, что завершение работы endmqbrk перед окончанием работы менеджера: endmqbrk -m QMgrName.
Работу издателя можно продемонстрировать с помощью программы amqsgama, предложенной в SupportPacs MA0C в качестве теста. Эта программа издателя из перечня спортивных тем для подписки (табл.10.1) помещает на Sport/Soccer/Event/ - в таблице данная тема выделена курсивом) и проверяет ответы
спорт/футбол/* спорт/теннис/* спорт/баскетбол/* |
спорт/футбол/расписание игр спорт/футбол/события спорт/футбол/обзоры |
Формат запуска программы:
amqsgama TeamName1 TeamName2 QMgrName
Результаты работы программы, моделирующей случайным образом забивание голов той или иной командой, выглядит следующим образом (рис. 10.3):
(рис 10.3) Результаты работы программы издателяДля работы программы необходимо создание очереди: SAMPLE.BROKER.RESULTS.STREAM.
Именно в эту очередь поступают сообщения от издателя. Необходимо также, чтобы был запущен amqsgama, чтобы отобразить результаты игры полностью. Все используемые функции в программе служат для подключения к менеджеру
Блочная структура программы выглядит следующим образом.
Подключение к менеджеру брокера (MQCONN)
Открытие очереди потока брокера (MQOPEN)
Инициализация таймера матча
Генерация MQRFH для публикации события о
начале матча
Добавление имен команд в данные
Помещение публикации в очередь потоков
Начало цикла по времени матча:
засыпание на случайный период
попытка забить гол (50% вероятность)
генерация публикации о забитом голе
(RFH для ScoreUpdate)
случайный выбор команды, забившей гол
добавление имени команды в данные для
публикации
помещение публикации в очередь потоков
Окончание цикла по времени матча
Генерация MQRFH для публикации события о
конце матча
Добавление имен команд для публикации
Помещение публикации в очередь потоков
Закрытие очереди потока брокера (MQCLOSE)
Отключение от менеджера брокера (MQDISC)
Программа amqsgama имеет следующий код:
/*************************************************************************************/
/* Имя программы: AMQSGAMA */
/* Описание: Основанная на модели Publish/Subscribe программа */
/* моделирует результаты футбольного матча и */
/* отправляет их от издателя к брокеру */
/* Statement: Licensed Materials - Property of IBM */
/* SupportPac MA0E */
/* (C) Copyright IBM Corp. 1999 */
/*************************************************************************************/
#include <stdlib.h>
#include <stdio.h>
#include <string.h>
#include <time.h>
#include <cmqc.h> /* MQI */
#include <cmqpsc.h> /* MQI Publish/Subscribe */
#include <windows.h>
#if MQAT_DEFAULT == MQAT_WINDOWS_NT
#define msSleep(time) \
Sleep(time)
#elif MQAT_DEFAULT == MQAT_UNIX
#define msSleep(time) \
{ \
struct timeval tval; \
tval.tv_sec = time / 1000; \
tval.tv_usec = (time % 1000) * 1000; \
select(0, NULL, NULL, NULL, tval); \
}
#endif
#define STREAM "SAMPLE.BROKER.RESULTS.STREAM"
#define TOPIC_PREFIX "Sport/Soccer/Event/"
#define MATCH_STARTED "MatchStarted"
#define MATCH_ENDED "MatchEnded"
#define SCORE_UPDATE "ScoreUpdate"
#define MATCH_LENGTH 30000 /* 30 Second match length */
#define REAL_TIME_RATIO 333
#define AVERAGE_NUM_OF_GOALS 5
#define DEFAULT_MESSAGE_SIZE 512 /* Maximum buffer size for a message */
static const MQRFH DefaultMQRFH = {MQRFH_DEFAULT};
typedef struct
{
MQCHAR32 Team1;
MQCHAR32 Team2;
} Match_Teams, *pMatch_Teams;
void BuildMQRFHeader( PMQBYTE pStart
, PMQLONG pDataLength
, MQCHAR TopicType[] );
void PutPublication( MQHCONN hConn
, MQHOBJ hObj
, PMQBYTE pMessage
, MQLONG messageLength
, PMQLONG pCompCode
, PMQLONG pReason );
int main(int argc, char **argv)
{
MQHCONN hConn = MQHC_UNUSABLE_HCONN;
MQHOBJ hObj = MQHO_UNUSABLE_HOBJ;
MQLONG CompCode;
MQLONG Reason;
MQOD od = { MQOD_DEFAULT };
MQLONG Options;
PMQBYTE pMessageBlock = NULL;
MQLONG messageLength;
MQLONG timeRemaining;
MQLONG delay;
PMQCHAR pScoringTeam;
pMatch_Teams pTeams;
MQCHAR32 team1;
MQCHAR32 team2;
char QMName[MQ_Q_MGR_NAME_LENGTH+1] = "";
MQLONG randomNumber;
MQLONG ConnReason;
/* Проверка аргументов программы */
if( (argc < 3)||(argc > 4)||(strlen(argv[1]) > 31)||(strlen(argv[2]) > 31) )
{
printf("Usage: amqsgam team1 team2 <QManager>\n");
printf(" Maximum 31 characters per team name,\n");
printf(" no spaces or '\"' characters allowed.\n");
exit(0);
}
else
{
strcpy(team1, argv[1]);
strcpy(team2, argv[2]);
}
/* Использовать default queue manager или заданный в зависимости от наличия аргумена */
if (argc > 3) strcpy(QMName, argv[3]);
MQCONN( QMName, hConn, CompCode, ConnReason );
if( CompCode == MQCC_FAILED )
{
printf("MQCONN failed with CompCode %d and Reason %d\n", CompCode, ConnReason);
}
else if( ConnReason == MQRC_ALREADY_CONNECTED )
{
CompCode = MQCC_OK;
}
if( CompCode == MQCC_OK )
{
strncpy(od.ObjectName, STREAM, (size_t)MQ_Q_NAME_LENGTH);
Options = MQOO_OUTPUT + MQOO_FAIL_IF_QUIESCING;
MQOPEN( hConn, od, Options, hObj, CompCode, Reason );
if( CompCode != MQCC_OK )
{
printf("MQOPEN failed to open \"%s\"\nwith CompCode %d and Reason %d\n",
od.ObjectName, CompCode, Reason);
}
}
if( CompCode == MQCC_OK )
{
srand( (unsigned)(time( NULL ))
+ (unsigned)(team1[0] + team2[(strlen(team2) - 1)]) );
timeRemaining = MATCH_LENGTH;
messageLength = DEFAULT_MESSAGE_SIZE;
pMessageBlock = (PMQBYTE)malloc(messageLength);
if( pMessageBlock == NULL )
{
printf("Unable to allocate storage\n");
}
else
{
if( CompCode == MQCC_OK )
{
/* создание MQRFH для публикации о начале матча */
BuildMQRFHeader( pMessageBlock, messageLength, MATCH_STARTED );
pTeams = (pMatch_Teams)(pMessageBlock + messageLength);
strcpy(pTeams->Team1, team1);
strcpy(pTeams->Team2, team2);
messageLength += sizeof(Match_Teams);
printf("Match between %s and %s\n", team1, team2);
/* помещение сообщения (публикации) в очередь потока */
PutPublication( hConn, hObj, pMessageBlock, messageLength, CompCode, Reason );
if( CompCode != MQCC_OK )
{
printf("MQPUT failed with CompCode %d and Reason %d\n", CompCode, Reason);
}
else
{
/* Моделирование попытки забить один из 5 голов (каждые 30 сек) */
while( (timeRemaining > 0)(CompCode == MQCC_OK) )
{
randomNumber = rand();
delay = REAL_TIME_RATIO
+ ( (randomNumber * MATCH_LENGTH)
/ (RAND_MAX * AVERAGE_NUM_OF_GOALS));
if( delay > timeRemaining ) delay = timeRemaining;
msSleep(delay);
timeRemaining -= delay;
if( timeRemaining > 0 )
{
if( (randomNumber % 2) == 0 ) /* Шанс забить гол - 50 процентов */
{
messageLength = DEFAULT_MESSAGE_SIZE;
BuildMQRFHeader( pMessageBlock , messageLength, SCORE_UPDATE );
printf("GOAL! ");
pScoringTeam = (PMQCHAR)pMessageBlock + messageLength;
if( rand() < (RAND_MAX/2) )
{
strcpy(pScoringTeam, team1);
printf(team1);
}
else
{
strcpy(pScoringTeam, team2);
printf(team2);
}
printf(" scores after %d minutes\n", ((MATCH_LENGTH - timeRemaining)/REAL_TIME_RATIO));
messageLength += sizeof(MQCHAR32);
/* помещение сообщения о забитом голе в очередь потока */
PutPublication( hConn, hObj, pMessageBlock, messageLength, CompCode, Reason );
if( CompCode != MQCC_OK )
printf("MQPUT failed with CompCode %d and Reason %d\n", CompCode, Reason);
}
}
} /* конец цикла по времени окончания матча ( timeRemaining ) */
if( CompCode == MQCC_OK )
{
printf("Full time\n");
messageLength = DEFAULT_MESSAGE_SIZE;
BuildMQRFHeader( pMessageBlock , messageLength , MATCH_ENDED );
pTeams = (pMatch_Teams)(pMessageBlock + messageLength);
strcpy(pTeams->Team1, team1);
strcpy(pTeams->Team2, team2);
messageLength += sizeof(Match_Teams);
/* помещение сообщения о конце матча в очередь потока */
PutPublication( hConn, hObj, pMessageBlock, messageLength, CompCode, Reason );
if( CompCode != MQCC_OK )
printf("MQPUT failed with CompCode %d and Reason %d\n", CompCode, Reason);
}
}
}
free( pMessageBlock );
} /* end of else (pMessageBlock != NULL) */
}
if( hObj != MQHO_UNUSABLE_HOBJ )
{
MQCLOSE( hConn , hObj, MQCO_NONE, CompCode , Reason );
if( CompCode != MQCC_OK )
printf("MQCLOSE failed with CompCode %d and Reason %d\n", CompCode, Reason);
}
if( (hConn != MQHC_UNUSABLE_HCONN) (ConnReason != MQRC_ALREADY_CONNECTED) )
{
MQDISC( hConn, CompCode, Reason );
if( CompCode != MQCC_OK )
printf("MQDISC failed with CompCode %d and Reason %d\n", CompCode, Reason);
}
return(0);
}
/* end of main Function */
/* Function Name : BuildMQRFHeader */
void BuildMQRFHeader( PMQBYTE pStart, PMQLONG pDataLength, MQCHAR TopicType[] )
{
PMQRFH pRFHeader = (PMQRFH)pStart;
PMQCHAR pNameValueString;
memset((PMQBYTE)pStart, 0, *pDataLength);
memcpy( pRFHeader, DefaultMQRFH, (size_t)MQRFH_STRUC_LENGTH_FIXED);
memcpy( pRFHeader->Format, MQFMT_STRING, (size_t)MQ_FORMAT_LENGTH);
pRFHeader->CodedCharSetId = MQCCSI_INHERIT;
pNameValueString = (MQCHAR *)pRFHeader + MQRFH_STRUC_LENGTH_FIXED;
strcpy(pNameValueString, MQPS_COMMAND_B);
strcat(pNameValueString, MQPS_PUBLISH);
strcat(pNameValueString, MQPS_PUBLICATION_OPTIONS_B);
strcat(pNameValueString, MQPS_NO_REGISTRATION);
strcat(pNameValueString, MQPS_TOPIC_B);
strcat(pNameValueString, TOPIC_PREFIX);
strcat(pNameValueString, TopicType);
*pDataLength = MQRFH_STRUC_LENGTH_FIXED + ((strlen(pNameValueString)+15)/16)*16;
pRFHeader->StrucLength = *pDataLength;
}
/* Function Name : PutPublication */
void PutPublication( MQHCONN hConn, MQHOBJ hObj, PMQBYTE pMessage,
MQLONG messageLength, PMQLONG pCompCode, PMQLONG pReason )
{
MQPMO pmo = { MQPMO_DEFAULT };
MQMD md = { MQMD_DEFAULT };
memcpy(md.Format, MQFMT_RF_HEADER, (size_t)MQ_FORMAT_LENGTH);
md.MsgType = MQMT_DATAGRAM;
md.Persistence = MQPER_PERSISTENT;
pmo.Options |= MQPMO_NEW_MSG_ID;
MQPUT( hConn, hObj, md, pmo, messageLength, pMessage, pCompCode, pReason );
}
В качестве комментария следует отметить, что функция BuildMQRFHeader формирует значения по умолчанию для заголовка MQRFH, устанавливает параметры format и CCSID пользовательских данных. В строку NameValueString добавляются команды, тема и опции для публикации и она выравнивается на 16-ти байтовую границу. StrucLength в MQRFH устанавливается как общая длина. Входными параметрами функции являются pStart – начало блока сообщения, TopicType[] – строка с именем темы. Входным и выходным параметром одновременно является pDataLength – размер блока сообщения при входе и размер выходного блока информации.
Функция PutPublication формирует сообщение для вывода в очередь брокера с помощью команды . Входными параметрами функции являются hConn – идентификатор менеджера для команды MQHCONN, hObj – идентификатор очереди, pMessage – идентификатор на начало блока сообщения, messageLength – длина данных в сообщении. Выходными параметрами функции являются pCompCode и pReason – коды завершения команды .
Работу amqsresa из состава SupportPacs MA0C, которая подписывается у
amqsresa QMgrName
где QmgrName – имя менеджера очередей, на котором запущен
Необходимая тема подписки ( TOPIC = "Sport/Soccer/*" ) и очереди для
Для работы программы необходимо создание очередей:
runmqsc QMgrName define qlocal(SAMPLE.BROKER.RESULTS.STREAM) define qlocal(RESULTS.SERVICE.SAMPLE.QUEUE) define qlocal(SYSTEM.BROKER.CONTROL.QUEUE) end
Результаты работы программы amqsresa, получающей сообщения от
(рис 10.4) Результаты работы программы подписчикаОбъем программы
Подключение к менеджеру брокера (MQCONN )
Открытие очереди потока брокера и подписчика
(MQOPEN )
Генерация MQRFH и подписка на все события
Ожидание появления сообщений в очереди
подписчика (до 3 минут)
Извлечение из очереди MQGET всех публикаций
Отбор публикации по теме "Sport/Soccer/*"
Обработка и отображение результатов публикации
в зависимости от событий (начало матча,
конец матча, изменение счета)
Выход из цикла по концу матча или по таймеру
Закрытие очереди потока брокера и подписчика
(MQCLOSE)
Отключение от менеджера брокера (MQDISC)
В комментарии к алгоритму следует отметить, что подписка на все события (в алгоритме выделено курсивом) вряд ли объяснима, скорее это сделано в учебных целях.
В заключение лекции можно привести график времени доставки публикации (мсек) в зависимости от количества
(рис 10.5) Зависимость времени доставки публикации от количества подписчиков Механизм
Поставщика информации принято называть
Информация, посылаемая как сообщение, характеризуется темой.
(рис 10.1) Схема взаимодействия брокеров, издателей и подписчиковОбщий подход работы механизма
Сценарий
Сценарий
Сценарий
Механизм
Для работы с механизмом
Регион/СекторРынка/Компания
Конкретные темы выглядят следующим образом:
NewYork/InformationTechnology/IBM
Этот синтаксис разрешает в строках иметь символы:
* - означающий ноль или любое количество
произвольных символов;
? – означающий один произвольный символ.
Данный синтаксис позволит в дальнейшем определить, например, такие темы для подписки на рынке акций как:
* получение всех цен на акции
со всех рынков мира;
London/* получение всех цен на акции
с Лондонской биржи;
NewYork/Banks/* получение цен акций всех
нью-йоркских банков;
*/*/IBM получение всех цен на акции
компании ИБМ на всех рынках мира.
Символ ? в строках, например, 'ABC%*D' означает 'ABC*D'.
В числе основных функций для работы с механизмом
RegPub |
Регистрация издателя (Register Publisher) |
RegSub |
Регистрация |
Publish |
Публикация |
ReqUpdate |
Запрос от издателя к |
DeletePub |
Удаление публикации (Delete Publication) |
DeregPub |
Отказ от регистрации издателя (Deregister Publisher) |
DeregSub |
Отказ от регистрации |
Более подробно эти функции будут рассмотрены на примере ниже, а описание их форматов можно найти в документации , .
Для осуществления работы издателя и
typedef struct tagMQRFH {
MQCHAR4 StrucId;
/* Идентификатор структуры */
MQLONG Version;
/* Номер версии структуры */
MQLONG StrucLength;
/* Общая длина MQRFH */
MQLONG Encoding;
/* Data encoding */
MQLONG CodedCharSetId;
/* Идентификатор множества перекодировки */
MQCHAR8 Format;
/* Имя формата */
MQLONG Flags;
/* Флаги */
} MQRFH;
Этот заголовок также включает строку NameValueString, в которой приложение публикации или подписки помещает команды, которые должен выполнить StrucLength в заголовке определяет длину структуры заголовка, включительно с переменной длины NameValueString в конце структуры. Поля Encoding, CodedCharSetId и Format описывают структуру данных в заголовке публикации.
Задача создания собственных приложений Издателя и
(рис 10.2) Модель Издатель/ПодписчикПриложение Издателя должно публиковать символьные строки, задаваемые пользователем на определенные темы. Следовательно, такая программа с именем publisher может иметь следующий формат запуска:
publisher Topic Command QMgrName PubQueue
где Topic – тема публикации, символьная строка длиной не более 256 байт; Command – команды Издателя для RegPub, Publish, ReqUpdate, DeletePub, DeregPub ), PubQueue - очередь издателя. Текст публикации вводится в командном окне с запущенной программой.
Приложение subscriber может иметь формат запуска:
subscriber Topic QMgrName SubQueue
где Topic – тема SubQueue - очередь Topic отображается в командном окне с запущенной программой subscriber.
Для работы приложений требуется создание очередей: одна у Издателя и одна у Publisher_queue и Subscriber_queue. Publisher_queue. Subscriber_queue. Создание этих очередей осуществляется простой командой :
runmqsc QMgrName define qlocal(имя очереди) end
Следует рассмотреть потоки сообщений и то, как сообщения по определенной теме находят своего
SYSTEM.BROKER.CONTROL.QUEUE, содержащие TestTopic. SYSTEM.BROKER.*.TestTopic и оно попадает в потоковую очередь SYSTEM.BROKER.DEFAULT.STREAM
Сообщения
Также как и сообщения
Основные шаги, из которых складывается создание Приложения издателя. Таких шагов для программы publisher будет 7.
Шаг 1. Для того чтобы положить сообщение в очередь потоков для помещения сообщений.
Шаг 2. Для публикации необходимо сформировать MQRFH структуру (см. руководство по MQRFH должны быть изменены по сравнению со значениями по умолчанию перед публикацией или в тот момент, когда мы определим значения этих полей.
Шаг 3. Сразу за структурой MQRFH должна следовать символьная строка NameValueString. Указатель pNameValueString определяет начальную позицию NameValueString в теле сообщения.
Шаг 4. Содержимое строки NameValueString должно включать всю необходимую информацию о публикации и всю необходимую последовательность команд для формирования непрерывной символьной строки нашей публикации. Содержимое строки NameValueString позволяет
Шаг 5. Поле StrucLength должно содержать длину MQRFH структуры и сопровождающей его строки NameValueString. Длина MQRFH фиксирована и длина NameValueString – переменная. Значение StrucLength позволяет не применять разделитель конца строки в NameValueString, хотя и он может быть использован, если это необходимо приложению- MQRFH и NameValueString выравниваются на границе 16 байт автоматически.
Шаг 6. Данные о публикации (в данном случае строка данных для публикации) помещаются сразу после MQFRH и NameValueString структуры, указатель pUserData должен быть установлен на начальную позицию строки введенных данных о публикации.
Шаг 7. Осуществляется процедура тестирования работы программы. Программа компилируется, устраняются ошибки компиляции. Она вызывается, например, командой:
publisher TestTopic Publish QM_broker Publisher_queue
В командном окне программы publisher вводится сообщение по теме: Hello
На этом этапе еще нет приложения
Publisher_queue и исправить это; есть ли успешные сообщения от SYSTEM.DEFAULT.LOCAL.STREAM, используя MQ explorer или команду amqsbcg и убедиться, что в формате сообщения о публикации нет ошибок.Основные шаги, из которых складывается создание Приложения подписчика.
Шаг 1. Прежде всего, командой необходимо открыть очередь
Шаг 2. На этом шаге необходимо зарегистрировать интересы в виде темы и послать SendBrokerCommand, одним из аргументов которой является командная строка для помещения в NameValueString командного сообщения. Эта функция SendBrokerCommand должна быть аналогичной по коду, как и для приложения Издателя.
Шаг 3. Теперь, когда . Первое, что необходимо сделать при получении сообщения, это проверить, что оно в нужном формате MQRFH.
Шаг 4. После распознавания формата MQRFH сообщения, можно выделить порцию наиболее интересного сообщения. В данном случае нет необходимости смотреть на строку NameValueString так как интересуют данные, которые следуют за этой строкой, они задаются указателем pUserData и могут быть напечатаны на экране или выведены в файл.
Шаг 5. На этом последнем шаге необходимо отказаться от регистрации и, как это делалось раньше на шаге 2, вызвать функцию SendBrokerCommand, добавив соответствующую команду для отказа от регистрации. Теперь, так же как и раньше, надо скомпилировать код и затем исправить ошибки компиляции. После этого программа готова для работы, ее можно выполнить:
subscriber TestTopic QM_broker Subscriber_queue
Этот запуск должен отобразить информацию об успешной регистрации у publisher на тему TestTopic с сообщением, например, PublisherReady. В результате, в программе- PublisherReady.
Если этого не произойдет, то необходимо найти ошибку.
Получив успешно одну публикацию от одного издателя, можно расширить эксперимент и попробовать запустить много программ-издателей и программ-
Теперь после знакомства с технологией
После этого можно стартовать
strmqbrk -m QMgrName
Для отображения состояния dspmqbrk:
dspmqbrk -m QMgrName
В ответ появится следующее сообщение:
WebSphere MQ message broker for queue manager QMgrName running
Теперь
На каждом менеджере может быть стартован только один
В менеджере есть необходимые
runmqsc QMgrName display qlocal(SYSTEM.BROKER.*) end
Следует сразу отметить, что завершение работы endmqbrk перед окончанием работы менеджера: endmqbrk -m QMgrName.
Работу издателя можно продемонстрировать с помощью программы amqsgama, предложенной в SupportPacs MA0C в качестве теста. Эта программа издателя из перечня спортивных тем для подписки (табл.10.1) помещает на Sport/Soccer/Event/ - в таблице данная тема выделена курсивом) и проверяет ответы
спорт/футбол/* спорт/теннис/* спорт/баскетбол/* |
спорт/футбол/расписание игр спорт/футбол/события спорт/футбол/обзоры |
Формат запуска программы:
amqsgama TeamName1 TeamName2 QMgrName
Результаты работы программы, моделирующей случайным образом забивание голов той или иной командой, выглядит следующим образом (рис. 10.3):
(рис 10.3) Результаты работы программы издателяДля работы программы необходимо создание очереди: SAMPLE.BROKER.RESULTS.STREAM.
Именно в эту очередь поступают сообщения от издателя. Необходимо также, чтобы был запущен amqsgama, чтобы отобразить результаты игры полностью. Все используемые функции в программе служат для подключения к менеджеру
Блочная структура программы выглядит следующим образом.
Подключение к менеджеру брокера (MQCONN)
Открытие очереди потока брокера (MQOPEN)
Инициализация таймера матча
Генерация MQRFH для публикации события о
начале матча
Добавление имен команд в данные
Помещение публикации в очередь потоков
Начало цикла по времени матча:
засыпание на случайный период
попытка забить гол (50% вероятность)
генерация публикации о забитом голе
(RFH для ScoreUpdate)
случайный выбор команды, забившей гол
добавление имени команды в данные для
публикации
помещение публикации в очередь потоков
Окончание цикла по времени матча
Генерация MQRFH для публикации события о
конце матча
Добавление имен команд для публикации
Помещение публикации в очередь потоков
Закрытие очереди потока брокера (MQCLOSE)
Отключение от менеджера брокера (MQDISC)
Программа amqsgama имеет следующий код:
/*************************************************************************************/
/* Имя программы: AMQSGAMA */
/* Описание: Основанная на модели Publish/Subscribe программа */
/* моделирует результаты футбольного матча и */
/* отправляет их от издателя к брокеру */
/* Statement: Licensed Materials - Property of IBM */
/* SupportPac MA0E */
/* (C) Copyright IBM Corp. 1999 */
/*************************************************************************************/
#include <stdlib.h>
#include <stdio.h>
#include <string.h>
#include <time.h>
#include <cmqc.h> /* MQI */
#include <cmqpsc.h> /* MQI Publish/Subscribe */
#include <windows.h>
#if MQAT_DEFAULT == MQAT_WINDOWS_NT
#define msSleep(time) \
Sleep(time)
#elif MQAT_DEFAULT == MQAT_UNIX
#define msSleep(time) \
{ \
struct timeval tval; \
tval.tv_sec = time / 1000; \
tval.tv_usec = (time % 1000) * 1000; \
select(0, NULL, NULL, NULL, tval); \
}
#endif
#define STREAM "SAMPLE.BROKER.RESULTS.STREAM"
#define TOPIC_PREFIX "Sport/Soccer/Event/"
#define MATCH_STARTED "MatchStarted"
#define MATCH_ENDED "MatchEnded"
#define SCORE_UPDATE "ScoreUpdate"
#define MATCH_LENGTH 30000 /* 30 Second match length */
#define REAL_TIME_RATIO 333
#define AVERAGE_NUM_OF_GOALS 5
#define DEFAULT_MESSAGE_SIZE 512 /* Maximum buffer size for a message */
static const MQRFH DefaultMQRFH = {MQRFH_DEFAULT};
typedef struct
{
MQCHAR32 Team1;
MQCHAR32 Team2;
} Match_Teams, *pMatch_Teams;
void BuildMQRFHeader( PMQBYTE pStart
, PMQLONG pDataLength
, MQCHAR TopicType[] );
void PutPublication( MQHCONN hConn
, MQHOBJ hObj
, PMQBYTE pMessage
, MQLONG messageLength
, PMQLONG pCompCode
, PMQLONG pReason );
int main(int argc, char **argv)
{
MQHCONN hConn = MQHC_UNUSABLE_HCONN;
MQHOBJ hObj = MQHO_UNUSABLE_HOBJ;
MQLONG CompCode;
MQLONG Reason;
MQOD od = { MQOD_DEFAULT };
MQLONG Options;
PMQBYTE pMessageBlock = NULL;
MQLONG messageLength;
MQLONG timeRemaining;
MQLONG delay;
PMQCHAR pScoringTeam;
pMatch_Teams pTeams;
MQCHAR32 team1;
MQCHAR32 team2;
char QMName[MQ_Q_MGR_NAME_LENGTH+1] = "";
MQLONG randomNumber;
MQLONG ConnReason;
/* Проверка аргументов программы */
if( (argc < 3)||(argc > 4)||(strlen(argv[1]) > 31)||(strlen(argv[2]) > 31) )
{
printf("Usage: amqsgam team1 team2 <QManager>\n");
printf(" Maximum 31 characters per team name,\n");
printf(" no spaces or '\"' characters allowed.\n");
exit(0);
}
else
{
strcpy(team1, argv[1]);
strcpy(team2, argv[2]);
}
/* Использовать default queue manager или заданный в зависимости от наличия аргумена */
if (argc > 3) strcpy(QMName, argv[3]);
MQCONN( QMName, hConn, CompCode, ConnReason );
if( CompCode == MQCC_FAILED )
{
printf("MQCONN failed with CompCode %d and Reason %d\n", CompCode, ConnReason);
}
else if( ConnReason == MQRC_ALREADY_CONNECTED )
{
CompCode = MQCC_OK;
}
if( CompCode == MQCC_OK )
{
strncpy(od.ObjectName, STREAM, (size_t)MQ_Q_NAME_LENGTH);
Options = MQOO_OUTPUT + MQOO_FAIL_IF_QUIESCING;
MQOPEN( hConn, od, Options, hObj, CompCode, Reason );
if( CompCode != MQCC_OK )
{
printf("MQOPEN failed to open \"%s\"\nwith CompCode %d and Reason %d\n",
od.ObjectName, CompCode, Reason);
}
}
if( CompCode == MQCC_OK )
{
srand( (unsigned)(time( NULL ))
+ (unsigned)(team1[0] + team2[(strlen(team2) - 1)]) );
timeRemaining = MATCH_LENGTH;
messageLength = DEFAULT_MESSAGE_SIZE;
pMessageBlock = (PMQBYTE)malloc(messageLength);
if( pMessageBlock == NULL )
{
printf("Unable to allocate storage\n");
}
else
{
if( CompCode == MQCC_OK )
{
/* создание MQRFH для публикации о начале матча */
BuildMQRFHeader( pMessageBlock, messageLength, MATCH_STARTED );
pTeams = (pMatch_Teams)(pMessageBlock + messageLength);
strcpy(pTeams->Team1, team1);
strcpy(pTeams->Team2, team2);
messageLength += sizeof(Match_Teams);
printf("Match between %s and %s\n", team1, team2);
/* помещение сообщения (публикации) в очередь потока */
PutPublication( hConn, hObj, pMessageBlock, messageLength, CompCode, Reason );
if( CompCode != MQCC_OK )
{
printf("MQPUT failed with CompCode %d and Reason %d\n", CompCode, Reason);
}
else
{
/* Моделирование попытки забить один из 5 голов (каждые 30 сек) */
while( (timeRemaining > 0)(CompCode == MQCC_OK) )
{
randomNumber = rand();
delay = REAL_TIME_RATIO
+ ( (randomNumber * MATCH_LENGTH)
/ (RAND_MAX * AVERAGE_NUM_OF_GOALS));
if( delay > timeRemaining ) delay = timeRemaining;
msSleep(delay);
timeRemaining -= delay;
if( timeRemaining > 0 )
{
if( (randomNumber % 2) == 0 ) /* Шанс забить гол - 50 процентов */
{
messageLength = DEFAULT_MESSAGE_SIZE;
BuildMQRFHeader( pMessageBlock , messageLength, SCORE_UPDATE );
printf("GOAL! ");
pScoringTeam = (PMQCHAR)pMessageBlock + messageLength;
if( rand() < (RAND_MAX/2) )
{
strcpy(pScoringTeam, team1);
printf(team1);
}
else
{
strcpy(pScoringTeam, team2);
printf(team2);
}
printf(" scores after %d minutes\n", ((MATCH_LENGTH - timeRemaining)/REAL_TIME_RATIO));
messageLength += sizeof(MQCHAR32);
/* помещение сообщения о забитом голе в очередь потока */
PutPublication( hConn, hObj, pMessageBlock, messageLength, CompCode, Reason );
if( CompCode != MQCC_OK )
printf("MQPUT failed with CompCode %d and Reason %d\n", CompCode, Reason);
}
}
} /* конец цикла по времени окончания матча ( timeRemaining ) */
if( CompCode == MQCC_OK )
{
printf("Full time\n");
messageLength = DEFAULT_MESSAGE_SIZE;
BuildMQRFHeader( pMessageBlock , messageLength , MATCH_ENDED );
pTeams = (pMatch_Teams)(pMessageBlock + messageLength);
strcpy(pTeams->Team1, team1);
strcpy(pTeams->Team2, team2);
messageLength += sizeof(Match_Teams);
/* помещение сообщения о конце матча в очередь потока */
PutPublication( hConn, hObj, pMessageBlock, messageLength, CompCode, Reason );
if( CompCode != MQCC_OK )
printf("MQPUT failed with CompCode %d and Reason %d\n", CompCode, Reason);
}
}
}
free( pMessageBlock );
} /* end of else (pMessageBlock != NULL) */
}
if( hObj != MQHO_UNUSABLE_HOBJ )
{
MQCLOSE( hConn , hObj, MQCO_NONE, CompCode , Reason );
if( CompCode != MQCC_OK )
printf("MQCLOSE failed with CompCode %d and Reason %d\n", CompCode, Reason);
}
if( (hConn != MQHC_UNUSABLE_HCONN) (ConnReason != MQRC_ALREADY_CONNECTED) )
{
MQDISC( hConn, CompCode, Reason );
if( CompCode != MQCC_OK )
printf("MQDISC failed with CompCode %d and Reason %d\n", CompCode, Reason);
}
return(0);
}
/* end of main Function */
/* Function Name : BuildMQRFHeader */
void BuildMQRFHeader( PMQBYTE pStart, PMQLONG pDataLength, MQCHAR TopicType[] )
{
PMQRFH pRFHeader = (PMQRFH)pStart;
PMQCHAR pNameValueString;
memset((PMQBYTE)pStart, 0, *pDataLength);
memcpy( pRFHeader, DefaultMQRFH, (size_t)MQRFH_STRUC_LENGTH_FIXED);
memcpy( pRFHeader->Format, MQFMT_STRING, (size_t)MQ_FORMAT_LENGTH);
pRFHeader->CodedCharSetId = MQCCSI_INHERIT;
pNameValueString = (MQCHAR *)pRFHeader + MQRFH_STRUC_LENGTH_FIXED;
strcpy(pNameValueString, MQPS_COMMAND_B);
strcat(pNameValueString, MQPS_PUBLISH);
strcat(pNameValueString, MQPS_PUBLICATION_OPTIONS_B);
strcat(pNameValueString, MQPS_NO_REGISTRATION);
strcat(pNameValueString, MQPS_TOPIC_B);
strcat(pNameValueString, TOPIC_PREFIX);
strcat(pNameValueString, TopicType);
*pDataLength = MQRFH_STRUC_LENGTH_FIXED + ((strlen(pNameValueString)+15)/16)*16;
pRFHeader->StrucLength = *pDataLength;
}
/* Function Name : PutPublication */
void PutPublication( MQHCONN hConn, MQHOBJ hObj, PMQBYTE pMessage,
MQLONG messageLength, PMQLONG pCompCode, PMQLONG pReason )
{
MQPMO pmo = { MQPMO_DEFAULT };
MQMD md = { MQMD_DEFAULT };
memcpy(md.Format, MQFMT_RF_HEADER, (size_t)MQ_FORMAT_LENGTH);
md.MsgType = MQMT_DATAGRAM;
md.Persistence = MQPER_PERSISTENT;
pmo.Options |= MQPMO_NEW_MSG_ID;
MQPUT( hConn, hObj, md, pmo, messageLength, pMessage, pCompCode, pReason );
}
В качестве комментария следует отметить, что функция BuildMQRFHeader формирует значения по умолчанию для заголовка MQRFH, устанавливает параметры format и CCSID пользовательских данных. В строку NameValueString добавляются команды, тема и опции для публикации и она выравнивается на 16-ти байтовую границу. StrucLength в MQRFH устанавливается как общая длина. Входными параметрами функции являются pStart – начало блока сообщения, TopicType[] – строка с именем темы. Входным и выходным параметром одновременно является pDataLength – размер блока сообщения при входе и размер выходного блока информации.
Функция PutPublication формирует сообщение для вывода в очередь брокера с помощью команды . Входными параметрами функции являются hConn – идентификатор менеджера для команды MQHCONN, hObj – идентификатор очереди, pMessage – идентификатор на начало блока сообщения, messageLength – длина данных в сообщении. Выходными параметрами функции являются pCompCode и pReason – коды завершения команды .
Работу amqsresa из состава SupportPacs MA0C, которая подписывается у
amqsresa QMgrName
где QmgrName – имя менеджера очередей, на котором запущен
Необходимая тема подписки ( TOPIC = "Sport/Soccer/*" ) и очереди для
Для работы программы необходимо создание очередей:
runmqsc QMgrName define qlocal(SAMPLE.BROKER.RESULTS.STREAM) define qlocal(RESULTS.SERVICE.SAMPLE.QUEUE) define qlocal(SYSTEM.BROKER.CONTROL.QUEUE) end
Результаты работы программы amqsresa, получающей сообщения от
(рис 10.4) Результаты работы программы подписчикаОбъем программы
Подключение к менеджеру брокера (MQCONN )
Открытие очереди потока брокера и подписчика
(MQOPEN )
Генерация MQRFH и подписка на все события
Ожидание появления сообщений в очереди
подписчика (до 3 минут)
Извлечение из очереди MQGET всех публикаций
Отбор публикации по теме "Sport/Soccer/*"
Обработка и отображение результатов публикации
в зависимости от событий (начало матча,
конец матча, изменение счета)
Выход из цикла по концу матча или по таймеру
Закрытие очереди потока брокера и подписчика
(MQCLOSE)
Отключение от менеджера брокера (MQDISC)
В комментарии к алгоритму следует отметить, что подписка на все события (в алгоритме выделено курсивом) вряд ли объяснима, скорее это сделано в учебных целях.
В заключение лекции можно привести график времени доставки публикации (мсек) в зависимости от количества
(рис 10.5) Зависимость времени доставки публикации от количества подписчиков Для получения официальных документов о завершении программы дополнительного профессионального образования (удостоверения о повышении квалификации, дипломов о профессиональной переподготовке и MBA) необходимо предоставить:
Внимание! Вы можете не заказывать доставку бумажной версии официального документы, а скачать его в электронном виде и распечатать самостоятельно. Информация о выданном документе в течение 1 месяца загружается в Федеральную информационную систему «Федеральный реестр сведений о документах об образовании и (или) о квалификации, документах об обучении» - ФИС ФРДО.
Доступ на новый сайт осуществляется с использованием адреса электронной почты, который был указан вами при регистрации на "старом". Мы постарались перенести все ваши данные с прежнего ресурса, однако не исключена вероятность потери части информации.
При возникновении проблемы со входом, воспользуйтесь функцией сброса пароля
Если вы обнаружите несоответствия, пожалуйста, сообщите нам.