【问题标题】:Sending stream (e.g. file) from actor to actor从演员到演员发送流(例如文件)
【发布时间】:2017-07-27 19:41:23
【问题描述】:

我想将数据从一个actor中的流发送到另一个actor中的流。最后,这应该与远程参与者一起工作。使用 Akka.NET Streams 这应该是一件容易的事,但也许我误解了它。

这是 SenderActor 的一部分:

var stream = new FileStream(FILENAME, FileMode.Open);
var materializer = Context.Materializer();
var source = StreamConverters.FromInputStream(() => stream, CHUNK_SIZE);
var result = source.To(Sink.ActorRef<ByteString>(this.receiver, new Messages.StreamComplete()))
    .Run(materializer);

注意:this.receiver 是一个 IActorRef,数据应该发送到该地址。

现在 ReceiverActor 会在流结束后获取所有 ByteString 消息和 Messages.StreamCompleted

如何轻松地将它放在 ReceiverActor 中?在最好的情况下再次作为Stream

ReceiverActor中,我尝试将所有ByteString 消息发送到Source,而Source 应该填写MemoryStream

class ReceiverActor : ReceiveActor
{
    private readonly IActorRef streamReceiver;
    private readonly MemoryStream stream;

    public ReceiverActor()
    {
        this.stream = new MemoryStream();
        this.streamReceiver = Source.ActorRef<ByteString>(128, Akka.Streams.OverflowStrategy.Fail)
            .To(StreamConverters.FromOutputStream(() => this.stream, true))
            .Run(Context.Materializer());
        Context.Watch(this.streamReceiver);

        Receive<ByteString>((message) => ReceivedStreamChunk(message));
        Receive<Messages.StreamComplete>((message) => ReceivedStreamComplete(message));
        Receive<Terminated>((message) => ReceivedTerminated(message));
        ReceiveAny((message) => ReceivedAnyMessage(message));
    }
    private void ReceivedTerminated(Terminated message)
    {
        Console.WriteLine($"[receiver] {message.ActorRef.Path.ToStringWithoutAddress()} terminated, local stream length {(this.stream.CanRead ? this.stream.Length : -1)}");
    }
    private void ReceivedStreamComplete(Messages.StreamComplete message)
    {
        Console.WriteLine($"[receiver] got signaled that the stream completed, local stream length {this.stream.Length}");
    }
    private void ReceivedStreamChunk(object message)
    {
        Console.WriteLine($"[receiver] got chunk, previous stream length {this.stream.Length}");
        this.streamReceiver.Forward(message);
    }
    private void ReceivedAnyMessage(object message)
    {
        Console.WriteLine($"[receiver] got message {message.GetType().FullName}");
    }
}

但是MemoryStream 是异步填充的,当streamReceiver 终止时它会关闭流,因此我无法获取数据。

如何正确检索流?


更新我让它在本地工作:

感谢 Horusiath 在 Akka.NET 的 gitter 频道中的输入,我能够直接收到 ByteStrings,而不需要在那里使用 Akka.Stream。

Receive<ByteString>((message) => ReceivedStreamChunk(message));

还有:

private void ReceivedStreamChunk(ByteString message)
{
    var bytes = message.ToArray();
    targetStream.Write(bytes, 0, bytes.Length);
    Console.WriteLine($"[receiver] got chunk, now stream length {this.stream.Length}");
}

请注意,Akka.NET 1.3 会将方法 WriteTo(Stream)WriteToAsync(Stream, CancellationToken) 添加到 ByteString

这仍然不适用于远程演员,因为接收演员系统收到此错误(序列化器是 Hyperion):

错误[没有为此对象定义无参数构造函数。]

其实ByteString有一个无参数的构造函数但它是protected

我认为ByteString 不可序列化?

【问题讨论】:

    标签: c# akka.net akka.net-streams


    【解决方案1】:

    通过在源和接收器之间插入转换流,将ByteString 转换为byte[] 时,我让它与远程参与者一起工作。

    例子:

    FileIO.FromFile(...).Via(Flow.Create&lt;ByteString&gt;().Select(x =&gt; x.ToArray())).To(...)

    【讨论】:

    • 我来这个问题是因为我有同样的需求:将文件从一个演员发送到另一个远程演员。我也找到了这个答案:stackoverflow.com/a/49545111/1060314 特别是 一般来说,使用 actor refs 的源和接收器并未设计为通过远程连接工作 - 它们不包括消息重试,如果某些流可能会导致系统死锁控制消息不会被传入。这似乎是 1.4 中的 StreamsRefs 解决的。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-08
    • 2015-06-29
    • 1970-01-01
    • 2015-09-25
    • 1970-01-01
    • 2017-01-24
    相关资源
    最近更新 更多