アドバンスドキューイングの使用

EDBPostgresAdvancedServerAdvancedQueuingは、AdvancedServerデータベースにメッセージキューイングとメッセージ処理を提供します。ユーザー定義のメッセージはキューに格納されます。キューのコレクションは、キューテーブルに格納されます。依存するキューを作成する前に、まずキューテーブルを作成する必要があります。

On the server side, procedures in the DBMS_AQADM package create and manage message queues and queue tables. Use the DBMS_AQ package to add or remove messages from a queue, or register or unregister a PL/SQL callback procedure. For more information about DBMS_AQ and DBMS_AQADM, click here.

クライアント側では、アプリケーションはEDB.NETドライバーを使用してメッセージをエンキュー/デキューします。

メッセージのエンキューまたはデキュー

サーバー側の設定

.NETアプリケーションでアドバンストキューイング機能を使用するには、まずユーザー定義型、キューテーブル、およびキューを作成してから、データベースサーバーでキューを開始する必要があります。EDB-PSQLを呼び出して、AdvancedServerホストデータベースに接続します。コマンドラインで次のSPLコマンドを使用します。

ユーザー定義タイプの作成

RAWデータ型を指定するには、ユーザー定義型を作成する必要があります。次の例は、 myxml という名前のユーザー定義型の作成を示しています。

CREATE TYPE myxml AS (value XML) ;

キューテーブルの作成

キューテーブルは、同じペイロードタイプの複数のキューを保持できます。次の例は、 MSG_QUEUE_TABLE という名前のテーブルの作成を示しています。

EXEC DBMS_AQADM.CREATE_QUEUE_TABLE
      (queue_table => 'MSG_QUEUE_TABLE',
       queue_payload_type => 'myxml',
       comment => 'Message queue table');
END;

キューの作成

次の例は、テーブル MSG_QUEUE_TABLE 内に MSG_QUEUE という名前のキューを作成する方法を示しています。

BEGIN
DBMS_AQADM.CREATE_QUEUE ( queue_name => 'MSG_QUEUE', queue_table => 'MSG_QUEUE_TABLE', comment => 'This queue contains pending messages.');
END;

キューの開始

キューが作成されたら、コマンドラインで次のSPLコードを呼び出して、EDBデータベースのキューを開始します。

BEGIN
DBMS_AQADM.START_QUEUE
(queue_name => 'MSG_QUEUE');
END;

クライアント側の例

ユーザー定義型を作成し、続いてキューテーブルとキューを作成したら、キューを開始します。次に、EDB.Netドライバーを使用してメッセージをエンキューまたはデキューできます。

メッセージをキューに入れる:

.NETアプリケーションでメッセージをエンキューするには、以下を行う必要があります。

  1. EnterpriseDB.EDBClient 名前空間をインポートします。

  2. キューの名前を渡し、 EDBAQQueue のインスタンスを作成します。

  3. エンキューメッセージを作成し、ペイロードを定義します。

  4. queue.Enqueue メソッドを呼び出します。

次のコードリストは、 Queue.enqueue メソッドの使用を示しています:

using EnterpriseDB.EDBClient;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading.Tasks;

namespace AQXml
{
    class MyXML
    {
        public string value { get; set; }
    }
    class Program
    {
        static void Main(string[] args)
        {
            int messagesToSend = 1;
            if (args.Length > 0 && !string.IsNullOrEmpty(args[0]))
            {
                messagesToSend = int.Parse(args[0]);
            }
            for (int i = 0; i < 5; i++)
            {
                EnqueMsg("test message: " + i);
            }
        }

              private static EDBConnection GetConnection()
        {
            string connectionString = "Server=127.0.0.1;Host=127.0.0.1;Port=5444;User Id=enterprisedb;Password=test;Database=edb;Timeout=999";
            EDBConnection connection = new EDBConnection(connectionString);
            connection.Open();
            return connection;
        }


