Streaming
Use streaming when producers should publish data independently of consumers: telemetry, state, events, sensor samples, images, audio, and pipeline outputs. A topic identifies the data channel; readers subscribe to one or more topics.
Write a value
The following examples publish through MQTT so all three languages can participate in the same system.
- Python
- C++
- TypeScript
from luxai.magpie.transport import MqttConnection, MqttStreamWriter
connection = MqttConnection("mqtt://broker.example.com:1883")
connection.connect()
writer = MqttStreamWriter(connection)
writer.write({"temperature": 22.5, "unit": "C"}, topic="robot/sensors/temperature")
writer.close()
connection.disconnect()
#include <magpie/frames/primitive_frames.hpp>
#include <magpie/transport/mqtt_connection.hpp>
#include <magpie/transport/mqtt_stream_writer.hpp>
using namespace magpie;
auto connection = std::make_shared<MqttConnection>("mqtt://broker.example.com:1883");
connection->connect();
MqttStreamWriter writer(connection);
writer.write(StringFrame("online"), "robot/status");
writer.close();
connection->disconnect();
import {MqttConnection, MqttStreamWriter} from '@luxai-qtrobot/magpie'
const connection = new MqttConnection('mqtt://broker.example.com:1883')
await connection.connect()
const writer = new MqttStreamWriter(connection)
await writer.write({temperature: 22.5, unit: 'C'}, 'robot/sensors/temperature')
writer.close()
await connection.disconnect()
Read a topic
- Python
- C++
- TypeScript
from luxai.magpie.transport import MqttConnection, MqttStreamReader
connection = MqttConnection("mqtt://broker.example.com:1883")
connection.connect()
reader = MqttStreamReader(connection, topic="robot/sensors/#")
try:
while True:
value, topic = reader.read(timeout=5.0)
print(topic, value)
finally:
reader.close()
connection.disconnect()
MqttStreamReader reader(connection, "robot/sensors/#");
std::unique_ptr<Frame> frame;
std::string topic;
if (reader.read(frame, topic, 5.0)) {
Logger::info("received " + topic + " as " + frame->name());
}
import {MqttConnection, MqttStreamReader, TimeoutError} from '@luxai-qtrobot/magpie'
const connection = new MqttConnection('mqtt://broker.example.com:1883')
await connection.connect()
const reader = new MqttStreamReader(connection, {topic: 'robot/sensors/#'})
try {
const [value, topic] = await reader.read(5)
console.log(topic, value)
} catch (error) {
if (!(error instanceof TimeoutError)) throw error
}
Topic design
Use topic names that describe data, not consumers. robot/camera/color is more reusable than dashboard/camera. A practical hierarchy starts broad and becomes specific:
robot/{robot_id}/status
robot/{robot_id}/sensors/temperature
robot/{robot_id}/camera/color
robot/{robot_id}/audio/microphone
MQTT readers support + for one level and # for the remaining hierarchy. ZeroMQ topic filtering is prefix-oriented. Keep the portable part of your topic design simple when readers may change transport.
Backpressure and slow consumers
Streams should not allow an unlimited backlog. Choose queue sizes for the meaning of the data:
- For live camera or robot state, prefer the newest item and allow old samples to be dropped.
- For audit events, use a durable transport and application-level persistence.
- For commands that require a result, use RPC, not streaming.
Always close readers and writers during shutdown. Close a shared connection only after every component using it has stopped.