【问题标题】:How can I extract values from an Rx Observable?如何从 Rx Observable 中提取值?
【发布时间】:2016-05-05 22:17:36
【问题描述】:

我正在尝试将一些 ReactiveX 概念集成到现有项目中,认为这可能是一种很好的做法,也是一种使某些任务更简洁的方法。

我打开一个文件,从它的行创建一个 Observable,然后进行一些过滤,直到得到我想要的行。现在,我想使用 re.search() 从其中两行中提取一些信息以返回特定组。我一生都无法弄清楚如何从 Observable 中获取这些值(不将它们分配给全局变量)。

train = 'ChooChoo'

with open(some_file) as fd:
    line_stream = Observable.from_(fd.readlines())

a_stream = line_stream.skip_while(
        # Begin at dictionary
        lambda x: 'config = {' not in x
    ).skip_while(
        # Begin at train key
        lambda x: "'" + train.lower() + "'" not in x
    ).take_while(
        # End at closing brace of dict value
        lambda x: '}' not in x
    ).filter(
        # Filter sdk and clang lines only
        lambda x: "'sdk'" in x or "'clang'" in x
    ).subscribe(lambda x: match_some_regex(x))

代替该流末尾的.subscribe(),我尝试使用.to_list() 来获取一个列表,我可以“以正常方式”对其进行迭代,但它只返回一个类型的值:

<class 'rx.anonymousobservable.AnonymousObservable'>

我在这里做错了什么?

我见过的每个 Rx 示例都只打印结果。如果我希望它们在我可以同步使用的数据结构中怎么办?

【问题讨论】:

  • 你只能subscribe的回调中获取结果,这就是异步observables的全部意义——收集的结果不会同步可用,你不要不要让他们“out”这么多“in”!根据the docstring(强调我的),to_list “返回一个可观察序列,其中包含单个元素以及包含源序列所有元素的列表。”
  • 好的,那么...如何将该列表中的两行分配给本地命名空间中的变量 A 和 B?
  • 哪个本地命名空间?如果您指的是您还分配了a_stream 的那个,那么您没有!在控制返回到该范围时,无法保证异步过程已完成,这就是您观察它的原因
  • 啊,我明白了。一旦我观察到我需要的值(包含“sdk”和“clang”的两行),我是否需要以某种方式取消订阅或终止 Observable 以将这些行放入列表中?想象一个类似的场景,您正在异步查看击键并希望返回前 2 个击键的字母。从那里开始必须可以同步操作,对吧?
  • 是的,但是在回调中,不在您启动进程的范围内。目前尚不清楚您认为在这种情况下您从 observable 中获得了什么好处,只需使用 itertools

标签: python reactivex rx-py


【解决方案1】:

在短期内,我使用 itertools 实现了我想要的功能(正如 @jonrsharpe 所建议的那样)。这个问题一直萦绕在我的脑海里,所以我今天回过头来想通了。

这不是 Rx 的一个很好的例子,因为它只使用一个线程,但至少现在我知道如何在需要时突破“monad”。

#!/usr/bin/env python

from __future__ import print_function
from rx import *

def my_on_next(item):
    print(item, end="", flush=True)

def my_on_error(throwable):
    print(throwable)

def my_on_completed():
    print('Done')
    pass

def main():
    foo = []

    # Create an observable from a list of numbers
    a = Observable.from_([14, 9, 5, 2, 10, 13, 4])

    # Keep only the even numbers
    b = a.filter(lambda x: x % 2 == 0)

    # For every item, call a function that appends the item to a local list
    c = b.map(lambda x: foo.append(x))
    c.subscribe(lambda x: x, my_on_error, my_on_completed)

    # Use the list outside the monad!
    print(foo)

if __name__ == "__main__":
    main()

这个例子相当做作,所有的中间可观察对象都不是必需的,但它表明你可以轻松地做我最初描述的事情。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-07-07
    • 1970-01-01
    • 1970-01-01
    • 2016-03-15
    相关资源
    最近更新 更多