news 2026/9/25 5:47:12

RocketRide Python SDK 数据投喂实战:send / send_files / pipe 全解析

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
RocketRide Python SDK 数据投喂实战:send / send_files / pipe 全解析

【免费下载链接】rocketride-server

High-performance AI pipeline engine with a C++ core and 50+ Python-extensible nodes. Build, debug, and scale LLM workflows with 13+ model providers, 8+ vector databases, and agent orchestration, all from your IDE. Includes VS Code extension, TypeScript/Python SDKs, and Docker deployment.

项目地址:https://gitcode.com/gh_mirrors/ro/rocketride-server
点击查看免费下载

本指南以 RocketRide 官方 Python 文档 Sending Data 为核心,系统讲解如何将数据送入运行中的 AI 管道:一次性内存数据发送(send())、并发文件上传(send_files())以及分块流式传输(pipe()),并结合 client-python 源码 与测试用例,深入每一层实现细节。读完本文,你将掌握三种投喂方式的适用场景、参数语义、进度事件机制、异常处理策略,并能够编写健壮的生产级数据发送代码。

前置条件:数据只能发给特定来源的管道

RocketRide 中"发送数据"的目标对象,是已经在运行中的管道实例。管道启动后返回一个 token,所有数据与控制调用都以它为寻址依据(详见 Running Pipelines):

result = await client.use(filepath='pipeline.pipe') token = result['token']

关于数据入口,有一个必须先搞清楚的事实:send()/send_files()/pipe()只能投递给source为webhook或dropper的管道。如果你的管道来源是chat,请改用client.chat(),否则数据无法按预期进入处理流程。这一点在 DataMixin 源码 与官方文档中均有明确说明。

RocketRide Python SDK 是 async-first 的,基于asyncio与websockets构建,因此以下所有方法都需要在异步环境中调用(如asyncio.run(main()))。

一次性发送:send()

当你要发送的数据已经完整存在于内存中时,send()是最直接的路径。它的语义是:打开一条 pipe → 写入一次 → 关闭并返回管道处理结果。一次调用完成整条生命周期:

result = await client.send(token, 'Hello, pipeline!', objinfo={'name': 'greeting.txt'}, mimetype='text/plain')

参数语义与默认行为

参数含义默认值 / 说明
token管道 token(来自client.use())必填,寻址目标管道
data待发送内容,str或bytes必填;传入其他类型会抛出ValueError
objinfo数据的元信息字典,如{'name': 'greeting.txt'}可选,默认{}
mimetype数据的 MIME 类型可选,省略时以application/octet-stream发送
on_sse服务端事件回调可选,用于接收本次传输的 SSE 事件

两点容易被忽视的细节(原文档明确强调):

  • 没有自动检测。省略mimetype时负载一律按application/octet-stream发送,SDK 不会根据objinfo['name']的扩展名猜测类型。源码中DataPipe.__init__的self._mime_type = mime_type or 'application/octet-stream'正是这一行为的直接实现(mixins/data.py)。
  • 字符串自动转字节。send()内部会把str用 UTF-8 编码为bytes;只有str/bytes两种类型被接受(mixins/data.py)。

send()内部做了什么

从源码看,send()并不是一个独立的传输协议,而是pipe()的语法糖(mixins/data.py):

  1. 调用self.pipe(token, objinfo_with_size(...), mimetype, on_sse=on_sse)创建一个临时DataPipe;
  2. await pipe.open()打开连接(服务端分配pipe_id);
  3. await pipe.write(buffer)写入整块数据;
  4. await pipe.close()关闭并返回处理结果(PIPELINE_RESULT)。

其中_objinfo_with_size会向objinfo注入size字段,且保证不为 0(源码注释说明:值为 0 时会被解析过滤器当作"空"跳过,因此最小取 1)。如果中途任何一步抛错,send()会在finally语义中尽力关闭已打开的 pipe,避免泄漏连接。

返回结构:PIPELINE_RESULT

