ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

RocketRide Python SDK 数据投喂实战:send / send_files / pipe 全解析 【免费下载链接】rocketride-serverHigh-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 Pipelinesresult await client.use(filepathpipeline.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}, mimetypetext/plain)参数语义与默认行为参数含义默认值 / 说明token管道 token来自client.use()必填寻址目标管道data待发送内容str或bytes必填传入其他类型会抛出ValueErrorobjinfo数据的元信息字典如{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_sseon_sse)创建一个临时DataPipeawait pipe.open()打开连接服务端分配pipe_idawait pipe.write(buffer)写入整块数据await pipe.close()关闭并返回处理结果PIPELINE_RESULT。其中_objinfo_with_size会向objinfo注入size字段且保证不为 0源码注释说明值为 0 时会被解析过滤器当作空跳过因此最小取 1。如果中途任何一步抛错send()会在finally语义中尽力关闭已打开的 pipe避免泄漏连接。返回结构PIPELINE_RESULTsend()返回PIPELINE_RESULT定义见 types/data.py这是一个TypedDict基础字段为name本次处理结果的唯一标识UUID 格式字符串path文件路径上下文直接数据发送时通常为空字符串objectId被处理对象的唯一追踪 IDUUIDresult_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-streammixins/data.py。两个必须知道的硬性约束原文档特别提醒了两条需要 API key。send_files()要求客户端配置了 API keyself._apikey否则直接抛出RuntimeError(API key is required for file uploads)mixins/data.py。缺失文件抛ValueError。任何条目的路径经os.path.isfile()校验失败时抛出ValueError(fFile 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_sizewrite数据分块传输中filepath、bytes_sent、file_sizeclose传输完成进入处理filepath、bytes_sent、file_sizecomplete上传处理全部成功以上字段 upload_timeresulterror上传或处理失败以上字段 error字符串订阅方式通过add_monitor({token: token}, [apaevt_status_upload])建立监控订阅事件会到达你的on_event回调参见 Events 一节remove_monitor用于反订阅两者均按引用计数管理。返回结构UPLOAD_RESULT每个文件的返回条目是UPLOAD_RESULTtypes/data.pyactionopen | 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_exceptionsTrue)保证每个文件的任务独立收尾失败的条目带着error字段进入结果列表。分块流式传输pipe()当数据增量到达实时流、逐行日志、生成器产出或大到无法整体放进内存时send()不再适用pipe()是你的工具。一条流式上传的完整生命周期是open → write一次或多次→ closeclose()返回处理结果。pipe await client.pipe(token, objinfo{name: large.csv}, mime_typetext/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核心契约DataPipemixins/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()成功后由服务端分配的整数 IDmixins/data.py。分块建议约 1 MB服务端管道以 1 MB 左右的分块读取文件效果最佳send_files的实现也印证了这一点过大或过小的块都可能影响吞吐与内存占用。异步上下文管理器推荐写法DataPipe同时实现了__aenter__/__aexit__进入自动open()退出自动close()异常路径也能保证收尾async with await client.pipe(token, mime_typeapplication/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, objinfoNone, mime_typeNone, providerNone, on_sseNone)。open()的瞬态错误重试并发压力下例如 CI 中一次性打开大量管道数据监听器绑定尚未完成时open()可能遇到瞬态Connect call failed。SDK 对此做了有条件的单次重试mixins/data.py仅当错误消息精确匹配Connect call failed时才重试一次_PIPE_OPEN_RETRY_ATTEMPTS 2退避0.25sConnection 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, inputNone)可以通过当前 pipe 调用管道节点上的tool_functionmixins/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.pyclose()时则尽力反订阅并注销回调。完整的方法签名细节可查阅 API reference 中关于DataPipe的条目。如何选择一张决策表你手头有什么使用哪个方法内存中的字符串或 bytes一次发完send()磁盘上的多个文件需要每个文件的结果 进度事件send_files()分块/增量数据或超大负载pipe()目标是 chat 来源的管道chat()补充两条工程经验内存中的一次性数据优先send()它内部就是pipe 三连省去手动管理生命周期的负担且自带异常清理多文件场景优先send_files()并发 自动分块 事件进度三者齐备若还需要更细粒度的进度 UI订阅apaevt_status_upload事件即可拿到bytes_sent/file_size实时渲染。错误处理要点所有数据投递方法的异常语义可以总结为ValueErrorsend()收到非str/bytes数据send_files()收到空列表、非法条目、或文件不存在File not found: …pipe.write()收到非bytes缓冲。RuntimeErrorsend_files()缺少 API key对未打开/已关闭的 pipe 执行非法操作。PipeException服务端拒绝 open / write / closemessage保留服务端原文hint携带排障清单code归类失败类型。完整异常体系可参见 错误处理文档。从源码结构看这些异常统一继承自rocketride.core.exceptionsexceptions.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-serverHigh-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 SDK 实战指南用 rocketride 客户端构建、运行与部署 AI PipelineRocketRide Python SDK 实战指南用 rocketride 客户端构建、运行与部署 AI Pipeline 导读 本文以 RocketRidPyJNIus核心功能解析autoclass如何实现Python与Java的桥梁PyJNIus核心功能解析autoclass如何实现Python与Java的桥梁 PyJNIus是一个强大的Python库它通过autoclass功能实现了跨平台移动开发MOOTDX量化投资指南Python通达信数据接口实战解析MOOTDX量化投资指南Python通达信数据接口实战解析 还在为量化投资数据获取而烦恼吗面对复杂的API接口和繁琐的数据处理流程很多量化爱好者常常在起步金融科技数据分析上一篇终极Marlin 3D打印机固件教程从入门到精通的完整指南下一篇【亲测免费】 ppInk 项目使用教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表