C++ サーバーと C# WPF UI の 2 つのアプリケーションがあります。C++ コードは、ZeroMQ メッセージング [PUB/SUB] サービスを介して (どこからでも誰からでも) 要求を受け取ります。私は自分の C# コードをバック テストに使用し、「バック テスト」を作成して実行しています。これらのバック テストは、多くの「単体テスト」で構成でき、それぞれが C++ サーバーから何千ものメッセージを送受信します。
現在、個々のバック テストはうまく機能しており、それぞれ数千のリクエストとキャプチャを含む N 個の単体テストを送信できます。私の問題は建築です。別のバック テスト (最初のテストに続いて) をディスパッチすると、ポーリング スレッドがキャンセルおよび破棄されないために、イベント サブスクリプションが 2 回行われるという問題が発生します。これにより、誤った出力が発生します。これは些細な問題のように思えるかもしれませんが (一部の人にとってはそうかもしれません)、現在の構成でこのポーリング タスクをキャンセルするのは面倒です。いくつかのコード...
私のメッセージブローカークラスはシンプルで、次のようになります
public class MessageBroker : IMessageBroker<Taurus.FeedMux>, IDisposable
{
private Task pollingTask;
private NetMQContext context;
private PublisherSocket pubSocket;
private CancellationTokenSource source;
private CancellationToken token;
private ManualResetEvent pollerCancelled;
public MessageBroker()
{
this.source = new CancellationTokenSource();
this.token = source.Token;
StartPolling();
context = NetMQContext.Create();
pubSocket = context.CreatePublisherSocket();
pubSocket.Connect(PublisherAddress);
}
public void Dispatch(Taurus.FeedMux message)
{
pubSocket.Send(message.ToByteArray<Taurus.FeedMux>());
}
private void StartPolling()
{
pollerCancelled = new ManualResetEvent(false);
pollingTask = Task.Run(() =>
{
try
{
using (var context = NetMQContext.Create())
using (var subSocket = context.CreateSubscriberSocket())
{
byte[] buffer = null;
subSocket.Options.ReceiveHighWatermark = 1000;
subSocket.Connect(SubscriberAddress);
subSocket.Subscribe(String.Empty);
while (true)
{
buffer = subSocket.Receive();
MessageRecieved.Report(buffer.ToObject<Taurus.FeedMux>());
if (this.token.IsCancellationRequested)
this.token.ThrowIfCancellationRequested();
}
}
}
catch (OperationCanceledException)
{
pollerCancelled.Set();
}
}, this.token);
}
private void CancelPolling()
{
source.Cancel();
pollerCancelled.WaitOne();
pollerCancelled.Close();
}
public IProgress<Taurus.FeedMux> MessageRecieved { get; set; }
public string PublisherAddress { get { return "tcp://127.X.X.X:6500"; } }
public string SubscriberAddress { get { return "tcp://127.X.X.X:6501"; } }
private bool disposed = false;
protected virtual void Dispose(bool disposing)
{
if (!disposed)
{
if (disposing)
{
if (this.pollingTask != null)
{
CancelPolling();
if (this.pollingTask.Status == TaskStatus.RanToCompletion ||
this.pollingTask.Status == TaskStatus.Faulted ||
this.pollingTask.Status == TaskStatus.Canceled)
{
this.pollingTask.Dispose();
this.pollingTask = null;
}
}
if (this.context != null)
{
this.context.Dispose();
this.context = null;
}
if (this.pubSocket != null)
{
this.pubSocket.Dispose();
this.pubSocket = null;
}
if (this.source != null)
{
this.source.Dispose();
this.source = null;
}
}
disposed = true;
}
}
public void Dispose()
{
Dispose(true);
GC.SuppressFinalize(this);
}
~MessageBroker()
{
Dispose(false);
}
}
バックテストの「エンジン」は、各バック テストを実行するために使用します。最初に、各テスト (ユニット テスト)Dictionary
を含むと、各テストの C++ アプリケーションにディスパッチするメッセージを作成します。Test
DispatchTests
方法は、こちら
private void DispatchTests(ConcurrentDictionary<Test, List<Taurus.FeedMux>> feedMuxCollection)
{
broker = new MessageBroker();
broker.MessageRecieved = new Progress<Taurus.FeedMux>(OnMessageRecieved);
testCompleted = new ManualResetEvent(false);
try
{
// Loop through the tests.
foreach (var kvp in feedMuxCollection)
{
testCompleted.Reset();
Test t = kvp.Key;
t.Bets = new List<Taurus.Bet>();
foreach (Taurus.FeedMux mux in kvp.Value)
{
token.ThrowIfCancellationRequested();
broker.Dispatch(mux);
}
broker.Dispatch(new Taurus.FeedMux()
{
type = Taurus.FeedMux.Type.PING,
ping = new Taurus.Ping() { event_id = t.EventID }
});
testCompleted.WaitOne(); // Wait until all messages are received for this test.
}
testCompleted.Close();
}
finally
{
broker.Dispose(); // Dispose the broker.
}
}
最後のPING
メッセージは、終了したことを C++ に伝えるためのものです。次に、強制的に待機させて、C++ コードからすべての戻り値を受け取る前に次の [ユニット] テストがディスパッチされないようにしますManualResetEvent
。
C++ が PING メッセージを受信すると、メッセージをそのまま送り返します。受信したメッセージを 経由で処理し、ユニット テストを続行できるようにOnMessageRecieved
を設定するように PING から指示されます。ManualResetEvent.Set()
"次の方"...
private async void OnMessageRecieved(Taurus.FeedMux mux)
{
string errorMsg = String.Empty;
if (mux.type == Taurus.FeedMux.Type.MSG)
{
// Do stuff.
}
else if (mux.type == Taurus.FeedMux.Type.PING)
{
// Do stuff.
// We are finished reciving messages for this "unit test"
testCompleted.Set();
}
}
私の問題は。broker.Dispose()
、最後に上記がヒットしないことです。バックグラウンド スレッドで実行される finally ブロックが実行されるとは限りません
上記の取り消し線のテキストは、私がコードをいじったためです。子が完了する前に親スレッドを停止していました。ただし、まだ問題があります...
Nowbroker.Dispose()
が正しく呼び出され、呼び出されます。このメソッドでは、複数のサブスクリプションを回避するためbroker.Dispose()
に、ポーラー スレッドをキャンセルし、正しく破棄しようとします。Task
スレッドをキャンセルするには、CancelPolling()
メソッドを使用します
private void CancelPolling()
{
source.Cancel();
pollerCancelled.WaitOne(); <- Blocks here waiting for cancellation.
pollerCancelled.Close();
}
しかし、StartPolling()
方法では
while (true)
{
buffer = subSocket.Receive();
MessageRecieved.Report(buffer.ToObject<Taurus.FeedMux>());
if (this.token.IsCancellationRequested)
this.token.ThrowIfCancellationRequested();
}
ThrowIfCancellationRequested()
が呼び出されることはなく、スレッドがキャンセルされることもないため、適切に破棄されることはありません。ポーラー スレッドがsubSocket.Receive()
メソッドによってブロックされています。
今、私が望むものを達成する方法が明確ではありません。メッセージをポーリングするために使用されるスレッド以外のスレッドでbroker.Dispose()
/を呼び出す必要があり、キャンセルを強制する方法もあります。PollerCancel()
スレッドの中止は、私がどうしてもやりたいことではありません。
基本的に、次のバック テストを実行する前に適切に破棄したいのですが、broker
これを正しく処理し、ポーリングを分割して別のアプリケーション ドメインで実行するにはどうすればよいですか?
ハンドラー内で破棄しようとしOnMessageRecived
ましたが、これは明らかにポーラーと同じスレッドで実行され、追加のスレッドを呼び出さずにブロックする方法ではありません。
私が望むものを達成するための最良の方法は何ですか?私が従うことができるこの種のケースのパターンはありますか?
御時間ありがとうございます。