You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何创建输入输出解耦的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独立运行,并妥善处理流的启动、运行与终止逻辑。具体步骤如下:

  1. 单独启动线程运行输入处理Sink,使其与输出Source并行执行。
  2. 输出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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.26 01:24:14