        private static string ByteArrayToString(byte[] byteArray)
        {
            // Sanity check if it's null so we don't incur overhead of an exception
            if (byteArray == null)
            {
                return string.Empty;
            }
            try
            {
                StringBuilder hex = new StringBuilder(byteArray.Length * 2);
                foreach (byte b in byteArray)
                {
                    hex.AppendFormat("{0:x2}", b);
                }

                return hex.ToString().ToUpper();
            }
            catch
            {
                return string.Empty;
            }
        }

        private static bool EnqueMsg(string msg)
        {
            EDBConnection con = GetConnection();
            using (EDBAQQueue queue = new EDBAQQueue("MSG_QUEUE", con))
            {
                queue.MessageType = EDBAQMessageType.Xml;
                EDBTransaction txn = queue.Connection.BeginTransaction();
                QueuedEntities.Message queuedMessage = new QueuedEntities.Message() { MessageText = msg };

                try
                {
                    string rootElementName = queuedMessage.GetType().Name;
                    if (rootElementName.IndexOf(".") != -1)
                    {
                        rootElementName = rootElementName.Split('.').Last();
                    }

                    string xml = new Utils.XmlFragmentSerializer<QueuedEntities.Message>().Serialize(queuedMessage);
                    EDBAQMessage queMsg = new EDBAQMessage();
                    queMsg.Payload = new MyXML { value = xml };
                    queue.MessageType = EDBAQMessageType.Udt;
                    queue.UdtTypeName = "myxml";
                    queue.Enqueue(queMsg);
                    var messageId = ByteArrayToString((byte[])queMsg.MessageId);
                    Console.WriteLine("MessageID: " + messageId);
                    txn.Commit();
                    queMsg = null;
                    xml = null;
                    rootElementName = null;
                    return true;
                }
                catch (Exception ex)
                {
                    txn?.Rollback();
                    Console.WriteLine("Failed to enqueue message.");
                    Console.WriteLine(ex.ToString());
                    return false;
                }
                finally
                {
                    queue?.Connection?.Dispose();
                }
            }
        }

    }
}

メッセージをキューから取り出す

.NETアプリケーションでメッセージをデキューするには、次のことを行う必要があります。

  1. EnterpriseDB.EDBClient 名前空間をインポートします。

  2. キューの名前を渡し、 EDBAQQueue のインスタンスを作成します。

  3. queue.Dequeue メソッドを呼び出します。

using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading.Tasks;
using EnterpriseDB.EDBClient;

namespace DequeueXML
{
    class MyXML
    {
        public string value { get; set; }
    }
    class Program
    {
        static void Main(string[] args)
        {
            DequeMsg();
        }


        private static EDBConnection GetConnection()
        {
            string connectionString = "Server=127.0.0.1;Host=127.0.0.1;Port=5444;User Id=enterprisedb;Password=test;Database=edb;Timeout=999";
            EDBConnection connection = new EDBConnection(connectionString);
            connection.Open();
            return connection;
        }


