Первая программа будет достаточно простая и реализует так называемую модель "один к одному" или "точка-точка". Эта программа предназначена для чтения сообщений из очереди 1, записи их в очередь 2 и лог-файл на диске. Эта программа имеет практическое значение. Достаточно часто необходимо иметь файл переданных сообщений за определенный период времени, чтобы быстро ответить на вопрос "Было ли передано сообщение с такими идентификационными параметрами в теле сообщения:…"?
Автору приходилось сталкиваться с "плохим" стилем rewriter.exe и файл инициализации rewriter.ini, в котором 1-я строка – имя очереди для чтения, 2-я строка – имя очереди для записи, 3-я строка – имя лог-файла, как показано ниже.
QUEUE_INPUT QUEUE_OUTPUT C:\TEMP\rewriter.log
Разрабатываемая программа может быть представлена в следующей последовательности псевдокода:
MQCONN MQOPEN --> цикл чтения сообщений | (на основе gmo.WaitInterval): | MQGET | MQPUT |-- конец цикла MQCLOSE MQDISC
Ниже приводится листинг программы , , в случае проблем с
/* Листинг программы rewriter */
/**********************************************************************/
/* Program name: Rewriter */
/* Description: Rewriter C program pass messages to output queue */
/* Function: */
/* Rewriter is a sample C program to demonstrate the main MQI calls; */
/* each message is copied from the input queue to the output */
/* queue, and sends a report to the log file */
/**********************************************************************/
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
#include <signal.h>
#include <io.h>
/* includes for MQI */
#include <cmqc.h>
char queue1[48] = "";
char queue2[48] = "";
char logfilename[48] = "";
char logfilename2[48] = "";
char buf[48];
int queuenamelen;
time_t tmr;
FILE * fp;
FILE *fptr;
void cntrl_c_handler(int sig);
/* Declare MQI structures needed */
MQOD odG = {MQOD_DEFAULT}; /* Object Descriptor for GET */
MQOD odP = {MQOD_DEFAULT}; /* Object Descriptor for PUT */
MQOD odI = {MQOD_DEFAULT}; /* Object Descriptor for InitQ */
MQOD odR = {MQOD_DEFAULT}; /* Object Descriptor for report */
MQMD md = {MQMD_DEFAULT}; /* Message Descriptor */
MQGMO gmo = {MQGMO_DEFAULT}; /* get message options */
MQPMO pmo = {MQPMO_DEFAULT}; /* put message options */
MQTMC2 *trig; /* trigger message structure */
MQCHAR48 QManager; /* queue manager name */
MQHCONN Hcon; /* connection handle */
MQHOBJ Hobj; /* object handle, server queue */
MQHOBJ Hinq; /* handle for MQINQ */
MQHOBJ Hout; /* handle for MQPUT */
MQLONG O_options; /* MQOPEN options */
MQLONG C_options; /* MQCLOSE options */
MQLONG CompCode; /* completion code */
MQLONG Reason; /* reason code */
MQLONG CReason; /* reason code (MQCONN) */
MQBYTE buffer[8001]; /* message buffer */
MQLONG buflen; /* buffer length */
MQLONG messlen; /* message length received */
MQLONG Select[1]; /* attribute selectors */
MQLONG SelectValue[1]; /* value attribute selectors */
MQLONG char_count;
int main(int argc, char **argv)
{
strcpy(QManager, ""); /* Работаем с менеджером очередей по умолчанию */
if ( (fptr=fopen ("rewriter.ini","r" )) == NULL )
{printf("Cannot open rewriter.ini file" ); exit(1); }
else{ /* Открываем ini-файл и присваиваем значения переменным */
fgets(queue1, 48, fptr);
queuenamelen = strlen(queue1) - 1;
queue1[queuenamelen] = ' ';
fgets(queue2, 48, fptr);
queuenamelen = strlen(queue2) - 1;
queue2[queuenamelen] = ' ';
fgets(logfilename, 48, fptr);
queuenamelen = strlen(logfilename) - 1;
logfilename[queuenamelen] = ' ';
tmr = time(NULL);
strcpy ( buf, ctime(tmr));
buf[strlen(buf)-1]=0; // переход на новую строку
strncat (logfilename, buf,10);
strcpy(odG.ObjectName, queue1);
strcpy(odP.ObjectName, queue2);
fclose (fptr);
}
MQCONN(QManager, Hcon, CompCode, CReason);
if (CompCode == MQCC_FAILED)
{
printf("MQCONN ended with reason code %ld\n", CReason);
exit(CReason);
}
O_options = MQOO_INPUT_SHARED + MQOO_FAIL_IF_QUIESCING;
MQOPEN(Hcon, odG, O_options, Hobj, CompCode, Reason); /* открываем очередь для чтения - odG */
if (Reason != MQRC_NONE) { printf("MQOPEN (input) ended with reason code %ld\n", Reason); }
if (CompCode == MQCC_FAILED) { exit(Reason); }
O_options = MQOO_OUTPUT + MQOO_FAIL_IF_QUIESCING;
MQOPEN(Hcon, odP, O_options, Hout, CompCode, Reason); /* открываем очередь для записи - odP */
if (Reason != MQRC_NONE) { printf("MQOPEN (output) ended with reason code %ld\n", Reason); }
if (CompCode == MQCC_FAILED) { exit(Reason); }
fp = fopen (logfilename,"a");
if ( fp==NULL ){ printf("Cannot open log file %s\n", logfilename); }
printf("Rewriter(C) sending messages from %s to %s and to log-file %s \n
",odG.ObjectName, odP.ObjectName, logfilename);
/*****************************************************************************/
/* Читаем сообщения из QUEUE_INPUT и пишем в QUEUE_OUTPUT */
/* до тех пор пока не встретим сообщение об ошибке */
/*****************************************************************************/
buflen = sizeof(buffer) - 1;
while (CompCode == MQCC_OK)
{
gmo.Options = MQGMO_ACCEPT_TRUNCATED_MSG + MQGMO_WAIT;
gmo.WaitInterval = 3000; /* Ожидаем новые сообщения 3 секунды */
//gmo.WaitInterval = MQWI_UNLIMITED;
memcpy(md.MsgId, MQMI_NONE, sizeof(md.MsgId));
memcpy(md.CorrelId, MQMI_NONE, sizeof(md.CorrelId));
MQGET(Hcon, Hobj, md, gmo, buflen, buffer, messlen, CompCode, Reason);
if ((CompCode == MQCC_OK) || (CompCode == MQCC_WARNING))
{
buffer[messlen] = '\0'; /* заносим символ конец строки в буфер с прочитанным сообщением */
buflen = messlen;
MQPUT(Hcon, Hout, md, pmo, buflen, buffer, CompCode, Reason);
if ((CompCode == MQCC_OK) || (CompCode == MQCC_WARNING))
{
tmr = time(NULL);
strcpy ( buf, ctime(tmr));
buf[strlen(buf)-1]=0; // переход к новой строке
Reason = fprintf( fp, "%s: %s\n", buf, buffer );
}
} /* конец обработки входного сообщения */
} /* конец цикла чтения/записи сообщений функциями MQGET, MQPUT */
fclose (fp);
C_options = 0; /* нет никаких опций при закрытии */
MQCLOSE(Hcon, Hobj, C_options, CompCode, Reason); /* закрываем очередь для чтения */
if (Reason != MQRC_NONE)
{printf("MQCLOSE (input) ended with reason code %ld\n", Reason); }
MQCLOSE(Hcon, Hout, C_options, CompCode, Reason); /* закрываем очередь для записи */
if (Reason != MQRC_NONE)
{printf("MQCLOSE (output) ended with reason code %ld\n", Reason); }
if (CReason != MQRC_ALREADY_CONNECTED)
{
MQDISC(Hcon, CompCode, Reason);
if (Reason != MQRC_NONE){ printf("MQDISC ended with reason code %ld\n", Reason); }
}
return(0);
}
В данной версии мы выходим из цикла программы по опции gmo.WaitInterval = 3000, когда ожидаем сообщение в очереди в течении 3 сек, а его там нет (опция gmo.WaitInterval работает быстрее, чем если бы мы опрашивали очередь по собственному временному циклу). Другой вариант программы может быть таким. Задаем gmo.WaitInterval = MQWI_UNLIMITED ; что соответствует gmo.WaitInterval= -1. Программа будет крутиться "бесконечно" до тех пор, пока мы не остановим её принудительно, например, нажатием клавиш CNTRL_C (стандартный останов). В этом случае нужно добавить обработчик прерываний по нажатию CNTRL_C потому, что при таком выходе объекты очереди останутся не закрытыми и , в противном случае лог-файл не будет формироваться. Следует отметить, что размер массива buffer ограничивает длину сообщения 8Кб и при появлении сообщений большей длины следует увеличить размер буфера.
Программа rewriter.exe работает достаточно быстро и сравнительные скорости работы данного алгоритма при длине сообщения 1Кб на компьютере INTEL Pentium 1.8Ггц приведены в таблице ниже.
| Язык программы\тип очереди | Not Persistent | Persistent |
|---|---|---|
| 1000 сооб/сек | 400 сооб/сек | |
| Visual Basic 6.0 | 200 сооб/сек | 140 сооб/сек |
Увеличение длины сообщения не ведет к пропорциональному уменьшению скорости. Эти исследования читатель может проделать самостоятельно. Реальные приложения, работающие с базами данных, имеют скорость обработки сообщений в 3-4 раза меньше.
Возвращаясь к вопросу о стилях rewriter.ini файле приведет к зависанию программы и мучительному поиску причин такого зависания, не говоря о других более сложных ситуациях, например, когда очередь открыта эксклюзивно другим приложением.
Для версии программы gmo.WaitInterval = MQWI_UNLIMITED полезно сделать вывод на экран передаваемых сообщений, чтобы наблюдать динамику работы созданного интерфейса. Таких улучшений может быть достаточно много и мы рассмотрим две достаточно полезные модификации.
Программа rewriter может вызываться как
/* Код для вызова rewriter.exe
как MQSeries-триггер */
int main(int argc,
char **argv)
{ if (argc > 1)
{trig = (MQTMC2*)argv[1];
strncpy(odG.ObjectName,
trig->QName,
MQ_Q_NAME_LENGTH);
strncpy(queue1,
trig->QName,
MQ_Q_NAME_LENGTH);
strncpy(QManager,
trig->QMgrName,
MQ_Q_MGR_NAME_LENGTH);
strncpy(odP.ObjectName,
trig->UserData,
MQ_PROCESS_USER_DATA_LENGTH);
strncpy(queue2, trig->UserData,
MQ_PROCESS_USER_DATA_LENGTH);
strncpy(logfilename2,
trig->EnvData, 48);
}
Возможная модификация этого варианта - программа rewriter может вызываться с передачей параметров через командную строку и эту модификацию читатель может проделать самостоятельно.
rewriter может быть модифицирована в программу разветвитель mqsplitter.exe: чтение сообщений из очереди 1 и запись их в очередь 2, в очередь 3 и лог-файл на диске.Можно сделать программу mqsplitter.exe на разных языках, например, на Visual Basic 6.0 с интерфейсом, показанным на рис.9.1, и сравнить производительность программ на разных языках, реализующих один и тот же алгоритм. Такая задача будет хорошим лабораторным практикумом.
(рис 9.1) Интерфейс программы mqsplitter на VB6Модификацию программы mqsplitter.exe читателю предлагается сделать самостоятельно и одновременно проверить идею создания "вечного двигателя". Для создания "вечного двигателя" понадобиться изменить исходный mqsplitter.ini файл следующим образом:
|
-> |
|
Если в очереди QUEUE_INPUT будет хотя бы одно сообщение, то ваша программа будет посылать сообщения в очередь QUEUE_OUTPUT1 до тех пор, пока не будет остановлена. Можно заложить в ini-файл программы пятый параметр: время опроса очереди, измеряемое в миллисекундах. Такая программа окажется весьма полезной при тестировании интерфейсов.
Задачи
В этой задаче для нас интересна технология такого перехода. Совершенно очевидно, что переход от одной автоматизированной
В качестве конкретной задачи рассмотрим создание интерфейса по передаче клиентов из системы
Существует разные варианты решения поставленной задачи:
Первый путь ясен в технической реализации. С одной стороны при каждом обновлении таблицы client в АБС1 программа-обработчик срабатывает как триггер базы данных и помещает результаты оператора update в очередь, они приходят на АБС2, где своя программа-обработчик срабатывает как триггер очереди и помещает сообщение в таблицу client АБС2.
Блок 1 MQCONN; MQOPEN; Блок 2 MQBEGIN; MQGET; Блок 3 Begin Tran Блок 4 UPDATE CLIENT SET ... WHERE ... Блок 5 If Error = 0 then Commit Tran else Rollback Tran; Блок 6 If Error = 0 then MQCMIT else MQBACK; Блок 7 MQCLOSE; MQDISC;
В этой программе
Тем не менее, при небольшом числе интерфейсов рекомендуется вариант 1, как наиболее экономический, к рассмотрению которого мы и переходим. WebSphere BI Broker будет рассмотрен в лекции 12. К сожалению, рассмотрение средств
Сообщения
Наша очередная задача: на сервере 1 прочитать сообщение из входной очереди, положить её в очередь для отправки на сервер 2 как сообщение-запрос и дождаться прихода сообщения-ответа, как это показано на рис.9.2. Все это необходимо оформить в виде
(рис 9.2) Структура объектов WebSphere MQИтак, последовательность псевдокода представляется следующим образом (обратите внимание на блок 5 и опции MQMD):
Блок 1 MQCONN Блок 2 MQOPEN Блок 3 MQBEGIN Блок 4 MQGET (Input_queue) Блок 5 MQPUT (Output_queue, MQMD.MsgType = MQMT_REQUEST, MQMD.ReplyToQ = Reply_queue) Блок 6 MQGET (Reply_queue) Блок 7 If Reply time < 10 sec then MQCMIT else MQBACK; Блок 8 MQCLOSE Блок 9 MQDISC
Назовем нашу программу transmit.exe и файл инициализации transmit.ini, в котором 1-я строка – имя очереди для чтения, 2-я строка – имя очереди для записи, 3-я строка – имя очереди для ответа, 4-я строка – время ожидания ответа Reply_time = 3000мсек, как показано ниже.
QUEUE_INPUT QUEUE_OUTPUT QUEUE_REPLY 3000
Тип очереди Output_queue –
Ниже приводится листинг программы transmit.cpp для Microsoft Visual C++ ver.6.0. Для каждого сообщения и ) и об этом подробнее в лекции 11.
/* Листинг программы transmit */
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
#include <io.h>
#include <cmqc.h>
char queue_input[48] = "";
char queue_output[48] = "";
char queue_reply[48] = "";
char reply_time[48] = "";
char buf[48];
time_t tmr;
int queuenamelen;
FILE *fptr;
int main(int argc, char **argv)
{
MQOD odG = {MQOD_DEFAULT};
MQOD odP = {MQOD_DEFAULT};
MQOD odR = {MQOD_DEFAULT};
MQOD odI = {MQOD_DEFAULT};
MQMD md = {MQMD_DEFAULT};
MQBO mbo = {MQBO_DEFAULT};
MQGMO gmo = {MQGMO_DEFAULT};
MQPMO pmo = {MQPMO_DEFAULT};
MQCHAR48 QManager;
MQHCONN Hcon;
MQHOBJ Hobj;
MQHOBJ Hout;
MQHOBJ Hrep;
MQLONG O_options;
MQLONG C_options;
MQLONG CompCode;
MQLONG Reason;
MQLONG CReason;
MQBYTE buffer[8001];
MQLONG buflen;
MQLONG replylen;
MQLONG messlen;
static MQBYTE24 LastMsgId;
if ( (fptr=fopen ("transmit.ini","r" )) == NULL )
{printf("Cannot open transmit.ini file" );
exit(1);
}
else{
fgets(queue_input, 48, fptr);
queuenamelen = strlen(queue_input) - 1;
queue_input[queuenamelen] = ' ';
strcpy(odG.ObjectName, queue_input);
fgets(queue_output, 48, fptr);
queuenamelen = strlen(queue_output) - 1;
queue_output[queuenamelen] = ' ';
strcpy(odP.ObjectName, queue_output);
fgets(queue_reply, 48, fptr);
queuenamelen = strlen(queue_reply) - 1;
queue_reply[queuenamelen] = ' ';
strcpy(odR.ObjectName, queue_reply);
fgets(reply_time, 48, fptr);
queuenamelen = strlen(reply_time) - 1;
reply_time[queuenamelen] = ' ';
fclose (fptr);
}
strcpy(QManager, ""); /* Работаем с менеджером очередей по умолчанию */
MQCONN(QManager, Hcon, CompCode, CReason);
if (CompCode == MQCC_FAILED)
{
printf("MQCONN to %s ended with reason code %ld\n", QManager, CReason);
exit(CReason);
}
O_options = MQOO_BROWSE + MQOO_INPUT_SHARED ;
MQOPEN(Hcon, odG, O_options, Hobj, CompCode, Reason);
if (Reason != MQRC_NONE) { printf("MQOPEN %s ended with reason code %ld\n", queue_input, Reason); }
if (CompCode == MQCC_FAILED) { exit(Reason); }
O_options = MQOO_BROWSE + MQOO_INPUT_SHARED ;
MQOPEN(Hcon, odR, O_options, Hrep, CompCode, Reason);
if (Reason != MQRC_NONE) { printf("MQOPEN %s ended with reason code %ld\n", queue_reply, Reason); }
if (CompCode == MQCC_FAILED) { exit(Reason); }
O_options = MQOO_OUTPUT ;
MQOPEN(Hcon, odP, O_options, Hout, CompCode, Reason);
if (Reason != MQRC_NONE) { printf("MQOPEN %s ended with reason code %ld\n", queue_output, Reason); }
if (CompCode == MQCC_FAILED) { exit(Reason); }
while (CompCode == MQCC_OK)
{
buflen = sizeof(buffer) - 1;
memcpy(md.MsgId, MQMI_NONE, sizeof(md.MsgId));
memcpy(md.CorrelId, MQCI_NONE, sizeof(md.CorrelId));
gmo.Options = MQGMO_ACCEPT_TRUNCATED_MSG + MQGMO_WAIT + MQGMO_SYNCPOINT;
gmo.WaitInterval = 3000 ;
MQBEGIN (Hcon, mbo, CompCode, Reason);
MQGET(Hcon, Hobj, md, gmo, buflen, buffer, messlen, CompCode, Reason);
//if (Reason != MQRC_NONE) { printf("MQGET from %s ended with reason code %ld\n", queue_input, Reason); }
if ((CompCode == MQCC_OK) || (CompCode == MQCC_WARNING))
{
buffer[messlen] = '\0'; /* заносим символ конец строки в буфер с прочитанным сообщением */
buflen = messlen;
md.MsgType = MQMT_REQUEST;
md.Report = MQRO_EXCEPTION_WITH_DATA;
strncpy(md.ReplyToQ, queue_reply, MQ_Q_NAME_LENGTH);
memcpy(md.Format, MQFMT_STRING, MQ_FORMAT_LENGTH);
MQPUT(Hcon, Hout, md, pmo, buflen, buffer, CompCode, Reason);
if (Reason != MQRC_NONE)
{
printf("MQPUT to %s ended ended unsuccessfully with reason code %ld CompCode %ld\n", queue_output, Reason, CompCode );
MQBACK( Hcon, CompCode, Reason ) ;
CompCode = MQCC_FAILED ;
}
else
{
while (CompCode != MQCC_FAILED)
{
/** осуществляется проверка queue_reply **/
gmo.Options = MQGMO_ACCEPT_TRUNCATED_MSG + MQGMO_WAIT ;
gmo.WaitInterval = 3000 ;
memcpy(md.MsgId, MQMI_NONE, sizeof(md.MsgId));
memcpy(md.CorrelId, MQCI_NONE, sizeof(md.CorrelId));
MQGET(Hcon, Hrep, md, gmo, buflen, buffer, replylen, CompCode, Reason);
if (CompCode != MQCC_FAILED)
{
if (md.MsgType == MQMT_REPLY) /* report feedback */
{ printf("Transaction % s=> %s successfully: %s\n", queue_input, queue_output, buffer);
MQCMIT( Hcon, CompCode, Reason ) ;
}
else
{
printf("Transaction % s=> %s successfully, REPLY message not deliver, reason code %ld CompCode %ld\n", queue_input, queue_output, queue_reply, Reason, CompCode );
MQBACK( Hcon, CompCode, Reason ) ;
CompCode = MQCC_FAILED ;
}
}
if (Reason == MQRC_NO_MSG_AVAILABLE)
{
printf("Transaction % s=> %s UNsuccessfully, REPLY message not deliver\n", queue_input, queue_output );
MQBACK( Hcon, CompCode, Reason ) ;
CompCode = MQCC_FAILED ;
}
}
}
}
}
C_options = 0;
MQCLOSE(Hcon, Hobj, C_options, CompCode, Reason);
if (Reason != MQRC_NONE){printf("MQCLOSE %s ended with reason code %ld\n", queue_input, Reason); }
MQCLOSE(Hcon, Hout, C_options, CompCode, Reason);
if (Reason != MQRC_NONE){printf("MQCLOSE %s ended with reason code %ld\n", queue_output, Reason); }
MQCLOSE(Hcon, Hrep, C_options, CompCode, Reason);
if (Reason != MQRC_NONE){printf("MQCLOSE %s ended with reason code %ld\n", queue_reply, Reason); }
MQDISC(Hcon, CompCode, Reason);
if (Reason != MQRC_NONE){ printf("MQDISC ended with reason code %ld\n", Reason); }
return(0);
}
По тексту программы следует дать комментарии. Наличие опции
gmo.Options = MQGMO_SYNCPOINT;
подразумевает, что команда может не указываться. Операторы
md.MsgType = MQMT_REQUEST;
strncpy(md.ReplyToQ,
queue_reply,
MQ_Q_NAME_LENGTH);
определяют тип сообщения и очередь ответа, заданную в QUEUE_REPLY.
На очередь QUEUE_OUTPUT (или на удаленную очередь на другом менеджере) должна быть навешена программа-триггер, который возвращает сообщения типа . Если QUEUE_REPLY, то QUEUE_INPUT. такой же, как и исходного сообщения. В данной версии программы в целях упрощения отладки не проверяется это условие и читателю предлагается самостоятельно дописать этот фрагмент кода после отладки текущей версии программы. Работа с и будет рассмотрена подробнее в лекции 11.
Программу-триггер, которая "навешивается" на очередь QUEUE_OUTPUT (или на удаленную очередь) для формирования md.MsgType = MQMT_REPLY; ), читателю также предлагается сделать самостоятельно.
На данном примере мы познакомились с
положить сообщения в эти очереди. После помещает сообщения во все эти очереди, используя этот единственный
В версии
Рассмотрим этот механизм на примере задачи, когда distlist.exe, файл с текстом сообщения distlist.dat и файл инициализации distlist.ini, в котором 1-я строка – имя менеджера, 2-я и последующие строки – имена очередей, как показано ниже.
QM_ ALFA Queue_ Moscow Queue_ Kiev Queue_ Alma-Ata Queue_ SPetersburg Queue_ Novosibirsk Queue_ Saratov //last string must be blank
(рис 9.3) Механизм Distribution List для WebSphere MQНиже приводится листинг программы distlist.cpp для Microsoft Visual C++ ver.6.0.
/* Листинг программы distlist */
/* Program name: Distlist */
/* Description: Distlist C program pass messages to output queues */
/* by Distribution list for indicated Queue Manager */
/* distlist.ini file give list of queue and distlist.dat give file */
/* of message which copied to the output queue */
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <io.h>
#include <time.h>
#include <cmqc.h>
char queue[1000][48] ;
char buf[48];
int queuenamelen;
time_t tmr;
FILE * fp;
FILE *fptr;
static void print_usage(void);
static void print_responses( char * comment, PMQRR pRR, MQLONG NumQueues, PMQOR pOR);
int main(int argc, char **argv)
{
typedef enum {False, True} Bool;
MQOD od = {MQOD_DEFAULT}; /* Object Descriptor */
MQMD md = {MQMD_DEFAULT}; /* Message Descriptor */
MQPMO pmo = {MQPMO_DEFAULT}; /* put message options */
MQHCONN Hcon; /* connection handle */
MQHOBJ Hobj; /* object handle */
MQLONG O_options; /* MQOPEN options */
MQLONG C_options; /* MQCLOSE options */
MQLONG CompCode; /* completion code */
MQLONG OpenCode; /* MQOPEN completion code */
MQLONG Reason; /* reason code */
MQCHAR48 QManager; /* queue manager name */
MQLONG buflen; /* buffer length */
char buffer[101]; /* message buffer */
MQLONG Index ; /* Index into list of queues */
MQLONG NumQueues ; /* Number of queues */
PMQRR pRR=NULL; /* Pointer to response records */
PMQOR pOR=NULL; /* Pointer to object records */
Bool DisconnectRequired=False;/* Already connected switch */
Bool Connected=False; /* Connect succeeded switch */
typedef struct
{
MQBYTE24 MsgId;
MQBYTE24 CorrelId;
} PutMsgRec, *pPutMsgRec;
pPutMsgRec pPMR=NULL; /* Pointer to put msg records */
MQLONG PutMsgRecFields=MQPMRF_MSG_ID | MQPMRF_CORREL_ID;
/* Open ini file and setting value */
if ( (fptr=fopen ("distlist.ini","r" )) == NULL )
{printf("Cannot open distlist.ini file" );
print_usage();
exit(1); }
else{
fgets(QManager, 48, fptr);
queuenamelen = strlen(QManager) - 1;
QManager[queuenamelen] = ' ';
NumQueues = 0;
while (queuenamelen != 0)
{
fgets(queue[NumQueues], 48, fptr);
queuenamelen = strlen(queue[NumQueues]) - 1;
queue[NumQueues][queuenamelen] = ' ';
NumQueues++;
}
}
fclose (fptr);
--NumQueues; /* NumQueues - Number of Queue name */
/* Allocate response records, object records and put message records */
pRR = (PMQRR)malloc( NumQueues * sizeof(MQRR));
pOR = (PMQOR)malloc( NumQueues * sizeof(MQOR));
pPMR = (pPutMsgRec)malloc( NumQueues * sizeof(PutMsgRec));
if((NULL == pRR) || (NULL == pOR) || (NULL == pPMR))
{
printf("%s(%d) malloc failed\n", __FILE__, __LINE__);
exit(4);
}
/* Use parameters as the name of the target queues */
for( Index = 0 ; Index < NumQueues ; Index ++)
{
strncpy( (pOR+Index)->ObjectName, queue[Index], (size_t)MQ_Q_NAME_LENGTH);
strncpy( (pOR+Index)->ObjectQMgrName, QManager, (size_t)MQ_Q_MGR_NAME_LENGTH);
}
for( Index = 0 ; Index < NumQueues ; Index ++)
{
MQCONN((pOR+Index)->ObjectQMgrName, Hcon, ((pRR+Index)->CompCode), ((pRR+Index)->Reason));
if ((pRR+Index)->CompCode == MQCC_FAILED)
{
continue;
}
if ((pRR+Index)->CompCode == MQCC_OK)
{
DisconnectRequired = True ;
}
Connected = True;
break ;
}
/* Print any non zero responses */
print_responses("MQCONN", pRR, Index, pOR);
/* Print If failed to connect to queue manager then exit. */
if( False == Connected )
{
printf("Unable to connect to queue manager\n");
exit(3) ;
}
if ( (fp=fopen ("distlist.dat","r" )) == NULL )
{printf("Cannot open distlist.dat file" ); exit(2); }
else{
fgets(buffer, 100, fptr);
buflen = (MQLONG)strlen(buffer); /* length without null */
if (buffer[buflen-1] == '\n') /* last char is a new-line */
{
buffer[buflen-1] = '\0'; /* replace new-line with null */
--buflen; /* reduce buffer length */
}
}
fclose (fp);
tmr = time(NULL);
strcpy ( buf, ctime(tmr));
buf[strlen(buf)-5]=0;
printf("Distlist start send message to list queue %s\n", buf);
/* Open the target message queue for output */
od.Version = MQOD_VERSION_2 ;
od.RecsPresent = NumQueues ;
od.ObjectRecPtr = pOR;
od.ResponseRecPtr = pRR ;
O_options = MQOO_OUTPUT + MQOO_FAIL_IF_QUIESCING;
MQOPEN(Hcon, od, O_options, Hobj, OpenCode, Reason);
if (Reason == MQRC_MULTIPLE_REASONS)
{
print_responses("MQOPEN", pRR, NumQueues, pOR);
}
else
{
if (Reason != MQRC_NONE)
{
printf("MQOPEN returned CompCode=%d, Reason=%d\n", OpenCode, Reason);
}
}
/* Read message from the file
/* Loop until null line or end of file, or there is a failure */
CompCode = OpenCode; /* use MQOPEN result for initial test */
pmo.Version = MQPMO_VERSION_2 ;
pmo.RecsPresent = NumQueues ;
pmo.PutMsgRecPtr = pPMR ;
pmo.PutMsgRecFields = PutMsgRecFields ;
pmo.ResponseRecPtr = pRR ;
/* Put buffer to the message queue */
if (buflen > 0)
{
for( Index = 0 ; Index < NumQueues ; Index ++)
{
memcpy( (pPMR+Index)->MsgId, MQMI_NONE, sizeof((pPMR+Index)->MsgId));
memcpy( (pPMR+Index)->CorrelId, MQCI_NONE, sizeof((pPMR+Index)->CorrelId));
}
memcpy(md.Format, MQFMT_STRING, (size_t)MQ_FORMAT_LENGTH);
MQPUT(Hcon, Hobj, md, pmo, buflen, buffer, CompCode, Reason);
if (Reason == MQRC_MULTIPLE_REASONS)
{
print_responses("MQPUT", pRR, NumQueues, pOR);
}
else
{
if (Reason != MQRC_NONE)
{
printf("MQPUT returned CompCode=%d, Reason=%d\n", OpenCode, Reason);
}
}
tmr = time(NULL);
strcpy ( buf, ctime(tmr));
buf[strlen(buf)-5]=0; // strip new line
printf("Distlist finish send message to list queue %s\n", buf);
}
else /* satisfy end condition when empty line is read */
CompCode = MQCC_FAILED;
//}
if (OpenCode != MQCC_FAILED)
{
C_options = 0;
MQCLOSE(Hcon, Hobj, C_options, CompCode, Reason);
if (Reason != MQRC_NONE)
{
printf("MQCLOSE ended with reason code %d\n", Reason);
}
}
if (DisconnectRequired==True)
{
MQDISC(Hcon, CompCode, Reason);
if (Reason != MQRC_NONE)
{
printf("MQDISC ended with reason code %d\n", Reason);
}
}
if( NULL != pOR )
{
free( pOR ) ;
}
if( NULL != pRR )
{
free( pRR ) ;
}
if( NULL != pPMR )
{
free( pPMR ) ;
}
return(0);
}
static void print_usage(void)
{
printf("Distlist correct usage is:\n\n");
printf("Distlist Qmgr QName1 [QName2 [QName3 [...]]]\n\n");
}
static void print_responses( char * comment, PMQRR pRR, MQLONG NumQueues, PMQOR pOR)
{
MQLONG Index;
for( Index = 0 ; Index < NumQueues ; Index ++ )
{
if( MQCC_OK != (pRR+Index)->CompCode )
{
printf("%s for %.48s( %.48s) returned CompCode=%d, Reason=%d\n"
, comment
, (pOR+Index)->ObjectName
, (pOR+Index)->ObjectQMgrName
, (pRR+Index)->CompCode
, (pRR+Index)->Reason);
}
}
}
В завершение раздела можно сказать, что время работы механизма
| Количество очередей | Время работы distlist (сек) при ОП 512Мбт | Время работы distlist (сек) при ОП 1Гбт |
|---|---|---|
| 200 | 1 | 1 |
| 400 | 1 | 1 |
| 600 | 1 | 1 |
| 800 | 2 | 1 |
| 1000 | 3 | 2 |
| 1200 | 3 | 2 |
Таким образом,
Задачи, решаемые с помощью механизмов
Первая программа будет достаточно простая и реализует так называемую модель "один к одному" или "точка-точка". Эта программа предназначена для чтения сообщений из очереди 1, записи их в очередь 2 и лог-файл на диске. Эта программа имеет практическое значение. Достаточно часто необходимо иметь файл переданных сообщений за определенный период времени, чтобы быстро ответить на вопрос "Было ли передано сообщение с такими идентификационными параметрами в теле сообщения:…"?
Автору приходилось сталкиваться с "плохим" стилем rewriter.exe и файл инициализации rewriter.ini, в котором 1-я строка – имя очереди для чтения, 2-я строка – имя очереди для записи, 3-я строка – имя лог-файла, как показано ниже.
QUEUE_INPUT QUEUE_OUTPUT C:\TEMP\rewriter.log
Разрабатываемая программа может быть представлена в следующей последовательности псевдокода:
MQCONN MQOPEN --> цикл чтения сообщений | (на основе gmo.WaitInterval): | MQGET | MQPUT |-- конец цикла MQCLOSE MQDISC
Ниже приводится листинг программы , , в случае проблем с
/* Листинг программы rewriter */
/**********************************************************************/
/* Program name: Rewriter */
/* Description: Rewriter C program pass messages to output queue */
/* Function: */
/* Rewriter is a sample C program to demonstrate the main MQI calls; */
/* each message is copied from the input queue to the output */
/* queue, and sends a report to the log file */
/**********************************************************************/
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
#include <signal.h>
#include <io.h>
/* includes for MQI */
#include <cmqc.h>
char queue1[48] = "";
char queue2[48] = "";
char logfilename[48] = "";
char logfilename2[48] = "";
char buf[48];
int queuenamelen;
time_t tmr;
FILE * fp;
FILE *fptr;
void cntrl_c_handler(int sig);
/* Declare MQI structures needed */
MQOD odG = {MQOD_DEFAULT}; /* Object Descriptor for GET */
MQOD odP = {MQOD_DEFAULT}; /* Object Descriptor for PUT */
MQOD odI = {MQOD_DEFAULT}; /* Object Descriptor for InitQ */
MQOD odR = {MQOD_DEFAULT}; /* Object Descriptor for report */
MQMD md = {MQMD_DEFAULT}; /* Message Descriptor */
MQGMO gmo = {MQGMO_DEFAULT}; /* get message options */
MQPMO pmo = {MQPMO_DEFAULT}; /* put message options */
MQTMC2 *trig; /* trigger message structure */
MQCHAR48 QManager; /* queue manager name */
MQHCONN Hcon; /* connection handle */
MQHOBJ Hobj; /* object handle, server queue */
MQHOBJ Hinq; /* handle for MQINQ */
MQHOBJ Hout; /* handle for MQPUT */
MQLONG O_options; /* MQOPEN options */
MQLONG C_options; /* MQCLOSE options */
MQLONG CompCode; /* completion code */
MQLONG Reason; /* reason code */
MQLONG CReason; /* reason code (MQCONN) */
MQBYTE buffer[8001]; /* message buffer */
MQLONG buflen; /* buffer length */
MQLONG messlen; /* message length received */
MQLONG Select[1]; /* attribute selectors */
MQLONG SelectValue[1]; /* value attribute selectors */
MQLONG char_count;
int main(int argc, char **argv)
{
strcpy(QManager, ""); /* Работаем с менеджером очередей по умолчанию */
if ( (fptr=fopen ("rewriter.ini","r" )) == NULL )
{printf("Cannot open rewriter.ini file" ); exit(1); }
else{ /* Открываем ini-файл и присваиваем значения переменным */
fgets(queue1, 48, fptr);
queuenamelen = strlen(queue1) - 1;
queue1[queuenamelen] = ' ';
fgets(queue2, 48, fptr);
queuenamelen = strlen(queue2) - 1;
queue2[queuenamelen] = ' ';
fgets(logfilename, 48, fptr);
queuenamelen = strlen(logfilename) - 1;
logfilename[queuenamelen] = ' ';
tmr = time(NULL);
strcpy ( buf, ctime(tmr));
buf[strlen(buf)-1]=0; // переход на новую строку
strncat (logfilename, buf,10);
strcpy(odG.ObjectName, queue1);
strcpy(odP.ObjectName, queue2);
fclose (fptr);
}
MQCONN(QManager, Hcon, CompCode, CReason);
if (CompCode == MQCC_FAILED)
{
printf("MQCONN ended with reason code %ld\n", CReason);
exit(CReason);
}
O_options = MQOO_INPUT_SHARED + MQOO_FAIL_IF_QUIESCING;
MQOPEN(Hcon, odG, O_options, Hobj, CompCode, Reason); /* открываем очередь для чтения - odG */
if (Reason != MQRC_NONE) { printf("MQOPEN (input) ended with reason code %ld\n", Reason); }
if (CompCode == MQCC_FAILED) { exit(Reason); }
O_options = MQOO_OUTPUT + MQOO_FAIL_IF_QUIESCING;
MQOPEN(Hcon, odP, O_options, Hout, CompCode, Reason); /* открываем очередь для записи - odP */
if (Reason != MQRC_NONE) { printf("MQOPEN (output) ended with reason code %ld\n", Reason); }
if (CompCode == MQCC_FAILED) { exit(Reason); }
fp = fopen (logfilename,"a");
if ( fp==NULL ){ printf("Cannot open log file %s\n", logfilename); }
printf("Rewriter(C) sending messages from %s to %s and to log-file %s \n
",odG.ObjectName, odP.ObjectName, logfilename);
/*****************************************************************************/
/* Читаем сообщения из QUEUE_INPUT и пишем в QUEUE_OUTPUT */
/* до тех пор пока не встретим сообщение об ошибке */
/*****************************************************************************/
buflen = sizeof(buffer) - 1;
while (CompCode == MQCC_OK)
{
gmo.Options = MQGMO_ACCEPT_TRUNCATED_MSG + MQGMO_WAIT;
gmo.WaitInterval = 3000; /* Ожидаем новые сообщения 3 секунды */
//gmo.WaitInterval = MQWI_UNLIMITED;
memcpy(md.MsgId, MQMI_NONE, sizeof(md.MsgId));
memcpy(md.CorrelId, MQMI_NONE, sizeof(md.CorrelId));
MQGET(Hcon, Hobj, md, gmo, buflen, buffer, messlen, CompCode, Reason);
if ((CompCode == MQCC_OK) || (CompCode == MQCC_WARNING))
{
buffer[messlen] = '\0'; /* заносим символ конец строки в буфер с прочитанным сообщением */
buflen = messlen;
MQPUT(Hcon, Hout, md, pmo, buflen, buffer, CompCode, Reason);
if ((CompCode == MQCC_OK) || (CompCode == MQCC_WARNING))
{
tmr = time(NULL);
strcpy ( buf, ctime(tmr));
buf[strlen(buf)-1]=0; // переход к новой строке
Reason = fprintf( fp, "%s: %s\n", buf, buffer );
}
} /* конец обработки входного сообщения */
} /* конец цикла чтения/записи сообщений функциями MQGET, MQPUT */
fclose (fp);
C_options = 0; /* нет никаких опций при закрытии */
MQCLOSE(Hcon, Hobj, C_options, CompCode, Reason); /* закрываем очередь для чтения */
if (Reason != MQRC_NONE)
{printf("MQCLOSE (input) ended with reason code %ld\n", Reason); }
MQCLOSE(Hcon, Hout, C_options, CompCode, Reason); /* закрываем очередь для записи */
if (Reason != MQRC_NONE)
{printf("MQCLOSE (output) ended with reason code %ld\n", Reason); }
if (CReason != MQRC_ALREADY_CONNECTED)
{
MQDISC(Hcon, CompCode, Reason);
if (Reason != MQRC_NONE){ printf("MQDISC ended with reason code %ld\n", Reason); }
}
return(0);
}
В данной версии мы выходим из цикла программы по опции gmo.WaitInterval = 3000, когда ожидаем сообщение в очереди в течении 3 сек, а его там нет (опция gmo.WaitInterval работает быстрее, чем если бы мы опрашивали очередь по собственному временному циклу). Другой вариант программы может быть таким. Задаем gmo.WaitInterval = MQWI_UNLIMITED ; что соответствует gmo.WaitInterval= -1. Программа будет крутиться "бесконечно" до тех пор, пока мы не остановим её принудительно, например, нажатием клавиш CNTRL_C (стандартный останов). В этом случае нужно добавить обработчик прерываний по нажатию CNTRL_C потому, что при таком выходе объекты очереди останутся не закрытыми и , в противном случае лог-файл не будет формироваться. Следует отметить, что размер массива buffer ограничивает длину сообщения 8Кб и при появлении сообщений большей длины следует увеличить размер буфера.
Программа rewriter.exe работает достаточно быстро и сравнительные скорости работы данного алгоритма при длине сообщения 1Кб на компьютере INTEL Pentium 1.8Ггц приведены в таблице ниже.
| Язык программы\тип очереди | Not Persistent | Persistent |
|---|---|---|
| 1000 сооб/сек | 400 сооб/сек | |
| Visual Basic 6.0 | 200 сооб/сек | 140 сооб/сек |
Увеличение длины сообщения не ведет к пропорциональному уменьшению скорости. Эти исследования читатель может проделать самостоятельно. Реальные приложения, работающие с базами данных, имеют скорость обработки сообщений в 3-4 раза меньше.
Возвращаясь к вопросу о стилях rewriter.ini файле приведет к зависанию программы и мучительному поиску причин такого зависания, не говоря о других более сложных ситуациях, например, когда очередь открыта эксклюзивно другим приложением.
Для версии программы gmo.WaitInterval = MQWI_UNLIMITED полезно сделать вывод на экран передаваемых сообщений, чтобы наблюдать динамику работы созданного интерфейса. Таких улучшений может быть достаточно много и мы рассмотрим две достаточно полезные модификации.
Программа rewriter может вызываться как
/* Код для вызова rewriter.exe
как MQSeries-триггер */
int main(int argc,
char **argv)
{ if (argc > 1)
{trig = (MQTMC2*)argv[1];
strncpy(odG.ObjectName,
trig->QName,
MQ_Q_NAME_LENGTH);
strncpy(queue1,
trig->QName,
MQ_Q_NAME_LENGTH);
strncpy(QManager,
trig->QMgrName,
MQ_Q_MGR_NAME_LENGTH);
strncpy(odP.ObjectName,
trig->UserData,
MQ_PROCESS_USER_DATA_LENGTH);
strncpy(queue2, trig->UserData,
MQ_PROCESS_USER_DATA_LENGTH);
strncpy(logfilename2,
trig->EnvData, 48);
}
Возможная модификация этого варианта - программа rewriter может вызываться с передачей параметров через командную строку и эту модификацию читатель может проделать самостоятельно.
rewriter может быть модифицирована в программу разветвитель mqsplitter.exe: чтение сообщений из очереди 1 и запись их в очередь 2, в очередь 3 и лог-файл на диске.Можно сделать программу mqsplitter.exe на разных языках, например, на Visual Basic 6.0 с интерфейсом, показанным на рис.9.1, и сравнить производительность программ на разных языках, реализующих один и тот же алгоритм. Такая задача будет хорошим лабораторным практикумом.
(рис 9.1) Интерфейс программы mqsplitter на VB6Модификацию программы mqsplitter.exe читателю предлагается сделать самостоятельно и одновременно проверить идею создания "вечного двигателя". Для создания "вечного двигателя" понадобиться изменить исходный mqsplitter.ini файл следующим образом:
|
-> |
|
Если в очереди QUEUE_INPUT будет хотя бы одно сообщение, то ваша программа будет посылать сообщения в очередь QUEUE_OUTPUT1 до тех пор, пока не будет остановлена. Можно заложить в ini-файл программы пятый параметр: время опроса очереди, измеряемое в миллисекундах. Такая программа окажется весьма полезной при тестировании интерфейсов.
Задачи
В этой задаче для нас интересна технология такого перехода. Совершенно очевидно, что переход от одной автоматизированной
В качестве конкретной задачи рассмотрим создание интерфейса по передаче клиентов из системы
Существует разные варианты решения поставленной задачи:
Первый путь ясен в технической реализации. С одной стороны при каждом обновлении таблицы client в АБС1 программа-обработчик срабатывает как триггер базы данных и помещает результаты оператора update в очередь, они приходят на АБС2, где своя программа-обработчик срабатывает как триггер очереди и помещает сообщение в таблицу client АБС2.
Блок 1 MQCONN; MQOPEN; Блок 2 MQBEGIN; MQGET; Блок 3 Begin Tran Блок 4 UPDATE CLIENT SET ... WHERE ... Блок 5 If Error = 0 then Commit Tran else Rollback Tran; Блок 6 If Error = 0 then MQCMIT else MQBACK; Блок 7 MQCLOSE; MQDISC;
В этой программе
Тем не менее, при небольшом числе интерфейсов рекомендуется вариант 1, как наиболее экономический, к рассмотрению которого мы и переходим. WebSphere BI Broker будет рассмотрен в лекции 12. К сожалению, рассмотрение средств
Сообщения
Наша очередная задача: на сервере 1 прочитать сообщение из входной очереди, положить её в очередь для отправки на сервер 2 как сообщение-запрос и дождаться прихода сообщения-ответа, как это показано на рис.9.2. Все это необходимо оформить в виде
(рис 9.2) Структура объектов WebSphere MQИтак, последовательность псевдокода представляется следующим образом (обратите внимание на блок 5 и опции MQMD):
Блок 1 MQCONN Блок 2 MQOPEN Блок 3 MQBEGIN Блок 4 MQGET (Input_queue) Блок 5 MQPUT (Output_queue, MQMD.MsgType = MQMT_REQUEST, MQMD.ReplyToQ = Reply_queue) Блок 6 MQGET (Reply_queue) Блок 7 If Reply time < 10 sec then MQCMIT else MQBACK; Блок 8 MQCLOSE Блок 9 MQDISC
Назовем нашу программу transmit.exe и файл инициализации transmit.ini, в котором 1-я строка – имя очереди для чтения, 2-я строка – имя очереди для записи, 3-я строка – имя очереди для ответа, 4-я строка – время ожидания ответа Reply_time = 3000мсек, как показано ниже.
QUEUE_INPUT QUEUE_OUTPUT QUEUE_REPLY 3000
Тип очереди Output_queue –
Ниже приводится листинг программы transmit.cpp для Microsoft Visual C++ ver.6.0. Для каждого сообщения и ) и об этом подробнее в лекции 11.
/* Листинг программы transmit */
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
#include <io.h>
#include <cmqc.h>
char queue_input[48] = "";
char queue_output[48] = "";
char queue_reply[48] = "";
char reply_time[48] = "";
char buf[48];
time_t tmr;
int queuenamelen;
FILE *fptr;
int main(int argc, char **argv)
{
MQOD odG = {MQOD_DEFAULT};
MQOD odP = {MQOD_DEFAULT};
MQOD odR = {MQOD_DEFAULT};
MQOD odI = {MQOD_DEFAULT};
MQMD md = {MQMD_DEFAULT};
MQBO mbo = {MQBO_DEFAULT};
MQGMO gmo = {MQGMO_DEFAULT};
MQPMO pmo = {MQPMO_DEFAULT};
MQCHAR48 QManager;
MQHCONN Hcon;
MQHOBJ Hobj;
MQHOBJ Hout;
MQHOBJ Hrep;
MQLONG O_options;
MQLONG C_options;
MQLONG CompCode;
MQLONG Reason;
MQLONG CReason;
MQBYTE buffer[8001];
MQLONG buflen;
MQLONG replylen;
MQLONG messlen;
static MQBYTE24 LastMsgId;
if ( (fptr=fopen ("transmit.ini","r" )) == NULL )
{printf("Cannot open transmit.ini file" );
exit(1);
}
else{
fgets(queue_input, 48, fptr);
queuenamelen = strlen(queue_input) - 1;
queue_input[queuenamelen] = ' ';
strcpy(odG.ObjectName, queue_input);
fgets(queue_output, 48, fptr);
queuenamelen = strlen(queue_output) - 1;
queue_output[queuenamelen] = ' ';
strcpy(odP.ObjectName, queue_output);
fgets(queue_reply, 48, fptr);
queuenamelen = strlen(queue_reply) - 1;
queue_reply[queuenamelen] = ' ';
strcpy(odR.ObjectName, queue_reply);
fgets(reply_time, 48, fptr);
queuenamelen = strlen(reply_time) - 1;
reply_time[queuenamelen] = ' ';
fclose (fptr);
}
strcpy(QManager, ""); /* Работаем с менеджером очередей по умолчанию */
MQCONN(QManager, Hcon, CompCode, CReason);
if (CompCode == MQCC_FAILED)
{
printf("MQCONN to %s ended with reason code %ld\n", QManager, CReason);
exit(CReason);
}
O_options = MQOO_BROWSE + MQOO_INPUT_SHARED ;
MQOPEN(Hcon, odG, O_options, Hobj, CompCode, Reason);
if (Reason != MQRC_NONE) { printf("MQOPEN %s ended with reason code %ld\n", queue_input, Reason); }
if (CompCode == MQCC_FAILED) { exit(Reason); }
O_options = MQOO_BROWSE + MQOO_INPUT_SHARED ;
MQOPEN(Hcon, odR, O_options, Hrep, CompCode, Reason);
if (Reason != MQRC_NONE) { printf("MQOPEN %s ended with reason code %ld\n", queue_reply, Reason); }
if (CompCode == MQCC_FAILED) { exit(Reason); }
O_options = MQOO_OUTPUT ;
MQOPEN(Hcon, odP, O_options, Hout, CompCode, Reason);
if (Reason != MQRC_NONE) { printf("MQOPEN %s ended with reason code %ld\n", queue_output, Reason); }
if (CompCode == MQCC_FAILED) { exit(Reason); }
while (CompCode == MQCC_OK)
{
buflen = sizeof(buffer) - 1;
memcpy(md.MsgId, MQMI_NONE, sizeof(md.MsgId));
memcpy(md.CorrelId, MQCI_NONE, sizeof(md.CorrelId));
gmo.Options = MQGMO_ACCEPT_TRUNCATED_MSG + MQGMO_WAIT + MQGMO_SYNCPOINT;
gmo.WaitInterval = 3000 ;
MQBEGIN (Hcon, mbo, CompCode, Reason);
MQGET(Hcon, Hobj, md, gmo, buflen, buffer, messlen, CompCode, Reason);
//if (Reason != MQRC_NONE) { printf("MQGET from %s ended with reason code %ld\n", queue_input, Reason); }
if ((CompCode == MQCC_OK) || (CompCode == MQCC_WARNING))
{
buffer[messlen] = '\0'; /* заносим символ конец строки в буфер с прочитанным сообщением */
buflen = messlen;
md.MsgType = MQMT_REQUEST;
md.Report = MQRO_EXCEPTION_WITH_DATA;
strncpy(md.ReplyToQ, queue_reply, MQ_Q_NAME_LENGTH);
memcpy(md.Format, MQFMT_STRING, MQ_FORMAT_LENGTH);
MQPUT(Hcon, Hout, md, pmo, buflen, buffer, CompCode, Reason);
if (Reason != MQRC_NONE)
{
printf("MQPUT to %s ended ended unsuccessfully with reason code %ld CompCode %ld\n", queue_output, Reason, CompCode );
MQBACK( Hcon, CompCode, Reason ) ;
CompCode = MQCC_FAILED ;
}
else
{
while (CompCode != MQCC_FAILED)
{
/** осуществляется проверка queue_reply **/
gmo.Options = MQGMO_ACCEPT_TRUNCATED_MSG + MQGMO_WAIT ;
gmo.WaitInterval = 3000 ;
memcpy(md.MsgId, MQMI_NONE, sizeof(md.MsgId));
memcpy(md.CorrelId, MQCI_NONE, sizeof(md.CorrelId));
MQGET(Hcon, Hrep, md, gmo, buflen, buffer, replylen, CompCode, Reason);
if (CompCode != MQCC_FAILED)
{
if (md.MsgType == MQMT_REPLY) /* report feedback */
{ printf("Transaction % s=> %s successfully: %s\n", queue_input, queue_output, buffer);
MQCMIT( Hcon, CompCode, Reason ) ;
}
else
{
printf("Transaction % s=> %s successfully, REPLY message not deliver, reason code %ld CompCode %ld\n", queue_input, queue_output, queue_reply, Reason, CompCode );
MQBACK( Hcon, CompCode, Reason ) ;
CompCode = MQCC_FAILED ;
}
}
if (Reason == MQRC_NO_MSG_AVAILABLE)
{
printf("Transaction % s=> %s UNsuccessfully, REPLY message not deliver\n", queue_input, queue_output );
MQBACK( Hcon, CompCode, Reason ) ;
CompCode = MQCC_FAILED ;
}
}
}
}
}
C_options = 0;
MQCLOSE(Hcon, Hobj, C_options, CompCode, Reason);
if (Reason != MQRC_NONE){printf("MQCLOSE %s ended with reason code %ld\n", queue_input, Reason); }
MQCLOSE(Hcon, Hout, C_options, CompCode, Reason);
if (Reason != MQRC_NONE){printf("MQCLOSE %s ended with reason code %ld\n", queue_output, Reason); }
MQCLOSE(Hcon, Hrep, C_options, CompCode, Reason);
if (Reason != MQRC_NONE){printf("MQCLOSE %s ended with reason code %ld\n", queue_reply, Reason); }
MQDISC(Hcon, CompCode, Reason);
if (Reason != MQRC_NONE){ printf("MQDISC ended with reason code %ld\n", Reason); }
return(0);
}
По тексту программы следует дать комментарии. Наличие опции
gmo.Options = MQGMO_SYNCPOINT;
подразумевает, что команда может не указываться. Операторы
md.MsgType = MQMT_REQUEST;
strncpy(md.ReplyToQ,
queue_reply,
MQ_Q_NAME_LENGTH);
определяют тип сообщения и очередь ответа, заданную в QUEUE_REPLY.
На очередь QUEUE_OUTPUT (или на удаленную очередь на другом менеджере) должна быть навешена программа-триггер, который возвращает сообщения типа . Если QUEUE_REPLY, то QUEUE_INPUT. такой же, как и исходного сообщения. В данной версии программы в целях упрощения отладки не проверяется это условие и читателю предлагается самостоятельно дописать этот фрагмент кода после отладки текущей версии программы. Работа с и будет рассмотрена подробнее в лекции 11.
Программу-триггер, которая "навешивается" на очередь QUEUE_OUTPUT (или на удаленную очередь) для формирования md.MsgType = MQMT_REPLY; ), читателю также предлагается сделать самостоятельно.
На данном примере мы познакомились с
положить сообщения в эти очереди. После помещает сообщения во все эти очереди, используя этот единственный
В версии
Рассмотрим этот механизм на примере задачи, когда distlist.exe, файл с текстом сообщения distlist.dat и файл инициализации distlist.ini, в котором 1-я строка – имя менеджера, 2-я и последующие строки – имена очередей, как показано ниже.
QM_ ALFA Queue_ Moscow Queue_ Kiev Queue_ Alma-Ata Queue_ SPetersburg Queue_ Novosibirsk Queue_ Saratov //last string must be blank
(рис 9.3) Механизм Distribution List для WebSphere MQНиже приводится листинг программы distlist.cpp для Microsoft Visual C++ ver.6.0.
/* Листинг программы distlist */
/* Program name: Distlist */
/* Description: Distlist C program pass messages to output queues */
/* by Distribution list for indicated Queue Manager */
/* distlist.ini file give list of queue and distlist.dat give file */
/* of message which copied to the output queue */
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <io.h>
#include <time.h>
#include <cmqc.h>
char queue[1000][48] ;
char buf[48];
int queuenamelen;
time_t tmr;
FILE * fp;
FILE *fptr;
static void print_usage(void);
static void print_responses( char * comment, PMQRR pRR, MQLONG NumQueues, PMQOR pOR);
int main(int argc, char **argv)
{
typedef enum {False, True} Bool;
MQOD od = {MQOD_DEFAULT}; /* Object Descriptor */
MQMD md = {MQMD_DEFAULT}; /* Message Descriptor */
MQPMO pmo = {MQPMO_DEFAULT}; /* put message options */
MQHCONN Hcon; /* connection handle */
MQHOBJ Hobj; /* object handle */
MQLONG O_options; /* MQOPEN options */
MQLONG C_options; /* MQCLOSE options */
MQLONG CompCode; /* completion code */
MQLONG OpenCode; /* MQOPEN completion code */
MQLONG Reason; /* reason code */
MQCHAR48 QManager; /* queue manager name */
MQLONG buflen; /* buffer length */
char buffer[101]; /* message buffer */
MQLONG Index ; /* Index into list of queues */
MQLONG NumQueues ; /* Number of queues */
PMQRR pRR=NULL; /* Pointer to response records */
PMQOR pOR=NULL; /* Pointer to object records */
Bool DisconnectRequired=False;/* Already connected switch */
Bool Connected=False; /* Connect succeeded switch */
typedef struct
{
MQBYTE24 MsgId;
MQBYTE24 CorrelId;
} PutMsgRec, *pPutMsgRec;
pPutMsgRec pPMR=NULL; /* Pointer to put msg records */
MQLONG PutMsgRecFields=MQPMRF_MSG_ID | MQPMRF_CORREL_ID;
/* Open ini file and setting value */
if ( (fptr=fopen ("distlist.ini","r" )) == NULL )
{printf("Cannot open distlist.ini file" );
print_usage();
exit(1); }
else{
fgets(QManager, 48, fptr);
queuenamelen = strlen(QManager) - 1;
QManager[queuenamelen] = ' ';
NumQueues = 0;
while (queuenamelen != 0)
{
fgets(queue[NumQueues], 48, fptr);
queuenamelen = strlen(queue[NumQueues]) - 1;
queue[NumQueues][queuenamelen] = ' ';
NumQueues++;
}
}
fclose (fptr);
--NumQueues; /* NumQueues - Number of Queue name */
/* Allocate response records, object records and put message records */
pRR = (PMQRR)malloc( NumQueues * sizeof(MQRR));
pOR = (PMQOR)malloc( NumQueues * sizeof(MQOR));
pPMR = (pPutMsgRec)malloc( NumQueues * sizeof(PutMsgRec));
if((NULL == pRR) || (NULL == pOR) || (NULL == pPMR))
{
printf("%s(%d) malloc failed\n", __FILE__, __LINE__);
exit(4);
}
/* Use parameters as the name of the target queues */
for( Index = 0 ; Index < NumQueues ; Index ++)
{
strncpy( (pOR+Index)->ObjectName, queue[Index], (size_t)MQ_Q_NAME_LENGTH);
strncpy( (pOR+Index)->ObjectQMgrName, QManager, (size_t)MQ_Q_MGR_NAME_LENGTH);
}
for( Index = 0 ; Index < NumQueues ; Index ++)
{
MQCONN((pOR+Index)->ObjectQMgrName, Hcon, ((pRR+Index)->CompCode), ((pRR+Index)->Reason));
if ((pRR+Index)->CompCode == MQCC_FAILED)
{
continue;
}
if ((pRR+Index)->CompCode == MQCC_OK)
{
DisconnectRequired = True ;
}
Connected = True;
break ;
}
/* Print any non zero responses */
print_responses("MQCONN", pRR, Index, pOR);
/* Print If failed to connect to queue manager then exit. */
if( False == Connected )
{
printf("Unable to connect to queue manager\n");
exit(3) ;
}
if ( (fp=fopen ("distlist.dat","r" )) == NULL )
{printf("Cannot open distlist.dat file" ); exit(2); }
else{
fgets(buffer, 100, fptr);
buflen = (MQLONG)strlen(buffer); /* length without null */
if (buffer[buflen-1] == '\n') /* last char is a new-line */
{
buffer[buflen-1] = '\0'; /* replace new-line with null */
--buflen; /* reduce buffer length */
}
}
fclose (fp);
tmr = time(NULL);
strcpy ( buf, ctime(tmr));
buf[strlen(buf)-5]=0;
printf("Distlist start send message to list queue %s\n", buf);
/* Open the target message queue for output */
od.Version = MQOD_VERSION_2 ;
od.RecsPresent = NumQueues ;
od.ObjectRecPtr = pOR;
od.ResponseRecPtr = pRR ;
O_options = MQOO_OUTPUT + MQOO_FAIL_IF_QUIESCING;
MQOPEN(Hcon, od, O_options, Hobj, OpenCode, Reason);
if (Reason == MQRC_MULTIPLE_REASONS)
{
print_responses("MQOPEN", pRR, NumQueues, pOR);
}
else
{
if (Reason != MQRC_NONE)
{
printf("MQOPEN returned CompCode=%d, Reason=%d\n", OpenCode, Reason);
}
}
/* Read message from the file
/* Loop until null line or end of file, or there is a failure */
CompCode = OpenCode; /* use MQOPEN result for initial test */
pmo.Version = MQPMO_VERSION_2 ;
pmo.RecsPresent = NumQueues ;
pmo.PutMsgRecPtr = pPMR ;
pmo.PutMsgRecFields = PutMsgRecFields ;
pmo.ResponseRecPtr = pRR ;
/* Put buffer to the message queue */
if (buflen > 0)
{
for( Index = 0 ; Index < NumQueues ; Index ++)
{
memcpy( (pPMR+Index)->MsgId, MQMI_NONE, sizeof((pPMR+Index)->MsgId));
memcpy( (pPMR+Index)->CorrelId, MQCI_NONE, sizeof((pPMR+Index)->CorrelId));
}
memcpy(md.Format, MQFMT_STRING, (size_t)MQ_FORMAT_LENGTH);
MQPUT(Hcon, Hobj, md, pmo, buflen, buffer, CompCode, Reason);
if (Reason == MQRC_MULTIPLE_REASONS)
{
print_responses("MQPUT", pRR, NumQueues, pOR);
}
else
{
if (Reason != MQRC_NONE)
{
printf("MQPUT returned CompCode=%d, Reason=%d\n", OpenCode, Reason);
}
}
tmr = time(NULL);
strcpy ( buf, ctime(tmr));
buf[strlen(buf)-5]=0; // strip new line
printf("Distlist finish send message to list queue %s\n", buf);
}
else /* satisfy end condition when empty line is read */
CompCode = MQCC_FAILED;
//}
if (OpenCode != MQCC_FAILED)
{
C_options = 0;
MQCLOSE(Hcon, Hobj, C_options, CompCode, Reason);
if (Reason != MQRC_NONE)
{
printf("MQCLOSE ended with reason code %d\n", Reason);
}
}
if (DisconnectRequired==True)
{
MQDISC(Hcon, CompCode, Reason);
if (Reason != MQRC_NONE)
{
printf("MQDISC ended with reason code %d\n", Reason);
}
}
if( NULL != pOR )
{
free( pOR ) ;
}
if( NULL != pRR )
{
free( pRR ) ;
}
if( NULL != pPMR )
{
free( pPMR ) ;
}
return(0);
}
static void print_usage(void)
{
printf("Distlist correct usage is:\n\n");
printf("Distlist Qmgr QName1 [QName2 [QName3 [...]]]\n\n");
}
static void print_responses( char * comment, PMQRR pRR, MQLONG NumQueues, PMQOR pOR)
{
MQLONG Index;
for( Index = 0 ; Index < NumQueues ; Index ++ )
{
if( MQCC_OK != (pRR+Index)->CompCode )
{
printf("%s for %.48s( %.48s) returned CompCode=%d, Reason=%d\n"
, comment
, (pOR+Index)->ObjectName
, (pOR+Index)->ObjectQMgrName
, (pRR+Index)->CompCode
, (pRR+Index)->Reason);
}
}
}
В завершение раздела можно сказать, что время работы механизма
| Количество очередей | Время работы distlist (сек) при ОП 512Мбт | Время работы distlist (сек) при ОП 1Гбт |
|---|---|---|
| 200 | 1 | 1 |
| 400 | 1 | 1 |
| 600 | 1 | 1 |
| 800 | 2 | 1 |
| 1000 | 3 | 2 |
| 1200 | 3 | 2 |
Таким образом,
Задачи, решаемые с помощью механизмов
Для получения официальных документов о завершении программы дополнительного профессионального образования (удостоверения о повышении квалификации, дипломов о профессиональной переподготовке и MBA) необходимо предоставить:
Внимание! Вы можете не заказывать доставку бумажной версии официального документы, а скачать его в электронном виде и распечатать самостоятельно. Информация о выданном документе в течение 1 месяца загружается в Федеральную информационную систему «Федеральный реестр сведений о документах об образовании и (или) о квалификации, документах об обучении» - ФИС ФРДО.
Доступ на новый сайт осуществляется с использованием адреса электронной почты, который был указан вами при регистрации на "старом". Мы постарались перенести все ваши данные с прежнего ресурса, однако не исключена вероятность потери части информации.
При возникновении проблемы со входом, воспользуйтесь функцией сброса пароля
Если вы обнаружите несоответствия, пожалуйста, сообщите нам.