【问题标题】:Inline parsing of IObservable<byte>IObservable<byte>的内联解析
【发布时间】:2014-09-27 12:43:52
【问题描述】:

我有一个可观察的查询,它从我想要内联解析的流中生成IObservable&lt;byte&gt;。我希望能够根据数据源使用不同的策略来解析来自该序列的离散消息。请记住,我仍在 RX 的向上学习曲线上。我想出了一个解决方案,但不确定是否有办法使用开箱即用的运算符来完成此任务。

首先,我给IObservable写了如下扩展方法:

    public static IObservable<IList<T>> Parse<T>(
        this IObservable<T> source,
        Func<IObservable<T>, IObservable<IList<T>>> parsingFunction)
    {
        return parsingFunction(source);
    }

这允许我指定特定数据源使用的消息框架策略。一个数据源可能由一个或多个字节分隔,而另一个数据源可能由开始和停止块模式分隔,而另一个可能使用长度前缀策略。所以这里是我定义的分隔策略的一个例子:

public static class MessageParsingFunctions
{

    public static Func<IObservable<T>, IObservable<IList<T>>> Delimited<T>(T[] delimiter)
    {
        if (delimiter == null) throw new ArgumentNullException("delimiter");
        if (delimiter.Length < 1) throw new ArgumentException("delimiter must contain at least one element.");

        Func<IObservable<T>, IObservable<IList<T>>> parser =
            (source) =>
            {
                var shared = source.Publish().RefCount();

                var windowOpen = shared.Buffer(delimiter.Length, 1)
                    .Where(buffer => buffer.SequenceEqual(delimiter))
                    .Publish()
                    .RefCount();

                return shared.Buffer(windowOpen)
                    .Select(bytes =>
                        bytes
                        .Take(bytes.Count - delimiter.Length)
                        .ToList());

            };

        return parser;
    }
}

因此,最终,作为示例,我可以在序列中遇到字符串 '&lt;EOF&gt;' 的字节模式时,以以下方式使用代码解析序列中的离散消息:

var messages = ...operators that surface an IObservable<byte>
    .Parse(MessageParsingFunctions.Delimited(Encoding.ASCII.GetBytes("<EOF>")))
    ...further operators to package discrete messages along with additional metadata

问题:

  1. 是否有更直接的方法来使用开箱即用的运算符来完成此任务?
  2. 如果不是,是否最好将不同的解析函数(即 ParseDelimited、ParseLengthPrefixed 等)定义为本地扩展,而不是使用接受解析函数的更通用的 Parse 扩展方法?

提前致谢!

【问题讨论】:

  • 1.并不真地。 2.这只是对API设计的看法。我个人的回答(这有点符合 Rx 的设计)是“为什么不两者兼而有之?”通用Parse,以便开发人员可以轻松地使用新的选择器进行扩展,并将更具体的预定义扩展作为准备使用的解析器库。
  • 老实说,乍一看,您可能做得比您已经发布的解决方案差得多;任何“答案”最终都会成为一种Func&lt;IObservable,Chunk&gt; 构造,您可以在其中根据流上下文通过算法决定何时发布块。当我不在手机上时,我会尝试想出一个替代方案来进行比较。

标签: system.reactive


【解决方案1】:

看看Rxx Parsers。这是related lab。例如:

IObservable<byte> bytes = ...;

var parsed = bytes.ParseBinary(parser =>
  from next in parser
  let magicNumber = parser.String(Encoding.UTF8, 3).Where(value => value == "RXX")
  let header = from headerLength in parser.Int32
               from header in next.Exactly(headerLength)
               from headerAsString in header.Aggregate(string.Empty, (s, b) => s + " " + b)
               select headerAsString
  let message = parser.String(Encoding.UTF8)
  let entry = from length in parser.Int32
              from data in next.Exactly(length)
              from value in data.Aggregate(string.Empty, (s, b) => s + " " + b)
              select value
  let entries = from count in parser.Int32
                from entries in entry.Exactly(count).ToList()
                select entries
  select from _ in magicNumber.Required("The file's magic number is invalid.")
         from h in header.Required("The file's header is invalid.")
         from m in message.Required("The file's message is invalid.")
         from e in entries.Required("The file's data is invalid.")
         select new
         {
           Header = h,
           Message = m,
           Entries = e.Aggregate(string.Empty, (acc, cur) => acc + cur + Environment.NewLine)
         });

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-03-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-09-30
    • 2023-03-12
    • 2017-11-28
    相关资源
    最近更新 更多