Skip to main content

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.

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()

Read a topic​

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()

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.