Bri*_*lig 6 c# asynchronous producer-consumer
我正在使用生产者/消费者模式实现数据链路层.数据链路层有自己的线程和状态机,通过线路传输数据链路协议(以太网,RS-232 ......).物理层的接口表示为System.IO.Stream.另一个线程将消息写入数据链接对象并从中读取消息.
数据链接对象具有空闲状态,必须等待以下四种情况之一:
我很难找到最好的方法来实现这一点,而无需将通信分成读/写线程(从而大大增加了复杂性).以下是我如何获得4分中的3分:
// Read a byte from 'stream'. Timeout after 10 sec. Monitor the cancellation token.
stream.ReadTimeout = 10000;
await stream.ReadAsync(buf, 0, 1, cts.Token);
Run Code Online (Sandbox Code Playgroud)
要么
BlockingCollection<byte[]> SendQueue = new ...;
...
// Check for a message from network layer. Timeout after 10 seconds.
// Monitor cancellation token.
SendQueue.TryTake(out msg, 10000, cts.Token);
Run Code Online (Sandbox Code Playgroud)
我该怎么做才能阻止线程,等待所有四个条件?欢迎所有建议.我没有设置任何架构或数据结构.
编辑:********感谢大家的帮助.这是我的解决方案********
首先,我认为没有生产者/消费者队列的异步实现.所以我实现了类似于这个stackoverflow帖子的东西.
我需要一个外部和内部取消源来分别停止使用者线程并取消中间任务,类似于本文.
byte[] buf = new byte[1];
using (CancellationTokenSource internalTokenSource = new CancellationTokenSource())
{
CancellationToken internalToken = internalTokenSource.Token;
CancellationToken stopToken = stopTokenSource.Token;
using (CancellationTokenSource linkedCts =
CancellationTokenSource.CreateLinkedTokenSource(stopToken, internalToken))
{
CancellationToken ct = linkedCts.Token;
Task<int> readTask = m_stream.ReadAsync(buf, 0, 1, ct);
Task<byte[]> msgTask = m_sendQueue.DequeueAsync(ct);
Task keepAliveTask = Task.Delay(m_keepAliveTime, ct);
// Wait for at least one task to complete
await Task.WhenAny(readTask, msgTask, keepAliveTask);
// Next cancel the other tasks
internalTokenSource.Cancel();
try {
await Task.WhenAll(readTask, msgTask, keepAliveTask);
} catch (OperationCanceledException e) {
if (e.CancellationToken == stopToken)
throw;
}
if (msgTask.IsCompleted)
// Send the network layer message
else if (readTask.IsCompleted)
// Process the byte from the physical layer
else
Contract.Assert(keepAliveTask.IsCompleted);
// Send a keep alive message
}
}
Run Code Online (Sandbox Code Playgroud)
在这种情况下,我只会使用取消令牌进行取消。像保持活动计时器这样的重复超时最好表示为计时器。
因此,我将其建模为三个可取消的任务。首先,取消令牌:
所有通信均被网络层取消
CancellationToken token = ...;
Run Code Online (Sandbox Code Playgroud)
然后,三个并发操作:
收到一个字节
var readByteTask = stream.ReadAsync(buf, 0, 1, token);
Run Code Online (Sandbox Code Playgroud)
保活定时器已过期
var keepAliveTimerTask = Task.Delay(TimeSpan.FromSeconds(10), token);
Run Code Online (Sandbox Code Playgroud)
网络线程有消息可用
这个有点棘手。您当前的代码使用BlockingCollection<T>,它不兼容异步。我建议切换到TPL DataflowBufferBlock<T>或我自己的AsyncProducerConsumerQueue<T>,其中任何一个都可以用作异步兼容的生产者/消费者队列(意味着生产者可以是同步或异步的,消费者可以是同步或异步的)。
BufferBlock<byte[]> SendQueue = new ...;
...
var messageTask = SendQueue.ReceiveAsync(token);
Run Code Online (Sandbox Code Playgroud)
然后您可以使用Task.WhenAny来确定完成了哪些任务:
var completedTask = await Task.WhenAny(readByteTask, keepAliveTimerTask, messageTask);
Run Code Online (Sandbox Code Playgroud)
现在,您可以通过completedTask与其他结果进行比较并await对其进行 ing 来检索结果:
if (completedTask == readByteTask)
{
// Throw an exception if there was a read error or cancellation.
await readByteTask;
var byte = buf[0];
...
// Continue reading
readByteTask = stream.ReadAsync(buf, 0, 1, token);
}
else if (completedTask == keepAliveTimerTask)
{
// Throw an exception if there was a cancellation.
await keepAliveTimerTask;
...
// Restart keepalive timer.
keepAliveTimerTask = Task.Delay(TimeSpan.FromSeconds(10), token);
}
else if (completedTask == messageTask)
{
// Throw an exception if there was a cancellation (or the SendQueue was marked as completed)
byte[] message = await messageTask;
...
// Continue reading
messageTask = SendQueue.ReceiveAsync(token);
}
Run Code Online (Sandbox Code Playgroud)
| 归档时间: |
|
| 查看次数: |
629 次 |
| 最近记录: |