アドバンスドキューイングの使用¶
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アプリケーションでメッセージをエンキューするには、以下を行う必要があります。
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から設定されます。これは、新しいメッセージが現在のトランザクションの一部としてデキューされるかどうかを表します。 |
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 |
エンキュー/デキューは、現在のトランザクションの一部です。 |
注釈
上記のパラメータのデフォルトオプションを確認するには、``ここをクリックしてください<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では使用できません。