【问题标题】:How to use Reactive Extensions to parse a stream of characters from a serial port?如何使用 Reactive Extensions 解析来自串行端口的字符流?
【发布时间】:2014-06-26 07:19:02
【问题描述】:

我需要解析来自测试仪器的串行数据流,这似乎是响应式扩展的出色应用。

协议非常简单...每个“数据包”都是一个字母后跟数字。每种数据包类型的数字位数是固定的,但可能因数据包类型而异。例如

...A1234B123456C12...

我正在尝试将其分解为字符串的 Observable,例如

“A1234”“B123456”“C12”...

认为这很简单,但没有看到明显的解决方法(我对 LINQ 有一些经验,但对 Rx 很陌生)。

这是我到目前为止的代码,它从串行端口的 SerialDataReceived 事件中生成 Observable 字符。

        var serialData = Observable
                            .FromEventPattern<SerialDataReceivedEventArgs>(SerialPortMain, "DataReceived")
                            .SelectMany(_ => 
                            {
                                int dataLength = SerialPortMain.BytesToRead;
                                byte[] data = new byte[dataLength];
                                int nbrDataRead = SerialPortMain.Read(data, 0, dataLength);

                                if (nbrDataRead == 0)
                                    return  new char[0]; 

                                 var chars = Encoding.ASCII.GetChars(data);

                                 return chars; 
                            });

如何将serialData转换为String的Observable,其中每个字符串都是一个数据包?

