这有几个名字。一个是“数据流”(如“反应式编程” - 如果你查一下,这有点夸张的流行语),另一个是“信号模拟”(如模拟电信号开关)。我不知道 Erlang 中有一个框架,因为直接实现它非常简单。
消息排序问题可以自行解决,具体取决于您要如何编写。 Erlang 保证两个进程之间的消息顺序,因此只要消息在定义明确的通道中传输,这个系统范围的承诺就可以为您工作。如果您需要一些比直线更有趣的信号路径,您可以强制同步通信;虽然所有 Erlang 消息都是异步的,但是你可以在任何你想要的地方对 receive 引入同步阻塞。
如果您希望“B 星座”将消息传递给 C,但只有在其信号处理完全通过 B 的路径后,您可以创建一个信号管理器,将消息发送到 B1,并阻塞直到它收到B3 的输出,它将完成的消息传递给 C 并检查其框以获取来自 A 的下一件事:
a_loop(B) ->
receive {in, Data} -> B ! Data end,
a_loop(B).
% Note the two receives here -- we are blocking for the end of processing based
% on the known Ref we send out and expect to receive back in a message match.
b_manager(B1, C) ->
Ref = make_ref(),
receive Data -> B1 ! {Ref, Data} end,
receive {Ref, Result} -> C ! Result end,
b_manager(B1, C).
b_1(B2) ->
receive
{Ref, Data} ->
Mod1 = do_processing(Data),
B2 ! {Ref, Mod1}
end,
b_1(B2).
% Here you have as many "b_#" processes as you need...
b_2(B) ->
receive
{Ref, Data} ->
Result = do_other_processing(Data),
B ! {Ref, Result}
end,
b_2(B).
c_loop() ->
receive Result -> stuff(Result) end,
c_loop().
显然我大大地简化了一些事情——因为这显然不包括任何监督的概念——我什至没有说明你希望如何将它们链接在一起(以及这个小检查活性,您将 需要 将它们生成链接,因此如果有任何东西死亡,它们都会死亡 - 这可能正是您想要的 B 子集,因此您可以将其视为一个单元)。此外,您可能最终需要在某个地方(例如 A 处/之前或 B 处)安装油门。但基本上来说,这是一种让 B 阻塞直到它的处理段完成的方式传递消息的方式。
还有其他方法,例如 gen_event,但我发现它们不如编写处理管道的实际模拟灵活。至于如何实现这一点——我会将它组合成 OTP 管理器和 gen_fsm,因为这两个组件代表了与信号处理组件几乎完美的并行,而您的系统似乎旨在模仿。
为了发现您在 gen_fsms 中需要哪些状态以及如何将它们聚集在一起,我可能会在纯 Erlang 中以非常简单的方式制作几个小时的原型,以确保我确实了解问题,然后编写我适当的 OTP 主管和 gen_fsms。这可以确保我不会投入到一些 gen_foo 行为的殿堂中,而不是投入到实际解决我的问题上(无论如何,在它正确之前你必须至少写两次...... .).
希望这至少能给你一个开始解决问题的地方。无论如何,这在 Erlang 中是一件非常自然的事情——并且与语言和问题的工作方式非常接近,因此处理起来应该很有趣。