        private static string ByteArrayToString(byte[] byteArray)
        {
            // Sanity check if it's null so we don't incur overhead of an exception
            if (byteArray == null)
            {
                return string.Empty;
            }
            try
            {
                StringBuilder hex = new StringBuilder(byteArray.Length * 2);
                foreach (byte b in byteArray)
                {
                    hex.AppendFormat("{0:x2}", b);
                }

                return hex.ToString().ToUpper();
            }
            catch
            {
                return string.Empty;
            }
        }
        public static void DequeMsg(int waitTime = 10)
        {
            EDBConnection con = GetConnection();
            using (EDBAQQueue queueListen = new EDBAQQueue("MSG_QUEUE", con))
            {
                queueListen.UdtTypeName = "myxml";
                queueListen.DequeueOptions.Navigation = EDBAQNavigationMode.FIRST_MESSAGE;
                queueListen.DequeueOptions.Visibility = EDBAQVisibility.ON_COMMIT;
                queueListen.DequeueOptions.Wait = 1;
                EDBTransaction txn = null;

                while (1 == 1)
                    {

                    if (queueListen.Connection.State == System.Data.ConnectionState.Closed)
                    {
                        queueListen.Connection.Open();
                    }

                    string messageId = "Unknown";
                    try
                    {
                        // the listen function is a blocking function. It will Wait the specified waitTime or until a
                        // message is received.
                        Console.WriteLine("Listening...");
                        string v = queueListen.Listen(null, waitTime);
                        // If we are waiting for a message and we specify a Wait time,
                        // then if there are no more messages, we want to just bounce out.
                        if (waitTime > -1 && v == null)
                        {
                            Console.WriteLine("No message received during Wait period.");
                            Console.WriteLine();
                            continue;
                        }

                        // once we're here that means a message has been detected in the queue. Let's deal with it.
                        txn = queueListen.Connection.BeginTransaction();

                        Console.WriteLine("Attempting to dequeue message...");
                        // dequeue the message
                        EDBAQMessage deqMsg;
                        try
                        {
                            deqMsg = queueListen.Dequeue();
                        }
                        catch (Exception ex)
                        {
                            if (ex.Message.Contains("ORA-25228"))
                            {
                                Console.WriteLine("Message was not there.  Another process must have picked it up.");
                                Console.WriteLine();
                                txn.Rollback();
                                continue;
                            }
                            else
                            {
                                throw;
                            }
                        }

                        messageId = ByteArrayToString((byte[])deqMsg.MessageId);
                        if (deqMsg != null)
                        {
                            Console.WriteLine("Processing received message...");
                            // process the message payload
                            MyXML obj = new MyXML();
                            queueListen.Map<MyXML>(deqMsg.Payload, obj);

                            QueuedEntities.Message msg = new Utils.XmlFragmentSerializer<QueuedEntities.Message>().Deserialize(obj.value);

                            Console.WriteLine("Received Message:");
                            Console.WriteLine("MessageID: " + messageId);
                            Console.WriteLine("Message: " + msg.MessageText);
                            Console.WriteLine("Enqueue Time" + queueListen.MessageProperties.EnqueueTime);

                            txn.Commit();

                            Console.WriteLine("Finished processing message");
                            Console.WriteLine();

                        }
                        else
                        {
                            Console.WriteLine("Message was not dequeued.");
                        }
                    }
                    catch (Exception ex)
                    {
                        Console.WriteLine("Failed To dequeue or process the dequeued message.");
                        Console.WriteLine(ex.ToString());
                        Console.WriteLine();
                        if (txn != null)
                        {
                            txn.Rollback();
                            if (txn != null)
                            {
                                txn.Dispose();
                            }
                        }
                    }
                }
            }

        }
    }
}

EDBAQクラス

このアプリケーションでは、次のEDBAQクラスが使用されています。

EDBAQDequeueMode

EDBAQDequeueMode クラスは利用可能なすべてのデキューモードをリストします。

値

説明

ブラウズ

ロックせずにメッセージを読みます。

ロックされています

メッセージを読み取り、書き込みロックを取得します。

削除する

読んだ後にメッセージを削除します。これがデフォルト値です。

Remove_NoData

メッセージの受信を確認します。

EDBAQDequeueOptions

EDBAQDequeueOptions クラスはメッセージをデキューするときに利用できるオプションをリストします。

プロパティ

説明

ConsumerName

メッセージをデキューするコンシューマの名前。

DequeueMode

これはEDBAQDequeueModeから設定されます。これは、デキューオプションにリンクされたロック動作を表します。

ナビゲーション

これはEDBAQNavigationModeから設定されます。フェッチされるメッセージの位置を表します。

可視性

これはEDBAQVisibilityから設定されます。これは、新しいメッセージが現在のトランザクションの一部としてデキューされるかどうかを表します。

