【发布时间】:2011-12-29 20:53:30
【问题描述】:
大多数数据处理都可以设想为组件的管道,一个组件的输出馈入另一个组件的输入。一个典型的处理管道是:
reader | handler | writer
作为开始这个讨论的辅助,让我们考虑这个管道的面向对象的实现,其中每个段都是一个对象。 handler 对象包含对 reader 和 writer 对象的引用,并有一个 run 方法,如下所示:
define handler.run:
while (reader.has_next) {
data = reader.next
output = ...some function of data...
writer.put(output)
}
从示意图上看,依赖关系是:
reader <- handler -> writer
现在假设我想在阅读器和处理程序之间插入一个新的管道段:
reader | tweaker | handler | writer
同样,在这个 OO 实现中,tweaker 将是 reader 对象的包装器,tweaker 方法可能看起来像(在一些伪命令式代码中):
define tweaker.has_next:
return reader.has_next
define tweaker.next:
value = reader.next
result = ...some function of value...
return result
我发现这不是一个非常可组合的抽象。一些问题是:
-
tweaker只能用在handler的左侧,即我不能使用tweaker的上述实现来形成这个管道:阅读器 |处理程序 |调整器 |作家
-
我想利用管道的关联属性,让这个管道:
阅读器 |处理程序 |作家
可以表示为:
reader | p
其中p 是管道handler | writer。在这个 OO 实现中,我必须部分实例化 handler 对象
- 有点重述 (1),对象必须知道它们是“推”还是“拉”数据。
我正在寻找一个框架(不一定是 OO)来创建解决这些问题的数据处理管道。
我用Haskell 和functional programming 标记了它,因为我觉得函数式编程概念在这里可能有用。
作为一个目标,能够创建这样的管道会很好:
handler1
/ \
reader | partition writer
\ /
handler2
从某种角度来看,Unix shell 管道通过以下实现决策解决了很多此类问题:
管道组件在不同进程中异步运行
管道对象在“推杆”和“拉杆”之间调解传递数据;即,它们会阻止写入数据过快的写入器和尝试读取过快的读取器。
您使用特殊连接器
<和>将无源组件(即文件)连接到管道
我对在代理之间不使用线程或消息传递的方法特别感兴趣。也许这是最好的方法,但我想尽可能避免使用线程。
谢谢!
【问题讨论】:
-
也许您想产生几个线程,一个用于每个读取器、调整器、处理程序和写入器,并通过
Chans 进行通信?不过,我不是 100% 确定我理解顶级问题是什么…… -
到目前为止,最后一张图看起来像
reader >>> partition >>> handler1 *** handler2 >>> writer,但可能会有一些要求使它变得更复杂。 -
如果有帮助,我对
partition的想法是它会根据选择函数将输入数据发送到一个输出或另一个输出。 -
@user5402,可以做到这一点的箭头是
ArrowChoice的实例,partition运算符的 双重(仅使用arr就可以轻松进行分区,但是如果你不能重新加入就没有任何好处)是(|||)。
标签: haskell functional-programming pipe pipeline data-processing