如何创建输入输出解耦的ConduitT?类比Akka Streams fromSinkAndSource
问题:ConduitT中实现类似Akka Streams fromSinkAndSource的解耦流组合
在Akka Streams中,可通过fromSinkAndSource创建Flow,将输入输出解耦的Sink与Source组合。现需为Servant WebSocket提供一个ConduitT,要求数据摄入与输出生成过程完全解耦,由底层支持并发读写的Monad在多线程环境驱动。请问是否存在对应的ConduitT组合子或用法实现该需求?
场景补充
我正在实现基于ConduitT的游戏服务器(以FizzBuzz为例):游戏引擎计算结果会引入延迟,用户输入为整数指令,服务器仅在满足条件时输出字符串。以下是尝试的代码(无法正常运行):
#!/usr/bin/env stack -- stack --resolver lts-18.28 script --package conduit --package stm {-# LANGUAGE NumericUnderscores #-} -- 实现带延迟效果的FizzBuzz import Conduit import Control.Concurrent (forkIO, threadDelay) import Control.Concurrent.STM import Control.Monad import Data.Conduit.List as CL import Data.Void producer :: TQueue String -> ConduitT String String IO () producer q = repeatMC readFromQueue where readFromQueue = atomically $ readTQueue q consumer :: TQueue String -> ConduitT Int String IO () consumer q = awaitForever $ \elem -> lift $ fizz elem where fizz :: Int -> IO () fizz n | n `mod` 15 == 0 = delayedWrite 15 "FizzBuzz" | n `mod` 3 == 0 = delayedWrite 3 "Fizz" | n `mod` 5 == 0 = delayedWrite 5 "Buzz" | otherwise = return () delayedWrite :: Int -> String -> IO () delayedWrite i s = void . forkIO $ do threadDelay (i * 1_000_000) atomically $ writeTQueue q s flowSinkAndSource :: TQueue String -> ConduitT Int String IO () flowSinkAndSource q = consumer q .| producer q main :: IO () main = do let source = CL.sourceList [1 .. 30] sink = CL.mapM_ print in do q <- newTQueueIO runConduit $ source .| flowSinkAndSource q .| sink
原以为通过TQueue连接consumer与producer即可实现解耦,但运行时出现线程无限阻塞在STM事务的错误:
❯ ./fromSinkAndSource-in-conduit.hs fromSinkAndSource-in-conduit.hs: thread blocked indefinitely in an STM transaction
尝试修改写法后问题仍存在:
flowSinkAndSource :: TQueue String -> ConduitT Int String IO () flowSinkAndSource q = consumer q .| CL.sinkNull .| producer q
解决方案
问题根源
原代码存在两个核心问题:
- Conduit连接逻辑错误:
consumer q .| producer q试图将consumer的输出流传给producer的输入流,但consumer并未通过yield产生任何输出,仅异步写入TQueue。Conduit的线性执行逻辑会先跑完consumer,再执行producer,此时producer会因持续从空TQueue读取而阻塞。 - 未处理流终止逻辑:上游source([1..30])处理完成后,consumer结束,但producer仍持续等待TQueue内容,导致STM事务无限阻塞。
正确实现方式
要实现类似Akka Streams的解耦效果,需将输入处理Sink与输出生成Source独立运行,并妥善处理流的启动、运行与终止逻辑。具体步骤如下:
- 单独启动线程运行输入处理Sink,使其与输出Source并行执行。
- 输出Source从TQueue读取内容,同时监听输入流的终止信号,确保在输入结束后处理完剩余队列内容再终止。
修改后的代码:
#!/usr/bin/env stack -- stack --resolver lts-18.28 script --package conduit --package stm --package async {-# LANGUAGE NumericUnderscores #-} import Conduit import Control.Concurrent (threadDelay) import Control.Concurrent.STM import Control.Concurrent.Async (async, waitAnyCancel) import Control.Monad import Data.Conduit.List as CL import Data.Void -- 输入处理Sink:接收整数,异步计算后写入TQueue consumerSink :: TQueue String -> ConduitT Int Void IO () consumerSink q = awaitForever $ \n -> lift $ do let fizzStr | n `mod` 15 == 0 = Just "FizzBuzz" | n `mod` 3 == 0 = Just "Fizz" | n `mod` 5 == 0 = Just "Buzz" | otherwise = Nothing forM_ fizzStr $ \s -> void $ async $ do let delay = if n `mod`15 ==0 then 15 else if n`mod`3==0 then3 else5 threadDelay (delay * 1_000_000) atomically $ writeTQueue q s -- 输出生成Source:从TQueue读取内容,直到终止信号触发且队列空 producerSource :: TQueue String -> TVar Bool -> ConduitT () String IO () producerSource q isDone = loop where loop = do str <- lift $ atomically $ do done <- readTVar isDone -- 若输入已完成,优先读取队列剩余内容;否则等待新内容或终止信号 orElse (readTQueue q) $ do check done -- 队列空且已终止,返回空字符串触发退出 return "" if str /= "" then yield str >> loop else return () -- 组合解耦的Sink与Source,返回ConduitT Int String IO () flowSinkAndSource :: ConduitT Int String IO () flowSinkAndSource = do q <- lift newTQueueIO isDone <- lift $ newTVarIO False -- 启动异步线程处理输入 consumerAsync <- lift $ async $ runConduit $ awaitForever yield .| consumerSink q -- 监听输入流完成,标记终止信号 lift $ async $ do waitAnyCancel [consumerAsync] atomically $ writeTVar isDone True -- 输出流生成 producerSource q isDone -- 等待输入处理线程完成 lift $ waitAnyCancel [consumerAsync] main :: IO () main = do let source = CL.sourceList [1 .. 30] sink = CL.mapM_ print runConduit $ source .| flowSinkAndSource .| sink
关键改进点
- 并行解耦:通过
async单独启动输入处理线程,让输入消费与输出生成并行运行,摆脱Conduit线性连接的限制。 - 优雅终止:用
TVar Bool标记输入流完成状态,确保producer在输入结束后读完队列剩余内容再终止,避免阻塞。 - 类型匹配:consumerSink为
ConduitT Int Void IO ()(仅消费输入),producerSource为ConduitT () String IO ()(仅生成输出),严格实现输入输出解耦。
内容的提问来源于stack exchange,提问作者tonicebrian
相关产品推荐
相关产品推荐