Wait

検索基準に基づくメッセージの待機時間。

Msgid

メッセージ識別子。

相関

相関識別子。

DeqCondition

デキューの状態。ブール式です。

変換

メッセージをデキューする前に適用される変換。

DeliveryMode

デキューされたメッセージの配信モード。

EDBAQEnqueueOptions

EDBAQEnqueueOptions クラスは、メッセージをエンキューするときに利用できるオプションをリストします。

プロパティ

説明

可視性

これはEDBAQVisibilityから設定されます。これは、新しいメッセージが現在のトランザクションの一部としてエンキューされるかどうかを表します。

RelativeMsgid

相対メッセージ識別子。

シーケンス偏差

メッセージをデキューする必要があるときのシーケンス。

変換

メッセージをエンキューする前に適用される変換。

DeliveryMode

エンキューされたメッセージの配信モード。

EDBAQMessage

EDBAQMessage クラスは、エンキュー/デキューされるメッセージをリストします。

プロパティ

説明

ペイロード

キューに入れられる実際のメッセージ。

MessageId

キューに入れられたメッセージのID。

EDBAQMessageProperties

EDBAQMessageProperties は利用可能なメッセージプロパティをリストします。

プロパティ

説明

優先

メッセージの優先度。

ディレイ

メッセージがデキューできる期間の投稿。これは秒単位で指定されます。

有効期限

メッセージがデキューに利用できる期間。これは秒単位で指定されます。

相関

相関識別子。

試み

メッセージのデキューを試行した回数。

RecipientList

デフォルトのキューサブスクライバーを破棄する受信者リスト。

ExceptionQueue

未処理のメッセージを移動するキューの名前。

EnqueueTime

メッセージがエンキューされた時刻。

状態

デキュー中のメッセージの状態。

OriginalMsgid

最後のキューのメッセージID。

TransactionGroup

デキューされたメッセージのトランザクショングループ。

DeliveryMode

デキューされたメッセージの配信モード。

EDBAQMessageState

EDBAQMessageState クラスは、デキュー中のメッセージの状態を表します。

値

説明

期限切れ

メッセージは例外キューに移動されます。

処理済み

メッセージは処理され、保持されます。

準備ができて

メッセージを処理する準備ができました。

Waiting

メッセージは待機状態です。遅延に達していません。

EDBAQMessageType

EDBAQMessageType クラスはペイロードのタイプを表します。

値

説明

Raw

Rawメッセージタイプ。

注:現在、このペイロードタイプはサポートされていません。

UDT

ユーザー定義タイプのメッセージ。

XML

XMLタイプのメッセージ。

注:現在、このペイロードタイプはサポートされていません。

EDBAQNavigationMode

EDBAQNavigationMode クラスは、利用可能なさまざまなタイプのナビゲーションモードを表します。

値

説明

First_Message

検索語に一致する最初の利用可能なメッセージを返します。

Next_Message

検索項目に一致する次の使用可能なメッセージを返します。

Next_Transaction

次のトランザクショングループの最初のメッセージを返します。

EDBAQQueue

EDBAQQueue クラスはPostgreSQLデータベースで DMBS_AQ 機能を実行するSQLステートメントを表します。

プロパティ

説明

接続

使用する接続。

名前

キューの名前。

MessageType

このキューからエンキューまたはデキューされるメッセージタイプ。たとえば、EDBAQMessageType.Udtです。

UdtTypeName

メッセージタイプのユーザ定義タイプ名。

EnqueueOptions

使用するエンキューオプション。

DequeuOptions

使用するデキューオプション。

MessageProperties

使用するメッセージプロパティ。

EDBAQVisibility

EDBAQVisibility クラスは利用可能な可視性オプションを表します。

値

説明

Immediate

エンキュー/デキューは、進行中のトランザクションの一部ではありません。

On_Commit

エンキュー/デキューは、現在のトランザクションの一部です。

注釈