【问题标题】:frameworks for representing data processing as a pipeline将数据处理表示为管道的框架
【发布时间】: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

我发现这不是一个非常可组合的抽象。一些问题是:

  1. tweaker只能用在handler的左侧,即我不能使用tweaker的上述实现来形成这个管道:

    阅读器 |处理程序 |调整器 |作家

  2. 我想利用管道的关联属性,让这个管道:

    阅读器 |处理程序 |作家

可以表示为:

reader | p

其中p 是管道handler | writer。在这个 OO 实现中,我必须部分实例化 handler 对象

  1. 有点重述 (1),对象必须知道它们是“推”还是“拉”数据。

我正在寻找一个框架(不一定是 OO)来创建解决这些问题的数据处理管道。

我用Haskell 和functional programming 标记了它,因为我觉得函数式编程概念在这里可能有用。

作为一个目标,能够创建这样的管道会很好:

                     handler1
                   /          \
reader | partition              writer
                   \          /
                     handler2

从某种角度来看,Unix shell 管道通过以下实现决策解决了很多此类问题:

  1. 管道组件在不同进程中异步运行

  2. 管道对象在“推杆”和“拉杆”之间调解传递数据;即,它们会阻止写入数据过快的写入器和尝试读取过快的读取器。

  3. 您使用特殊连接器&lt; 和&gt; 将无源组件(即文件)连接到管道

我对在代理之间不使用线程或消息传递的方法特别感兴趣。也许这是最好的方法,但我想尽可能避免使用线程。

谢谢!

【问题讨论】:

  • 也许您想产生几个线程,一个用于每个读取器、调整器、处理程序和写入器,并通过Chans 进行通信?不过,我不是 100% 确定我理解顶级问题是什么……
  • 到目前为止,最后一张图看起来像reader &gt;&gt;&gt; partition &gt;&gt;&gt; handler1 *** handler2 &gt;&gt;&gt; writer,但可能会有一些要求使它变得更复杂。
  • 如果有帮助,我对partition 的想法是它会根据选择函数将输入数据发送到一个输出或另一个输出。
  • @user5402,可以做到这一点的箭头是 ArrowChoice 的实例,partition 运算符的 双重(仅使用 arr 就可以轻松进行分区,但是如果你不能重新加入就没有任何好处)是(|||)。

标签: haskell functional-programming pipe pipeline data-processing


【解决方案1】:

是的,arrows 几乎肯定是你的男人。

我怀疑您对 Haskell 相当陌生,仅基于您在问题中所说的内容。箭头可能看起来相当抽象,特别是如果您正在寻找的是“框架”。我知道我花了一段时间才真正了解箭头发生了什么。

因此,您可能会看着该页面并说“是的,这看起来像我想要的”,然后发现自己对如何开始使用箭头来解决问题感到迷茫。所以这里有一点指导,让你知道你在看什么。

Arrows 无法解决您的问题。相反,它们为您提供了一种可以用来表达问题的语言。你可能会发现一些预定义的箭头可以完成这项工作——也许是一些 kleisli 箭头——但最终你会想要实现一个箭头(预定义的只是让你容易实现它们的方法),它表达了“数据处理器”的意思。作为一个几乎微不足道的例子,假设你想通过简单的函数来实现你的数据处理器。你会写:

newtype Proc a b = Proc { unProc :: a -> b }

-- I believe Arrow has recently become a subclass of Category, so assuming that.

instance Category Proc where
    id = Proc (\x -> x)
    Proc f . Proc g = Proc (\x -> f (g x))

instance Arrow Proc where
    arr f = Proc f
    first (Proc f) = Proc (\(x,y) -> (f x, y))

这为您提供了使用各种箭头组合器(***)、(&amp;&amp;&amp;)、(&gt;&gt;&gt;) 等的机制,以及箭头符号,如果您正在做复杂的事情,这是相当不错的。因此,正如 Daniel Fischer 在评论中指出的那样,您在问题中描述的管道可以组成为:

reader >>> partition >>> (handler1 *** handler2) >>> writer

但很酷的是,处理器的含义由您决定。可以使用不同的处理器类型以类似的方式实现您提到的关于每个处理器派生线程的内容:

newtype Proc' a b = Proc (Source a -> Sink b -> IO ())

然后适当地实现组合器。

这就是您正在查看的内容:用于讨论组合进程的词汇表,其中有一些代码可以重用,但主要有助于在您实现这些组合子以定义有用的处理器时指导您的思考在您的域中。

我的第一个重要的 Haskell 项目是实现 arrow for quantum entanglement;那个项目让我真正开始理解 Haskell 的思维方式,这是我编程生涯的一个重要转折点。也许你的这个项目也会为你做同样的事情? :-)

【讨论】:

    【解决方案2】:

    感谢惰性求值,我们可以用 Haskell 中的普通函数组合来表达管道。这是一个计算文件中行的最大长度的示例:

    main = interact (show . maximum . map length . lines)
    

    这里的一切都是一个普通的函数,例如

    lines :: String -> [String]
    

    但是由于惰性求值,这些函数只能增量处理输入,并且只处理需要的量,就像 UNIX 管道一样。

    【讨论】:

      【解决方案3】:

      Haskell 的enumerator package 是一个很好的框架。它定义了三种类型的对象:

      1. 以块的形式生成数据的枚举器。
      2. 消耗数据块并在消耗足够多后返回值的迭代。
      3. 位于管道中间的枚举数。它们消耗块并产生块,可能有副作用。

      这三种类型的对象组成了一个流处理管道,你甚至可以在一个管道中拥有多个枚举器和迭代器(当一个完成时,下一个代替它)。从头开始编写这些对象之一可能很复杂,但是有很多组合器可用于将常规函数转换为数据流处理器。例如,这条管道从标准输入读取所有字符,使用函数toUpper将它们转换为大写,然后将它们写入标准输出:

      ET.enumHandle stdin $$ ET.map toUpper =$ ET.iterHandle stdout
      

      模块Data.Enumerator.Text 已导入为ET。

      【讨论】:

      【解决方案4】:

      Yesod 框架使用 conduit 包形式的 Haskell 管道库。

      【讨论】:

        猜你喜欢
        • 2021-07-15
        • 2019-07-15
        • 2021-12-16
        • 2011-06-14
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多