【问题标题】:Split IObservable<byte[]> to characters then to line将 IObservable<byte[]> 拆分为字符,然后拆分为行
【发布时间】:2015-07-04 03:29:15
【问题描述】:

Rx 很棒,但有时很难找到优雅的方式来做某事。这个想法很简单。我收到带有 byte[] 的事件,这个数组可能包含一行、多行或一行的一部分。我想要的是找到一种方法来获得 Line 的 IObservable 所以IObservable&lt;String&gt;,其中序列的每个元素都是一行。

几个小时后,我发现的最接近的解决方案非常难看,而且当然不起作用,因为 scan 在每个字符上触发 OnNext :

//Intermediate subject use to transform byte[] into char
var outputStream = new Subject<char>();
_reactiveSubcription = outputStream
    //Scan doesn't work it trigger OnNext on every char
    //Aggregate doesn't work neither as it doesn't return intermediate result
    .Scan(new StringBuilder(), (builder, c) => c == '\r' ? new StringBuilder() : builder.Append((char)c))
    .Subscribe(this);


Observable.FromEventPattern<ShellDataEventArgs>(shell, "DataReceived")
            //Data is a byte[]
            .Select(_ => _.EventArgs.Data)
            .Subscribe(array => array.ToObservable()
            //Convert into char
            .ForEach(c => outputStream.OnNext((char)c)));

注意:_reactiveSubcription 应该是IObservable&lt;String&gt;

在不考虑字符编码问题的情况下,我缺少什么来完成这项工作?

【问题讨论】:

  • 你在一个优雅的块中尝试做的事情在我的真实世界站中从未奏效。您如何知道 byte[] 已全部接收或它的结束字节在哪里?即,一条 2K 的消息是否会在一行中间的两个不同事件中以两个 1K 的块到达。我总是必须将字节弹出到理论上的无限队列中,并在这个其他构建的数组中查找我的 NewLine 或 CrLf,然后在 NewLines 和 CrLf 字符出现时将各个 Lines 拉出。如果没有 NewLine 或 CrLf 出现在我的 MaxLineSize 值中,我必须放置安全代码以清除构建的数组。

标签: c# system.reactive


【解决方案1】:

这对我有用。

首先,将 byte[] 转换为字符串并在 \r 上拆分字符串(Regex Split 保留分隔符)。

现在有一个字符串流,其中一些以\r 结尾。

然后是 Concat,以使它们保持有序。此外,由于strings 需要在下一步中处于热门状态,因此请发布它们。

var strings = bytes.
  Select(arr => (Regex.Split(Encoding.Default.GetString(arr, 0, arr.Length - 1), "(\r)")).
    Where(s=> s.Length != 0).
    ToObservable()).
  Concat().
  Publish().
  RefCount();

创建一个以\r 结尾的字符串时结束的字符串窗口。 strings 需要是热的,因为它同时用于窗口内容和窗口结束触发器。

var linewindows = strings.Window(strings.Where(s => s.EndsWith("\r")));

将每个窗口聚合成一个字符串。

var lines = linewindows.SelectMany(w => w.Aggregate((l, r) => l + r));

linesIObservable&lt;String&gt;,每个字符串包含一行。

为了测试这一点,我使用以下生成器生成 IObservable&lt;byte[]&gt;

var bytes = Observable.
Range(1, 10).
SelectMany(i => Observable.
    Return((byte)('A' + i)).
    Repeat(24).
    Concat(Observable.
        Return((byte)'\r'))).
Window(17).
SelectMany(w => w.ToArray());

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-07-27
    • 2019-01-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多