DORA Python API 开发实战指南:Node、Operator 与 DataflowBuilder 全解析
【免费下载链接】doraDORA (Dataflow-Oriented Robotic Architecture) is middleware designed to streamline and simplify the creation of AI-based robotic applications. It offers low latency, composable, and distributed dataflow capabilities. Applications are modeled as directed graphs, also referred to as pipelines.项目地址: https://gitcode.com/GitHub_Trending/do/dora
DORA(Dataflow-Oriented Robotic Architecture)是一套面向 AI 机器人应用的低延迟、可组合、分布式数据流中间件,应用被建模为有向图(pipeline)。本文基于仓库 docs/api-python.md 编写,系统讲解 DORA 的 Python 编程接口:如何用dora.Node编写自定义节点、如何在 dora runtime 进程内运行Operator、如何用DataflowBuilder在 Python 中程序化生成数据流 YAML,并延伸到 CUDA 张量零拷贝传输、服务/动作/流式通信模式与结构化日志。读完本文,你可以从零编写一个可运行的 Python 节点,并将其接入 dora daemon 管理的数据流图。
安装与运行环境要求
DORA 的 Python 客户端以dora-rs包发布,安装命令:
pip install dora-rs解释器兼容性:CPython 3.11 及以上。官方 wheel 以abi3-py311构建,单个 wheel 即可在所有后续 CPython 版本上运行,无需等待 dora 发版——3.11 是 dora 1.x 的下限(floor),并非锁定版本。不支持的构建:free-threaded 构建(python3.13t、python3.14t)不在 abi3 覆盖范围内,也没有为其发布 wheel。完整保证请参阅 Python 版本策略。
开发调试本仓库内的 Python 绑定(位于 apis/python/node)时,README 给出了本地构建方式:
uv venv --seed -p 3.11 uv pip install -e .类型提示需额外生成 stub 文件:
python generate_stubs.py dora dora/__init__.pyi maturin developNode API:自定义节点的核心接口
Node是编写自定义节点的首要接口,它连接到一个正在运行的数据流(dataflow),接收输入事件并发送输出:
from dora import Node从源码结构看,Node类由 Rust 通过 pyo3 导出(见 apis/python/node/src/lib.rs),内部持有DoraNode与一个Events事件流,并在构造时自动安装日志桥接(见下文“日志”小节)。
__init__(node_id=None)
创建一个新节点并连接正在运行的数据流:
# 标准用法:节点 ID 从 daemon 设置的环境变量读取 node = Node() # 动态用法:通过显式节点 ID 连接到正在运行的数据流 node = Node(node_id="my-dynamic-node")参数:
node_id(str,可选)——动态节点的显式 ID。省略时节点从 dora daemon 设置的环境变量读取身份。
异常:节点无法连接数据流时抛出RuntimeError。
从类型定义(dora/init.pyi)可以看到构造签名还支持daemon_port:Node("my-node", daemon_port=6789)可用于连接自定义端口上的 daemon。Rust 实现中(lib.rs),显式传入node_id走DoraNode::init_flexible(node_id)(带端口时使用DoraNode::builder()),省略时走DoraNode::init_from_env()。
接收事件:next/drain/try_recv/recv_async/is_empty
next(timeout=None)—— 从事件流取下一个事件,阻塞直到有事件可用或超时:
event = node.next() # 无限期阻塞 event = node.next(timeout=2.0) # 最多阻塞 2 秒timeout(float,可选)——最大等待秒数。- 返回:事件字典;当所有发送方都已关闭或超时到期时返回
None。
实现上(lib.rs),timeout会被转换为Duration;Rust 侧对 NaN、负数或无穷大的超时输入返回干净的ValueError而非 panic(见timeout_to_duration,lib.rs)。
drain()—— 不阻塞地取出所有缓冲事件:
events = node.drain() for event in events: print(event["type"])返回list[dict];无缓冲事件时返回空列表。
try_recv()—— 非阻塞接收,返回下一个缓冲事件(若有):
event = node.try_recv() if event is not None: print(event["type"])返回dict | None。
recv_async(timeout=None)—— 异步接收,配合asyncio使用:
event = await node.recv_async() event = await node.recv_async(timeout=5.0)- 超时到达时返回错误;所有发送方关闭时返回
None。 - 注意:此方法为实验特性,pyo3 的 async(Rust-Python FFI)集成仍在开发中。
is_empty()—— 检查事件流中是否还有缓冲事件:
if not node.is_empty(): event = node.try_recv()返回bool。
迭代支持
Node实现了__iter__与__next__(lib.rs),可以直接用for循环消费事件流:
for event in node: match event["type"]: case "INPUT": process(event["value"]) case "STOP": break迭代器每次迭代调用无超时的next();事件流关闭时产出None从而终止循环。
发送输出:send_output(output_id, data, metadata=None)
在输出通道上发送数据:
import pyarrow as pa # 发送原始字节 node.send_output("status", b"OK") # 发送 Apache Arrow 数组(可零拷贝) node.send_output("values", pa.array([1, 2, 3])) # 携带元数据发送 node.send_output("image", pa.array(pixels), {"camera_id": "front"})参数:
output_id(str)——数据流 YAML 中声明的输出名。data(bytes | pyarrow.Array)——负载。简单数据用bytes,需要零拷贝共享内存传输时用pyarrow.Array。metadata(dict,可选)——附加到消息的键值对。支持的值类型:bool、int、float、str、list[int]、list[float]、list[str]、datetime.datetime。
异常:data既不是bytes也不是pyarrow.Array时抛出RuntimeError。
Rust 侧的分发逻辑(send_payload,lib.rs)验证:PyBytes走send_output_bytes原始字节路径,pyarrow 数组走 Arrow 数组路径,其余类型直接报错。
零拷贝补充:send_output_raw(Python ≥ 3.11)。源码中还提供了可写缓冲区的零拷贝发送路径(lib.rs):预分配指定字节数的共享内存/对齐堆存储,通过 Python buffer 协议(memoryview、numpy.frombuffer、struct等)填充后再发送:
with node.send_output_raw("rgb_image", 1920 * 1080 * 3, {"width": 1920}) as buf: numpy.frombuffer(buf, numpy.uint8).reshape(1080, 1920, 3)[:] = frame # 退出 with 块时自动发送在 Python < 3.11 上该方法抛出NotImplementedError并给出升级提示(lib.rs)。
服务、动作与流式通信模式
Python 节点与 Rust 使用相同的元数据键约定(完整指南见 通信模式)。参数是键为字符串的普通 dict。
预定义元数据键:
| 键 | 说明 |
|---|---|
"request_id" | 服务请求/响应关联(UUID v7) |
"goal_id" | 动作目标标识(UUID v7) |
"goal_status" | 动作结果状态:"succeeded"、"aborted"或"canceled" |
"session_id" | 流式会话标识 |
"segment_id" | 会话内的流式分段(整数) |
"seq" | 流式分块序号(整数) |
"fin" | 流式分段最后一个分块(bool) |
"flush" | 丢弃输入上更旧的排队消息(bool) |
服务客户端示例:
import uuid # 发送带唯一 request_id 的请求 request_id = str(uuid.uuid7()) # Python 3.13+;旧版本用 uuid_utils 或 uuid.uuid4() node.send_output("request", data, {"request_id": request_id})服务端示例:
# 原样透传请求中的元数据(含 request_id) node.send_output("response", result, event["metadata"])动作客户端示例:
goal_id = str(uuid.uuid7()) node.send_output("goal", data, {"goal_id": goal_id})流式示例(用户中断时刷新下游队列):
# `metadata` 是参数名 -> 值的扁平 dict,没有外层 "parameters" 键 node.send_output( "text", data, metadata={ "session_id": session_id, "segment_id": 1, "seq": 0, "fin": False, "flush": True, }, )源码还提供send_service_request/send_service_response便捷包装(lib.rs):前者自动注入 UUID v7request_id并返回该 ID,后者校验透传的元数据中包含request_id。
日志:logging自动桥接与显式节点 API
Python 节点既可以使用内置logging模块(推荐),也可以使用显式节点 API。
Pythonlogging模块(自动桥接):
创建Node()时会自动安装一个 handler,把 Pythonlogging模块的路由指向 dora daemon,无需任何配置:
import logging from dora import Node node = Node() # 安装日志桥接 logging.info("Sensor initialized") # -> 结构化 "info" 日志 logging.warning("High temperature") # -> 结构化 "warn" 日志 logging.debug("Raw bytes: %s", data) # -> 结构化 "debug" 日志这些日志条目携带完整元数据(级别、消息、文件路径、行号),并支持 daemon 的min_log_level过滤、send_logs_as路由和dora/logs订阅者。
注意:不要在创建
Node()之前调用logging.basicConfig()。构造函数会设置桥接;先调用basicConfig()可能安装冲突的 handler。
实现细节(lib.rs):setup_logging将logging.basicConfig包装为默认使用HostHandler,后者调用 Rust 函数host_log把logging.LogRecord转码为tracing::Event(按 levelno 映射到 ERROR/WARN/INFO/DEBUG/TRACE),从而让 Python 日志进入 dora 的统一结构化日志管道。
显式节点 API:
log(level, message, target=None, fields=None)
发出带可选 target 和键值字段的结构化日志:
node.log("info", "Processing frame", target="vision") node.log("error", "Sensor timeout", fields={"sensor": "lidar", "retry": "3"})level(str)——日志级别:"error"、"warn"、"info"、"debug"或"trace"。message(str)——日志消息。target(str,可选)——目标模块或子系统名。fields(dict[str, str],可选)——结构化键值上下文字段。
log_error(message)、log_warn(message)、log_info(message)、log_debug(message)、log_trace(message)
常用级别的便捷方法,各自等价于node.log(level, message):
node.log_error("Connection failed") node.log_warn("Temperature elevated") node.log_info("Sensor initialized") node.log_debug("Raw bytes received") node.log_trace("Entering loop iteration")如何选择:
| 方法 | 结构化? | 字段? | 适用场景 |
|---|---|---|---|
logging.info() | 是 | 否 | 通用日志 |
node.log("info", msg, fields={...}) | 是 | 是 | 结构化上下文(sensor_id 等) |
node.log_info(msg) | 是 | 否 | 快速单行 |
print() | 否 | 否 | 遗留代码、快速调试 |
数据流内省:dataflow_descriptor/node_config/dataflow_id
dataflow_descriptor()—— 返回完整数据流描述符(解析后的数据流 YAML)为 Python 字典:
descriptor = node.dataflow_descriptor() print(descriptor["nodes"])node_config()—— 返回数据流描述符中该节点的配置块:
config = node.node_config() model_path = config.get("model", "default.pt")dataflow_id()—— 返回运行中数据流的唯一标识:
print(node.dataflow_id()) # 例如 "a1b2c3d4-..."Rust 实现中这三个方法分别对应dataflow_descriptor()(通过pythonize转换)、node_config()与dataflow_id(lib.rs)。
故障恢复:is_restart()与restart_count()
配合 dora 的故障容忍机制(见 docs/fault-tolerance.md),节点可判断自己是否是被重启的实例:
is_restart()—— 节点是否在之前退出或失败后被重启。用于决定恢复已保存状态还是全新启动:
if node.is_restart(): restore_checkpoint()restart_count()—— 返回节点被重启的次数:首次运行为0,第一次重启后为1,依此类推:
print(f"Restart #{node.restart_count()}")另外,timestamp()返回节点混合逻辑时钟(HLC)的当前 UTC 时间,可与 INPUT 事件的event["metadata"]["timestamp"]对算单事件处理延迟(lib.rs)。
ROS2 集成:merge_external_events(subscription)
将 ROS2 订阅流合并进节点主事件循环。调用后,ROS2 消息会以kind为"external"的事件到达:
from dora import Node, Ros2Context, Ros2Node, Ros2NodeOptions, Ros2Topic node = Node() ros2_context = Ros2Context() ros2_node = ros2_context.new_node("listener", Ros2NodeOptions()) topic = Ros2Topic("/chatter", "std_msgs/String", ros2_node) subscription = ros2_node.create_subscription(topic) node.merge_external_events(subscription) for event in node: if event["kind"] == "external": print("ROS2:", event["value"]) elif event["type"] == "INPUT": print("Dora:", event["id"])参数:
subscription(dora.Ros2Subscription)——通过 dora ROS2 桥创建的 ROS2 订阅。
实现上(lib.rs),merge_external_events将订阅流包装为 futures 流,并在节点内部把 dora 事件流与外部队列用merge_external_send合并;合并后的流上try_recv/drain不可用(返回空),is_empty保守返回False。完整的 ROS2 桥接说明见 docs/ros2-bridge.md。
事件字典结构
事件以普通 Python 字典返回,结构取决于事件类型。
INPUT—— 来自其他节点的输入消息:
{ "type": "INPUT", "id": "camera_image", # 数据流 YAML 中声明的输入 ID "kind": "dora", # "dora" 表示数据流事件,"external" 表示 ROS2 "value": <pyarrow.Array>, # 负载,Apache Arrow 数组 "metadata": { "timestamp": datetime, # UTC 时区 datetime.datetime "open_telemetry_context": "...", # 追踪上下文(若启用) ... # 任意用户元数据 }, }读取数据:
values = event["value"].to_pylist() # 转为 Python 列表 array = event["value"].to_numpy() # 转为 NumPy 数组INPUT_CLOSED—— 输入通道关闭(上游节点结束):
{ "type": "INPUT_CLOSED", "id": "camera_image", "kind": "dora", }STOP—— 数据流正在关闭:
{ "type": "STOP", "id": "MANUAL" | "ALL_INPUTS_CLOSED", # 停止原因 "kind": "dora", }ERROR—— runtime 发生错误:
{ "type": "ERROR", "error": "description of the error", "kind": "dora", }External(ROS2)—— 使用merge_external_events时,ROS2 消息到达为:
{ "kind": "external", "value": <pyarrow.Array>, # ROS2 消息转成的 Arrow 数组 }事件字典在 Rust 侧由PyEvent::to_py_dict构建(apis/python/operator/src/lib.rs),它同时给出了事件类型的完整集合:INPUT、INPUT_CLOSED、INPUT_RECOVERED、NODE_RESTARTED、RELOAD、PARAM_UPDATE、PARAM_DELETED、NODE_FAILED、ERROR、STOP,未来新增变体将以UNKNOWN兜底。
DoraStatus 枚举
作为 operator 的on_event方法返回值,控制事件循环:
from dora import DoraStatus| 值 | 含义 |
|---|---|
DoraStatus.CONTINUE | 继续处理事件(值0) |
DoraStatus.STOP | 停止当前 operator(值1) |
DoraStatus.STOP_ALL | 停止整个数据流(值2) |
该枚举定义在 dora/init.py,是 PythonEnum的子类。
Operator API:在 runtime 进程内运行的算子
Operator 运行在 dora runtime 进程内部(无独立 OS 进程)。它们被定义为名为Operator的 Python 类,包含on_event方法。
Operator 类(用户定义)
创建包含Operator类的 Python 文件:
from dora import DoraStatus class Operator: def __init__(self): # 在此初始化状态 self.count = 0 def on_event(self, dora_event, send_output) -> DoraStatus: if dora_event["type"] == "INPUT": self.count += 1 # 处理输入,可选地发送输出 send_output("result", b"processed", dora_event["metadata"]) return DoraStatus.CONTINUE方法:
__init__(self)—— 加载 operator 时调用一次。在此初始化状态或模型。on_event(self, dora_event, send_output) -> DoraStatus—— 每个到达事件都会调用。必须返回DoraStatus值。
on_event的参数:
dora_event(dict)——事件字典。send_output(可调用对象)——发送输出数据的回调(见下)。
runtime 还会在 operator 实例上设置self.dataflow_descriptor,值为解析后的数据流 YAML 字典。
send_output 回调
send_output回调随on_event传入,用于从 operator 发送数据:
send_output(output_id, data, metadata=None)参数:
output_id(str)——数据流 YAML 中声明的输出名。data(bytes | pyarrow.Array)——负载。metadata(dict,可选)——附加的元数据。传入dora_event["metadata"]可传播追踪上下文。
示例:
import pyarrow as pa from dora import DoraStatus class Operator: def on_event(self, dora_event, send_output) -> DoraStatus: if dora_event["type"] == "INPUT": result = pa.array([42], type=pa.int64()) send_output("output", result, dora_event["metadata"]) return DoraStatus.CONTINUEDataflowBuilder:用 Python 程序化构建数据流
from dora.builder import DataflowBuilder, Node, Operator, Output在 Python 中程序化生成数据流 YAML。完整实现见 apis/python/node/dora/builder.py,配套测试见 apis/python/node/tests/test_builder.py。
DataflowBuilder 类
__init__(name="dora-dataflow")
创建新的数据流构建器:
flow = DataflowBuilder("my-robot")name(str,可选)——数据流名称,默认为"dora-dataflow"。
add_node(id, **kwargs) -> Node
向数据流添加节点,返回用于进一步配置的Node对象:
sender = flow.add_node("sender")id(str)——唯一节点标识。**kwargs——透传到 YAML 的额外节点配置。
to_yaml(path=None) -> str | None
生成数据流的 YAML 表示。给定path时写入文件并返回None;否则返回 YAML 字符串:
# 写入文件 flow.to_yaml("dataflow.yml") # 获取字符串 yaml_str = flow.to_yaml()生成的 YAML 结构为{"nodes": [node.to_dict() ...]},与手写数据流文件(如 examples/python-dataflow/dataflow.yml)一致。
上下文管理器
DataflowBuilder支持with语句:
with DataflowBuilder("my-flow") as flow: flow.add_node("sender").path("sender.py") flow.to_yaml("dataflow.yml")Node 类(构建器)
由DataflowBuilder.add_node()返回。所有 setter 方法返回self以支持链式调用。
path(path) -> Node—— 设置节点可执行文件或脚本路径:
node.path("my_node.py")args(args) -> Node—— 设置节点命令行参数:
node.args("--verbose --port 8080")env(env) -> Node—— 设置节点环境变量:
node.env({"MODEL_PATH": "/models/yolo.pt"})build(command) -> Node—— 设置构建命令(启动前运行):
node.build("pip install -r requirements.txt")git(url, branch=None, tag=None, rev=None) -> Node—— 以 Git 仓库作为节点源码:
node.git("https://github.com/org/repo.git", branch="main")add_operator(operator) -> Node—— 为该节点附加一个Operator:
op = Operator("detector", python="object_detection.py") node.add_operator(op)add_output(output_id) -> Output—— 声明节点输出并返回Output引用,作为其他节点的输入源:
output = sender.add_output("data")add_input(input_id, source, queue_size=None, queue_policy=None) -> Node—— 订阅其他节点的输出:
# 使用 Output 对象 output = sender.add_output("data") receiver.add_input("data", output) # 使用字符串引用 receiver.add_input("tick", "dora/timer/millis/100") # 自定义队列大小 receiver.add_input("images", camera_output, queue_size=2) # 无损输入(队列满时阻塞发送方) receiver.add_input("commands", cmd_output, queue_size=100, queue_policy="backpressure")参数:
input_id(str)——该节点上的输入名。source(str | Output)——字符串("node_id/output_id")或Output对象。queue_size(int,可选)——该输入缓冲的最大消息数。queue_policy(str,可选)——"drop_oldest"(默认)或"backpressure"(缓冲至queue_size的 10 倍后才丢弃)。
从 builder.py 可以看到,queue_policy非drop_oldest/backpressure时会抛出ValueError,queue_size < 1同样报错;测试 test_builder.py 对此有专门覆盖。
to_dict() -> dict—— 返回节点的字典表示,供 YAML 序列化。
Output 类(构建器)
由Node.add_output()返回,表示节点输出的引用,作为add_input()的源:
output = sender.add_output("data") receiver.add_input("sensor_data", output) str(output) # "sender/data"Output.__str__实现为f"{node.id}/{output_id}"(builder.py)。
Operator 类(构建器)
定义嵌入节点 YAML 配置的 operator。
__init__(id, name=None, description=None, build=None, python=None, shared_library=None, send_stdout_as=None)
op = Operator( id="detector", python="object_detection.py", send_stdout_as="detection_text", )参数:
id(str)——唯一 operator 标识。name(str,可选)——显示名。description(str,可选)——人类可读描述。build(str,可选)——加载前运行的构建命令。python(str,可选)——Python operator 文件路径。shared_library(str,可选)——共享库 operator 路径。send_stdout_as(str,可选)——把 operator 的 stdout 作为该 ID 的输出路由。
to_dict() -> dict
返回字典表示,供 YAML 序列化。
CUDA 模块:GPU 张量零拷贝共享(可选扩展)
不属于 1.0 API。这些辅助函数在 1.0 前从
dora包移入了dora_tensor_pool扩展——它们服务于 tensor-pool 传输(位于扩展接缝之后),不受 dora 1.0 兼容性保证覆盖。从libraries/extensions/tensor-pool/python安装。
from dora_tensor_pool import torch_to_ipc_buffer, ipc_buffer_to_ipc_handle, open_ipc_handle这些工具用于节点间通过 CUDA IPC 进行零拷贝 GPU 张量共享。需要带 CUDA 支持的 PyTorch 和带 CUDA 支持的 Numba。
torch_to_ipc_buffer(tensor) -> tuple[pyarrow.Array, dict]
把 PyTorch CUDA 张量转换为包含 CUDA IPC handle 的 Arrow 数组,外加元数据字典。二者一起通过数据流发送即可共享 GPU 内存而不拷贝:
import torch import pyarrow as pa from dora import Node from dora_tensor_pool import torch_to_ipc_buffer node = Node() tensor = torch.randn(1024, 768, device="cuda") ipc_buffer, metadata = torch_to_ipc_buffer(tensor) node.send_output("gpu_data", ipc_buffer, metadata)tensor(torch.Tensor)——CUDA 张量。- 返回:IPC handle 的 int8 Arrow 数组,以及包含 shape、strides、dtype、size、offset 和 source 信息的元数据。
ipc_buffer_to_ipc_handle(handle_buffer, metadata) -> IpcHandle
从收到的 Arrow 缓冲和元数据重建 CUDA IPC handle:
from dora_tensor_pool import ipc_buffer_to_ipc_handle event = node.next() ipc_handle = ipc_buffer_to_ipc_handle(event["value"], event["metadata"])handle_buffer(pyarrow.Array)——来自event["value"]的 Arrow 数组。metadata(dict)——来自event["metadata"]的元数据。- 返回:
dora_tensor_pool.IpcHandle(cudaIpcMemHandle_t的轻量包装;调用.open()把 handle 映射进当前进程并获得设备指针,.close()释放)。
open_ipc_handle(ipc_handle, metadata) -> ContextManager[torch.Tensor]
打开 CUDA IPC handle 并产出 PyTorch 张量。以上下文管理器方式使用以保证正确清理:
from dora_tensor_pool import ipc_buffer_to_ipc_handle, open_ipc_handle event = node.next() ipc_handle = ipc_buffer_to_ipc_handle(event["value"], event["metadata"]) with open_ipc_handle(ipc_handle, event["metadata"]) as tensor: result = tensor * 2 # 直接在 GPU 上使用张量ipc_handle(IpcHandle)——ipc_buffer_to_ipc_handle返回的 handle。metadata(dict)——含 shape、strides、dtype 信息的元数据。- 返回:产出 CUDA
torch.Tensor的上下文管理器。
注意,Node上的register_tensor_pool/write_tensor_pool/read_tensor_pool/free_tensor_pool也属于该可选扩展(lib.rs),仅在以tensor-poolfeature 构建的 wheel 中存在。
快速开始:完整节点示例
一个接收图像、处理并发送结果的完整节点:
#!/usr/bin/env python3 """示例节点:接收消息、转换并发送输出。""" import logging import pyarrow as pa from dora import Node def main(): node = Node() for event in node: if event["type"] == "INPUT": input_id = event["id"] if input_id == "message": values = event["value"].to_pylist() number = values[0] # 创建含多个字段的 struct 数组 result = pa.StructArray.from_arrays( [ pa.array([number * 2]), pa.array([f"Message #{number}"]), ], names=["doubled", "description"], ) node.send_output("transformed", result) logging.info("Transformed message %d", number) elif event["type"] == "STOP": logging.info("Node stopping") break if __name__ == "__main__": main()运行:
dora run dataflow.yml配套的数据流 YAML 参考 examples/python-dataflow/dataflow.yml:其中sender声明输出message,transformer订阅sender/message并输出transformed,receiver同时订阅两个输入。更完整的 Python 数据流示例(含run.rs与动态数据流变体)见 examples/python-dataflow。
DataflowBuilder 完整示例
用代码构建数据流,替代手写 YAML:
#!/usr/bin/env python3 """构建一个简单的 sender -> receiver 数据流。""" from dora.builder import DataflowBuilder, Operator flow = DataflowBuilder("example-flow") # 添加定时驱动的发送节点 sender = flow.add_node("sender") sender.path("sender.py") tick_output = sender.add_output("message") # 添加订阅发送方的接收节点 receiver = flow.add_node("receiver") receiver.path("receiver.py") receiver.add_input("message", tick_output) # 添加带定时输入的节点 timed_node = flow.add_node("periodic") timed_node.path("periodic.py") timed_node.add_input("tick", "dora/timer/millis/100") # 添加带 operator 的节点 runtime_node = flow.add_node("runtime-node") op = Operator("detector", python="object_detection.py") runtime_node.add_operator(op) runtime_node.add_input("image", "camera/image") # 写入或打印 YAML flow.to_yaml("dataflow.yml") print(flow.to_yaml())该示例覆盖了本文介绍的全部构建器 API:path、add_output、add_input(Output对象与"dora/timer/millis/100"字符串两种 source 形态)、add_operator以及to_yaml的两种输出方式。运行数据流依旧使用dora run dataflow.yml,与手写 YAML 完全等效——这也是 DataflowBuilder 的核心价值:将数据流定义纳入 Python 工程化流程(变量、循环、类型检查),同时保持与 YAML 格式的 100% 兼容。
【免费下载链接】doraDORA (Dataflow-Oriented Robotic Architecture) is middleware designed to streamline and simplify the creation of AI-based robotic applications. It offers low latency, composable, and distributed dataflow capabilities. Applications are modeled as directed graphs, also referred to as pipelines.项目地址: https://gitcode.com/GitHub_Trending/do/dora
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考