send()返回PIPELINE_RESULT(定义见 types/data.py),这是一个TypedDict,基础字段为:

  • name:本次处理结果的唯一标识(UUID 格式字符串);
  • path:文件路径上下文,直接数据发送时通常为空字符串;
  • objectId:被处理对象的唯一追踪 ID(UUID);
  • result_types:可选。只有当你指定了 MIME 类型、管道确实执行了内容提取/处理时才会出现,它是一张"字段名 → 数据类型"的映射表。

result_types是理解返回结构的钥匙:管道可以返回任意动态字段,字段名和类型都由它声明。例如:

result_types = {"text": "text", "answers": "answers", "metadata": "json"}

对应result["text"](List[str]文本内容)、result["answers"](List[str]AI 生成回复)、result["metadata"](dict JSON 元数据)。集成测试 RocketRideClient_test.py 验证了无 MIME 发送时返回结构只有基础字段、result_types为None的行为;而指定'text/plain'后返回中会携带动态处理字段(测试见同文件 L350 起)。

文件并发上传:send_files()

当数据来自磁盘上的多个文件、且你需要"每个文件单独的结果 + 实时进度事件"时,send_files()是首选。它通过asyncio.gather将整个文件列表一次性全部并发上传(服务端负责排队),返回一个与输入一一对应的UPLOAD_RESULT列表(mixins/data.py)。

三种条目格式

列表中的每个条目可以是:

  • 纯路径字符串:'doc1.md'—— SDK 自动取文件名作为objinfo['name'],并用mimetypes.guess_type()猜测 MIME 类型;
  • (path, objinfo) 二元组:('doc3.json', {'tag': 'export'})—— 显式指定元信息,MIME 类型仍自动猜测;
  • (path, objinfo, mimetype) 三元组:('doc3.json', {'tag': 'export'}, 'application/json')—— 三个要素全部显式指定,最可控。
files = ['doc1.md', 'doc2.md', ('doc3.json', {'tag': 'export'}, 'application/json')] upload_results = await client.send_files(files, token) for r in upload_results: if r['action'] == 'complete': print('OK', r['filepath']) else: print('Failed', r['filepath'], r.get('error'))

注意:与前文send()不同,send_files()拥有自动 MIME 检测。省略 MIME 时会调用 Python 标准库mimetypes.guess_type(filepath),猜不到时兜底application/octet-stream(mixins/data.py)。

两个必须知道的硬性约束

原文档特别提醒了两条:

  1. 需要 API key。send_files()要求客户端配置了 API key(self._apikey),否则直接抛出RuntimeError('API key is required for file uploads')(mixins/data.py)。
  2. 缺失文件抛ValueError。任何条目的路径经os.path.isfile()校验失败时,抛出ValueError(f'File not found: {filepath}')(mixins/data.py)。空列表、非法元组长度、非法条目类型同样会抛ValueError。

传输细节与进度事件

send_files()内部对每个文件执行"管道式线性流程"(mixins/data.py):

  1. 创建并打开 pipe(阻塞等待服务端分配资源);
  2. 以1 MB 固定分块(chunk_size = 1024 * 1024)逐块write;
  3. 关闭 pipe 拿到处理结果;
  4. 汇总bytes_sent、upload_time等统计。

每个文件在传输过程中都会通过事件系统发出apaevt_status_upload事件,事件 body 携带filepath、bytes_sent、file_size字段。完整的动作序列为:

动作(action)阶段body 关键字段
open文件上传开始,pipe 已分配filepath、file_size
write数据分块传输中filepath、bytes_sent、file_size
close传输完成,进入处理filepath、bytes_sent、file_size
complete上传+处理全部成功以上字段 +upload_time+result
error上传或处理失败以上字段 +error字符串

订阅方式:通过add_monitor({'token': token}, ['apaevt_status_upload'])建立监控订阅,事件会到达你的on_event回调(参见 Events 一节;remove_monitor用于反订阅,两者均按引用计数管理)。

返回结构:UPLOAD_RESULT

每个文件的返回条目是UPLOAD_RESULT(types/data.py):

  • action:'open' | 'write' | 'close' | 'complete' | 'error',最终状态机落点;
  • filepath:原始文件路径;
  • bytes_sent/file_size:已传输字节数 / 文件总大小;
  • upload_time:该文件上传耗时(秒);
  • result:仅在action == 'complete'时出现,为PIPELINE_RESULT(文件上传总是携带 MIME,通常带result_types,可提取text/answers等处理字段);
  • error:仅在action == 'error'时出现。

