【问题标题】:rxjs operators to collect strings from a source and partially emit them according to a patternrxjs 运算符从源收集字符串并根据模式部分发出它们
【发布时间】:2018-10-15 13:01:07
【问题描述】:

所以我在一个支持蓝牙的项目上使用 typescript/RXJS/React-Native;我有一个从给定外围设备接收字符串的函数,但有一些我无法从这个外围设备中真正避免的警告。

首先,它通过某种命令模式与我交流;也就是说,每个命令的格式为/[a-zA-Z][^;]*;/(又名一个字母字符,后跟任意数量的字符,以分号结尾)。

这些命令后面可能有也可能没有应该被忽略的空格。此外,这些命令可能会或可能不会串联发送:a1234;bFGe4; 将是在同一消息中发送的两个命令a1234;bFGe4;。但是,有一个警告:如果命令足够长,则可能是不完整的。例如,我可能会收到两条消息c444a;X132124122412431,然后是1234124;,它们应该被翻译成两个单独的命令c444a;X1321241224124311234124;。这是由于硬件限制。

我设法使用一个 observable 来处理这种情况,该 observable 本身保持最后收到的字符串的结尾:

const ANY_MSG = /[a-zA-Z][^;]*;/g

const messages$ = new Observable<string>(sub => {
  let previousMsg = "";

  // monitor messages is the function w/ a callback that receives messages from the device
  monitorMessages((err, msg) => {
    if (err) {
      return sub.error(err);
    }

    const currMsg = previousMsg + msg
    const matches = currMsg.match(ANY_MSG)

    let lastIndex = 0
    if(matches) {
      for(const match of matches) {
        sub.next(match)
        lastIndex += match.length
      }
    }
    previousMsg = currMsg.slice(lastIndex)
  });
});

这个 observable 会按我的预期发出我的消息:如果它们被 monitorMessages 函数部分接收,它们用分号分隔并连接起来。

问题是,我觉得这个函数很难阅读和理解,如果它是通过 RXJS 的piped 函数编写的会更好。换句话说,我想做一些事情:

const messages$ = bindNodeCallback(monitorMessages).pipe(
  // ??????
);

但我不知道我需要在那里应用哪些运算符(或者即使我需要自己编写),因为发生以下情况:

  • 每当monitorMessages 发出消息时,将其分成两部分:消息的最后一个分号之前的部分和最后一个分号之后的部分
  • 从第一部分发出所有以分号分隔的命令,并存储第二部分
  • 每当monitorMessages 再次发出时,将存储的第二部分连接到新字符串并使其经历与以前相同的过程

这甚至可以完全通过 RXJS 操作符实现吗?也许我需要创建一个新的可观察对象,作为第二部分的“存储”?我只是觉得当前的方法(为 Observable 创建某种内部状态)非常奇怪,而且相对难以理解 atm。

【问题讨论】:

    标签: typescript rxjs operators observable


    【解决方案1】:

    我认为你的实现很好。但是,如果你真的想用操作符来实现它,你可以使用scan 来维护组合的可观察链中的一些状态,如下所示:

    const messages$ = bindNodeCallback(monitorMessages).pipe(
      scan((acc, received) => {
        const data = acc.remainder + received;
        const messages = data.match(ANY_MSG);
        if (messages) {
          const length = messsages.reduce((total, message) => total + message.length, 0);
          return { messages, remainder: data.slice(length) };
        }
        return { messages: [], remainder: data };
      }, { messages: [], remainder: "" }),
      mergeMap(({ messages }) => messages)
    );
    

    要发出您的消息数组,您可以使用mergeMap,返回数组 - 因为数组是ObservableInput

    【讨论】:

    • 这实际上正是我想要的。我假设第 3 行的 remainder 应该是 acc.remainder 对吧?我只是想将 Observable 创建(即“bindNodeCallback”)与其逻辑分开,这对我很有用。谢谢!
    • 是的,我已经确定了答案。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-11-25
    • 2017-05-07
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多