アドバンスドキューイングの使用¶
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アプリケーションでメッセージをエンキューするには、次のことを行う必要があります。
``EnterpriseDB.EDBClient``名前空間をインポートします。
キューの名前を渡し、``EDBAQQueue``のインスタンスを作成します。
エンキューメッセージを作成し、ペイロードを定義します。
``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アプリケーションでメッセージをデキューするには、次の手順を実行する必要があります。
``EnterpriseDB.EDBClient``名前空間をインポートします。
キューの名前を渡し、``EDBAQQueue``のインスタンスを作成します。
``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 |
エンキュー/デキューは、現在のトランザクションの一部です。 |
注釈
上記のパラメータのデフォルトオプションを確認するには、``ここをクリック``<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/> `_。
EDBAQ機能は、ユーザー定義タイプを使用してエンキュー/デキュー操作を呼び出します。``NoTypeLoading``はユーザー定義タイプをロードしないため、``Server Compatibility Mode = NoTypeLoading``はEDBAQで使用できません。