即便某个文件失败,send_files()也不会中断整体调用——asyncio.gather(..., return_exceptions=True)保证每个文件的任务独立收尾,失败的条目带着error字段进入结果列表。

分块流式传输:pipe()

当数据增量到达(实时流、逐行日志、生成器产出)或大到无法整体放进内存时,send()不再适用,pipe()是你的工具。一条流式上传的完整生命周期是open → write(一次或多次)→ close,close()返回处理结果。

pipe = await client.pipe(token, objinfo={'name': 'large.csv'}, mime_type='text/csv') await pipe.open() with open('large.csv', 'rb') as f: while True: chunk = f.read(64 * 1024) if not chunk: break await pipe.write(chunk) result = await pipe.close()

DataPipe核心契约

DataPipe(mixins/data.py)的行为约束如下:

  • write()只接受bytes。传入str会抛ValueError('Buffer must be bytes');非bytes一律拒绝。字符串请先.encode()(mixins/data.py)。
  • open()前不能write(),会抛RuntimeError('Pipe not opened')。
  • close()幂等:pipe 未打开或已关闭时直接返回{},不会重复关闭。
  • is_opened/pipe_id属性:pipe_id是open()成功后由服务端分配的整数 ID(mixins/data.py)。
  • 分块建议约 1 MB:服务端管道以 1 MB 左右的分块读取文件效果最佳(send_files的实现也印证了这一点),过大或过小的块都可能影响吞吐与内存占用。

异步上下文管理器(推荐写法)

DataPipe同时实现了__aenter__/__aexit__:进入自动open(),退出自动close(),异常路径也能保证收尾:

async with await client.pipe(token, mime_type='application/json') as pipe: await pipe.write(b'{"key": "value1"}') await pipe.write(b'{"key": "value2"}')

参数名差异:mimetype与mime_type

一个容易踩坑的细节:send()的参数叫mimetype,而pipe()的参数叫mime_type(下划线)。两者默认值一致——省略时均为application/octet-stream。从 pipe() 定义 可以看到签名是pipe(token, objinfo=None, mime_type=None, provider=None, on_sse=None)。

open()的瞬态错误重试

并发压力下(例如 CI 中一次性打开大量管道),数据监听器绑定尚未完成时open()可能遇到瞬态"Connect call failed"。SDK 对此做了有条件的单次重试(mixins/data.py):

  • 仅当错误消息精确匹配'Connect call failed'时才重试一次(_PIPE_OPEN_RETRY_ATTEMPTS = 2),退避0.25s;
  • "Connection refused"这类永久性失败不重试——它通常是remote节点配置错误,重试只会徒增延迟;
  • 重试耗尽后抛出PipeException,其中message原样保留服务端消息(适合直接展示给最终用户),hint附带开发者排障清单(管道未运行 / token 错误 / 管道 source 必须是 chat、webhook 或 dropper / MIME 与来源 lane 不匹配等)。

这些行为均由 test_pipe_open_retry.py 的 6 个用例逐条锁定:瞬态错误重试成功、重试耗尽后失败、非瞬态错误不重试、Connection refused不重试、畸形消息(非字符串 message)不崩溃、假值消息(0/False)不被吞掉。如果你的生产环境遇到open()偶发失败,可以优先对照这一清单排查。

通过 pipe 调用管道工具函数:DataPipe.tool()

DataPipe还提供了一个进阶能力:tool(tool, node_id='', input=None)可以通过当前 pipe 调用管道节点上的@tool_function(mixins/data.py)。它会复用本 pipe 已有的管道实例,避免从实例池重新借用带来的开销;node_id为空时向所有 tool-lane 节点广播,由第一个拥有该工具的节点处理。要求 pipe 已打开,否则抛RuntimeError。

SSE 事件订阅

pipe()与DataPipe本身都接受on_sse回调(async (type, data) => ...)。当open()成功后,SDK 会自动为这条 pipe 订阅['SSE']事件并把回调注册到对应pipe_id上(mixins/data.py);close()时则尽力反订阅并注销回调。完整的方法签名细节可查阅 API reference 中关于DataPipe的条目。

