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

EDB Postgres Advanced Server Advanced Queuingは、Advanced Serverデータベースのメッセージキューとメッセージ処理を提供します。ユーザー定義のメッセージはキューに保存されます。キューのコレクションはキューテーブルに格納されます。キューテーブルを作成してから、それに依存するキューを作成する必要があります。

サーバー側では、``DBMS_AQADM``パッケージのプロシージャがメッセージキューとキューテーブルを作成および管理します。``DBMS_AQ``パッケージを使用して、キューにメッセージを追加または削除したり、PL / SQLコールバックプロシージャを登録または登録解除します。``DBMS_AQ``および``DBMS_AQADM``の詳細については、``ここをクリック``<https://www.enterprisedb.com/docs/en/11.0/EPAS_BIP_Guide_v11/Database_Compatibility_for_Oracle_Developers_Built-in_Package_Guide.1.14.html#pID0E01HG0HA> `_。

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

メッセージをエンキューまたはデキューする

サービス側のセットアップ

.NETアプリケーションでAdvanced Queuing機能を使用するには、最初にユーザー定義型、キューテーブル、およびキューを作成してから、データベースサーバーでキューを開始する必要があります。 EDB-PSQLを呼び出して、Advanced Serverホストデータベースに接続します。コマンドラインで次の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から設定されます。現在のトランザクションの一部として、新しいメッセージがデキューされるかどうかを表します。

待って

検索条件ごとのメッセージの待機時間。

Msgid

メッセージ識別子。

相関関係

相関識別子。

DeqCondition

デキュー条件。これはブール式です。

変換

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

DeliveryMode

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

**EDBAQEnqueueOptions **

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

プロパティ

説明

可視性

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

RelativeMsgid

相対メッセージ識別子。

SequenceDeviation

メッセージをデキューするシーケンス。

変換

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

DeliveryMode

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

**EDBAQMessage **

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

プロパティ

説明

ペイロード

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

MessageId

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

**EDBAQMessageProperties **

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

プロパティ

説明

優先順位

メッセージの優先度。

遅延

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

有効期限

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

相関関係

相関識別子。

試行

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

RecipientList

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

ExceptionQueue

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

EnqueueTime

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

都道府県

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

OriginalMsgid

最後のキューのメッセージ識別子。

TransactionGroup

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

DeliveryMode

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

**EDBAQMessageState **

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

値

説明

期限切れ

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

処理済み

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

準備完了

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

待っています

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

**EDBAQMessageType **

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

値

説明

生

生のメッセージタイプ。

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

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``クラスは利用可能な可視性オプションを表します。

値

説明

すぐに

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

On_Commit

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

注釈