【免费下载链接】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.
本指南以 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):
- 调用
self.pipe(token, objinfo_with_size(...), mimetype, on_sse=on_sse)创建一个临时DataPipe; await pipe.open()打开连接(服务端分配pipe_id);await pipe.write(buffer)写入整块数据;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)。
两个必须知道的硬性约束
原文档特别提醒了两条:
- 需要 API key。
send_files()要求客户端配置了 API key(self._apikey),否则直接抛出RuntimeError('API key is required for file uploads')(mixins/data.py)。 - 缺失文件抛
ValueError。任何条目的路径经os.path.isfile()校验失败时,抛出ValueError(f'File not found: {filepath}')(mixins/data.py)。空列表、非法元组长度、非法条目类型同样会抛ValueError。
传输细节与进度事件
send_files()内部对每个文件执行"管道式线性流程"(mixins/data.py):
- 创建并打开 pipe(阻塞等待服务端分配资源);
- 以1 MB 固定分块(
chunk_size = 1024 * 1024)逐块write; - 关闭 pipe 拿到处理结果;
- 汇总
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.
相关推荐
RocketRide Python SDK 实战指南:用 `rocketride` 客户端构建、运行与部署 AI Pipeline
RocketRide Python SDK 实战指南:用 rocketride 客户端构建、运行与部署 AI Pipeline 导读 本文以 RocketRid
PyJNIus核心功能解析:autoclass如何实现Python与Java的桥梁
PyJNIus核心功能解析:autoclass如何实现Python与Java的桥梁 PyJNIus是一个强大的Python库,它通过autoclass功能实现了
跨平台移动开发MOOTDX量化投资指南:Python通达信数据接口实战解析
MOOTDX量化投资指南:Python通达信数据接口实战解析 还在为量化投资数据获取而烦恼吗?面对复杂的API接口和繁琐的数据处理流程,很多量化爱好者常常在起步
金融科技数据分析
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考