【问题讨论】:

    标签: c# serial-port system.reactive reactive-programming


    【解决方案1】:

    这里有一个稍微短一点的方法,与James' first solution 风格相同,使用类似的辅助方法:

    public static bool IsCompletePacket(string s)
    {
        switch (s[0])
        {
            case 'A':
                return s.Length == 5;
            case 'B':
                return s.Length == 6;
            case 'C':
                return s.Length == 7;
            default:
                throw new ArgumentException("Packet must begin with a letter");
        }
    }
    

    那么代码是:

    var packets = chars
        .Scan(string.Empty, (prev, cur) => char.IsLetter(cur) ? cur.ToString() : prev + cur)
        .Where(IsCompletePacket);
    

    Scan 部分构建以字母结尾的字符串,例如:

    A
    A1
    A12
    A123
    ...
    

    Where 然后只选择那些长度正确的。本质上它只是从 James' 中删除元组并使用字符串长度代替。

    【讨论】:

    • 是的,这是我书中的赢家。我将 IsCompletePacket 中的默认情况更改为返回 false,所以现在它只是忽略代码不支持的任何数据包。
    【解决方案2】:

    令人惊讶的繁琐!我已经解决了这几种方法:

    // helper method to get the packet length
    public int GetPacketLength(char c)
    {
        switch(c)
        {
            case 'A':
                return 5;
            case 'B':
                return 6;
            case 'C':
                return 7;
            default:
                throw new Exception("Unknown packet code");
        }
    }
    

    那么我们可以这样做:

    // chars is a test IObservable<char> 
    string[] messages = { "A1234", "B12345", "C123456" };
    var serialPort = Enumerable.Range(1, 10).ToObservable();
    var chars = serialPort.SelectMany((_, i) => messages[i % 3]);
    
    var packets = chars.Scan(
        Tuple.Create(string.Empty, -1),
        (acc, c) =>
            Char.IsLetter(c)
                ? Tuple.Create(c.ToString(), GetPacketLength(c) - 1)
                : Tuple.Create(acc.Item1 + c, acc.Item2 - 1))
        .Where(acc => acc.Item2 == 0)
        .Select(acc => acc.Item1)
        .Subscribe(Console.WriteLine);
    

    它的作用是这样的:

    • Scan 构建每个数据包并将其与数据包中剩余的字符数配对:例如("A",4) ("A1",3) ("A12",2) ("A123",1) ("A1234",0) ("B",5) ...
    • 然后我们知道剩下 0 个字符的对是我们需要的,所以我们使用 Where 过滤掉其余部分,并使用 Select 从对中过滤出结果

    另类

    这是另一种方法,功能较少。从风格的角度来看,我喜欢上面的代码,但下面的代码在内存方面更有效——万一它会有所作为。

    public static class ObservableExtensions
    {
        private const int MaxPacketLength = 7;
        private static Dictionary<char, int> PacketLengthTable =
            new Dictionary<char, int> { {'A', 5}, {'B', 6}, {'C', 7 } };
    
        public static IObservable<string> GetPackets(this IObservable<char> source)
        {
            return Observable.Create<string>(o =>
            {
                var currentPacketLength = 0;
                var buffer = new char[MaxPacketLength];
                var index = -1;
                return source.Subscribe(
                    c => {
                        if (Char.IsLetter(c))
                        {
                            currentPacketLength = PacketLengthTable[c];
                            buffer[0] = c;
                            index = 0;
                        }
                        else if(index >= 0)
                        {
                            index++;
                            buffer[index] = c;
                        }
                        if (index == currentPacketLength - 1)
                        {
                            o.OnNext(new string(buffer,0, currentPacketLength));
                            index = -1;
                        }
                    },
                    o.OnError,
                    o.OnCompleted);
            });
        }
    }
    

    而且可以这样使用:

    // chars is a test IObservable<char> 
    string[] messages = { "A1234", "B12345", "C123456" };
    var serialPort = Enumerable.Range(1, 10).ToObservable();
    var chars = serialPort.SelectMany((_, i) => messages[i % 3]);
    
    var packets = chars.GetPackets().Subscribe(Console.WriteLine);
    

    【讨论】:

    • 谢谢詹姆斯,这似乎工作。我选择了你的第二种方法。虽然有点“unRx-ey”,但我同意它更容易理解。
    【解决方案3】:

    这是我对这个问题的看法。

    static IObservable<string> Packets(IObservable<char> source)
    {
        return Observable.Create<string>(observer =>
        {
            var packet = new List<char>();
            Action emitPacket = () =>
            {
                if (packet.Count > 0)
                {
                    observer.OnNext(new string(packet.ToArray()));
                    packet.Clear();
                }
            };
            return source.Subscribe(
                c =>
                {
                    if (char.IsLetter(c))
                    {
                        emitPacket();
                    }
                    packet.Add(c);
                },
                observer.OnError,
                () =>
                {
                    emitPacket();
                    observer.OnCompleted();
                });
        });
    }
    

    如果输入是字符A425B90C2DX812,则输出是A425B90C2DX812

    请注意,您无需预先指定数据包长度或数据包类型(起始字母)。

    您可以使用更通用的扩展方法来实现相同的方法:

    static IObservable<IList<T>> GroupSequential<T>(
        this IObservable<T> source, Predicate<T> isFirst)
    {
        return Observable.Create<T>(observer =>
        {
            var group = new List<T>();
            Action emitGroup = () =>
            {
                if (group.Count > 0)
                {
                    observer.OnNext(group.ToList());
                    group.Clear();
                }
            };
            return source.Subscribe(
                item =>
                {
                    if (isFirst(item))
                    {
                        emitGroup();
                    }
                    group.Add(item);
                },
                observer.OnError,
                () =>
                {
                    emitGroup();
                    observer.OnCompleted();
                });
        });
    }
    

    Packets 的实现就是:

    static IObservable<string> Packets(IObservable<char> source)
    {
        return source
            .GroupSequential(char.IsLetter)
            .Select(x => new string(x.ToArray()));
    }
    

    【讨论】:

    • 我喜欢它适用于任何数据包长度的事实。但是在遵循逻辑时遇到了一些麻烦 - Rx 对我来说是一种非常不同的思维方式。
    • @TomBushell 它有助于“从里到外”阅读方法。实际情况是,当有人订阅Packets(source) 时,实际上只是订阅了source。 1. 每当source 发出char c 时,我们检查当前数据包是否完成:如果是,我们发出它并开始一个新数据包。我们总是将当前字符添加到当前数据包中。 2. 当source 发出错误时,我们只是转发它。 3. 当source 完成后,我们确保发出最后一个数据包(如果有的话)。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-11-08
    • 2013-11-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-01-24
    相关资源
    最近更新 更多