如何选择:一张决策表

你手头有什么使用哪个方法
内存中的字符串或 bytes,一次发完send()
磁盘上的多个文件,需要每个文件的结果 + 进度事件send_files()
分块/增量数据,或超大负载pipe()
目标是 chat 来源的管道chat()

补充两条工程经验:

  • 内存中的一次性数据优先send():它内部就是"pipe 三连",省去手动管理生命周期的负担,且自带异常清理;
  • 多文件场景优先send_files():并发 + 自动分块 + 事件进度三者齐备;若还需要更细粒度的进度 UI,订阅apaevt_status_upload事件即可拿到bytes_sent/file_size实时渲染。

错误处理要点

所有数据投递方法的异常语义可以总结为:

  • ValueError:send()收到非str/bytes数据;send_files()收到空列表、非法条目、或文件不存在('File not found: …');pipe.write()收到非bytes缓冲。
  • RuntimeError:send_files()缺少 API key;对未打开/已关闭的 pipe 执行非法操作。
  • PipeException:服务端拒绝 open / write / close,message保留服务端原文,hint携带排障清单,code归类失败类型。

完整异常体系可参见 错误处理文档。从源码结构看,这些异常统一继承自rocketride.core.exceptions(exceptions.py),捕获时按PipeException优先、其余按通用Exception兜底即可覆盖绝大多数生产场景。

小结

RocketRide Python SDK 的数据投喂能力围绕"管道 token + pipe 生命周期"这一核心模型展开:send()适合内存中的一次性负载,send_files()以 1 MB 分块并发上传磁盘文件并逐文件产出结果与进度事件,pipe()则以 open/write/close 三阶段覆盖增量与超大负载场景。三者的参数细节(mimetypevsmime_type、无自动检测、API key 依赖、字节强制、瞬态重试)都已在上文结合 DataMixin 源码 与测试用例逐一说明。实际使用时,先按决策表选对方法,再对照异常语义做好兜底,即可稳定地把数据送入 RocketRide 管道,交给 C++ 核心引擎处理。

【免费下载链接】rocketride-server

High-performance AI pipeline engine with a C++ core and 50+ Python-extensible nodes. Build, debug, and scale LLM workflows with 13+ model providers, 8+ vector databases, and agent orchestration, all from your IDE. Includes VS Code extension, TypeScript/Python SDKs, and Docker deployment.

项目地址:https://gitcode.com/gh_mirrors/ro/rocketride-server
点击查看免费下载
上一篇:终极Marlin 3D打印机固件教程:从入门到精通的完整指南
下一篇:【亲测免费】 ppInk 项目使用教程

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/25 5:46:17

用Winhance外部应用功能快速装机:WinGet一键安装常用软件指南

用Winhance外部应用功能快速装机:WinGet一键安装常用软件指南 【免费下载链接】Winhance-zh_CN A Chinese version of Winhance. C# application designed to optimize and customize your Windows experience. 项目地址: https://gitcode.com/gh_mirrors/wi/Winh…

作者头像 李华
网站建设 2026/9/25 5:45:32

自建CRM系统实战:从Excel迁移到Docker部署的完整指南

DeskcommCRM 是我们自己搭的一套客户管理系统,从决定自建到正式上线用了差不多三周时间。写这篇东西的原因很简单:我在配置系统、拉团队使用、处理数据迁移的过程中,看到太多人在“免费CRM”和“自建系统”之间来回纠结,也收到过不…

作者头像 李华
网站建设 2026/9/25 5:44:09

Atlas 300V部署YOLO:从硬件认知到环境搭建与模型转换实战

1. 从“atlas”到实际落地:先搞清楚它到底是个什么第一次看到“atlas”这个词,很多人会以为是个地图册,或者是某个希腊神话里的擎天巨神。但在AI算力、深度学习部署这个圈子里,atlas指的基本都是华为昇腾(Ascend&#…

作者头像 李华
网站建设 2026/9/25 5:41:30

AMD官网下载Vivado遇合规性失败?全流程排查与解决指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华