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

Conduit组合子示例解析:上游完整消息触发下游输出实现

Conduit组合子解析与服务器实现思路

让我来拆解这个Conduit组合子的实现逻辑,再给你讲讲对应的服务器代码该怎么写~

一、message组合子的实现逻辑

首先先看完整的组合子代码:

import qualified Data.ByteString as BS
import Data.Conduit
import Data.Conduit.Combinators
import Data.Conduit.Network

message :: Monad m => ConduitT BS.ByteString BS.ByteString m ()
message = loop
  where
    loop = await >>= maybe (return ()) go
    go x = if (BS.isSuffixOf "|" x)
           then yield (BS.init x) >> loop
           else leftover x

这个组合子的核心作用是从上游字节流中提取以|结尾的完整消息,去掉结尾的|后传递给下游,具体逻辑拆解如下:

  • 递归循环处理流:message直接绑定到loop函数,用递归的方式持续处理上游传来的数据块。
  • 等待上游数据:loop里的await会尝试从上游获取一个BS.ByteString数据块:
    • 如果上游已经没有数据(maybe的第一个分支),就直接退出循环,结束处理。
    • 如果拿到了数据块,就交给go函数处理。
  • 判断消息完整性:
    • go x首先检查当前数据块x是否以|结尾(BS.isSuffixOf "|" x):
      • 如果是完整消息:用BS.init x去掉结尾的|,然后通过yield把处理后的内容发送给下游,接着回到loop等待下一个消息。
      • 如果是不完整消息:用leftover x把这个数据块放回上游的缓冲区,等下一次await时会重新获取它,这样后续传来的新数据块就能和它拼接起来,继续判断是否构成完整消息。

二、服务器代码实现思路

基于这个message组合子,我们可以用Data.Conduit.Network快速实现一个TCP服务器,核心思路是:监听指定端口,对每个客户端连接,将客户端发来的字节流通过message组合子处理,再将结果返回给客户端(或做其他业务处理)。

完整服务器示例代码

import qualified Data.ByteString as BS
import Data.Conduit
import Data.Conduit.Combinators
import Data.Conduit.Network

message :: Monad m => ConduitT BS.ByteString BS.ByteString m ()
message = loop
  where
    loop = await >>= maybe (return ()) go
    go x = if (BS.isSuffixOf "|" x)
           then yield (BS.init x) >> loop
           else leftover x

main :: IO ()
main = runTCPServer (serverSettings 8080 HostAny) $ \appData -> do
  putStrLn "New client connected!"
  -- 构建数据处理管道:客户端输入 → 提取完整消息 → 返回给客户端
  appSource appData .| message .| appSink appData
  putStrLn "Client disconnected!"

代码解释

  • 启动TCP服务器:runTCPServer是Data.Conduit.Network提供的便捷函数,serverSettings 8080 HostAny表示监听8080端口,接受来自任意IP的连接。
  • 处理客户端连接:每个新连接会触发后面的匿名函数,appData包含了该连接的输入源(appSource,客户端发来的字节流)和输出sink(appSink,发送数据给客户端的通道)。
  • 数据管道:appSource appData .| message .| appSink appData是Conduit的管道语法,意思是:
    1. 从客户端读取字节流;
    2. 经过message组合子提取完整消息;
    3. 将处理后的完整消息发回给客户端。
  • 扩展方向:你可以在管道中添加更多组合子,比如在message之后加上mapM_ (liftIO . BS.putStrLn),这样就能在服务器控制台打印收到的完整消息;或者把消息写入文件、转发到其他服务等。

内容的提问来源于stack exchange,提问作者Kuznero

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.21 04:21:51