Skip to main content

A public ReductStore server with live ROS 2 robot data, a factory line, and labeled image datasets. Every example below reads the past hour, one record every 30 seconds.

Server
https://play.reduct.store
Read-only token
reductstore
Python SDK
pip install reduct-py

Each example below runs inside this setup, which reads the past hour:

import asyncio
from datetime import datetime, timedelta, timezone

from reduct import Client


async def main():
async with Client("https://play.reduct.store", api_token="reductstore") as client:
stop = datetime.now(timezone.utc)
start = stop - timedelta(hours=1)
... # one of the examples below


asyncio.run(main())

What's on the server

Robotics orion

Decode ROS messages to JSON

ReductROS decodes the stored ROS 2 messages on the server. Here, the boat's GPS fixes.

bucket = await client.get_bucket("orion")
when = {"$each_t": "30s", "#ext": {"ros": {"extract": {}}}}

async for record in bucket.query(
"Pablo05/sensor/gps/fix", start=start, stop=stop, when=when
):
fix = json.loads(await record.read_all())[0]
print(
fix["header"]["stamp"]["sec"],
fix["latitude"],
fix["longitude"],
)
Full script (ros_json.py)
import asyncio
import json
from datetime import datetime, timedelta, timezone

from reduct import Client

URL = "https://play.reduct.store"
TOKEN = "reductstore"


async def main():
async with Client(URL, api_token=TOKEN) as client:
bucket = await client.get_bucket("orion")
stop = datetime.now(timezone.utc)
start = stop - timedelta(hours=1)
when = {"$each_t": "30s", "#ext": {"ros": {"extract": {}}}}

async for record in bucket.query(
"Pablo05/sensor/gps/fix", start=start, stop=stop, when=when
):
fix = json.loads(await record.read_all())[0]
print(
fix["header"]["stamp"]["sec"],
fix["latitude"],
fix["longitude"],
)


asyncio.run(main())

Expected output

1790784087 42.35824 -71.08775
1790784313 42.35819 -71.08763
1790784477 42.35816 -71.08755
1790784654 42.35818 -71.08761
...
ReductROS raw messages →

Get camera frames as JPEG

Raw sensor_msgs/Image frames from the infrared camera, converted to JPEG on the server.

bucket = await client.get_bucket("orion")
when = {
"$each_t": "30s",
"#ext": {"ros": {"extract": {"encode": {"data": "jpeg"}}}},
}

async for record in bucket.query(
"right_ir/rotated/image_raw", start=start, stop=stop, when=when
):
image = json.loads(await record.read_all())[0]
with open(f"{record.timestamp}.jpg", "wb") as file:
file.write(base64.b64decode(image["data"]))
print(
f"{record.timestamp}.jpg", image["width"], image["height"]
)
Full script (ros_jpeg.py)
import asyncio
import base64
import json
from datetime import datetime, timedelta, timezone

from reduct import Client

URL = "https://play.reduct.store"
TOKEN = "reductstore"


async def main():
async with Client(URL, api_token=TOKEN) as client:
bucket = await client.get_bucket("orion")
stop = datetime.now(timezone.utc)
start = stop - timedelta(hours=1)
when = {
"$each_t": "30s",
"#ext": {"ros": {"extract": {"encode": {"data": "jpeg"}}}},
}

async for record in bucket.query(
"right_ir/rotated/image_raw", start=start, stop=stop, when=when
):
image = json.loads(await record.read_all())[0]
with open(f"{record.timestamp}.jpg", "wb") as file:
file.write(base64.b64decode(image["data"]))
print(
f"{record.timestamp}.jpg", image["width"], image["height"]
)


asyncio.run(main())

Expected output

1790784334847174.jpg 640 512
1790784498326544.jpg 640 512
1790784675706322.jpg 640 512
...
Infrared frame 1 from the orion boatInfrared frame 2 from the orion boatInfrared frame 3 from the orion boat
Encoding binary fields →

Filter by label

Only colour camera frames recorded while the GPS altitude label was between 110 and 120 m.

bucket = await client.get_bucket("orion")
when = {
"$each_t": "30s",
"$and": [{"&gps_z": {"$gt": 110}}, {"&gps_z": {"$lt": 120}}],
}

async for record in bucket.query(
"right_camera/image_color/compressed",
start=start,
stop=stop,
when=when,
):
print(record.timestamp, record.labels["gps_z"], record.size)
Full script (ros_labels.py)
import asyncio
from datetime import datetime, timedelta, timezone

from reduct import Client

URL = "https://play.reduct.store"
TOKEN = "reductstore"


async def main():
async with Client(URL, api_token=TOKEN) as client:
bucket = await client.get_bucket("orion")
stop = datetime.now(timezone.utc)
start = stop - timedelta(hours=1)
when = {
"$each_t": "30s",
"$and": [{"&gps_z": {"$gt": 110}}, {"&gps_z": {"$lt": 120}}],
}

async for record in bucket.query(
"right_camera/image_color/compressed",
start=start,
stop=stop,
when=when,
):
print(record.timestamp, record.labels["gps_z"], record.size)


asyncio.run(main())

Expected output

1790784328478195 118.0331 50676
1790784669343610 113.8821 50048
1790785497907116 119.8687 39300
...
Conditional queries →

Export an hour as MCAP

Camera, GPS, odometry, and TF merged into one MCAP file you can open in Foxglove.

TOPICS = [
"right_camera/image_color/compressed",
"Pablo05/sensor/gps/fix",
"Pablo05/odom",
"tf",
]

bucket = await client.get_bucket("orion")
when = {
"$each_t": "30s",
"#ext": {
"ros": {"export": {"format": "mcap", "duration": "1h"}}
},
}

async for record in bucket.query(
TOPICS, start=start, stop=stop, when=when
):
with open("orion.mcap", "wb") as file:
file.write(await record.read_all())
print("orion.mcap", record.content_type, record.size, "bytes")
Full script (ros_mcap.py)
import asyncio
from datetime import datetime, timedelta, timezone

from reduct import Client

URL = "https://play.reduct.store"
TOKEN = "reductstore"

TOPICS = [
"right_camera/image_color/compressed",
"Pablo05/sensor/gps/fix",
"Pablo05/odom",
"tf",
]


async def main():
async with Client(URL, api_token=TOKEN) as client:
bucket = await client.get_bucket("orion")
stop = datetime.now(timezone.utc)
start = stop - timedelta(hours=1)
when = {
"$each_t": "30s",
"#ext": {
"ros": {"export": {"format": "mcap", "duration": "1h"}}
},
}

async for record in bucket.query(
TOPICS, start=start, stop=stop, when=when
):
with open("orion.mcap", "wb") as file:
file.write(await record.read_all())
print("orion.mcap", record.content_type, record.size, "bytes")


asyncio.run(main())

Expected output

orion.mcap application/mcap 395404 bytes
Exporting MCAP →

Industrial factory

Find vibration above a threshold

1-second vibration chunks labeled with their RMS. The query returns only the bearing fault.

bucket = await client.get_bucket("factory")
when = {"$each_t": "30s", "&rms": {"$gt": 0.8}}

async for record in bucket.query(
"vibration", start=start, stop=stop, when=when
):
print(
record.timestamp,
record.labels["rms"],
record.labels["peak"],
record.size,
)
Full script (factory_vibration.py)
import asyncio
from datetime import datetime, timedelta, timezone

from reduct import Client

URL = "https://play.reduct.store"
TOKEN = "reductstore"


async def main():
async with Client(URL, api_token=TOKEN) as client:
bucket = await client.get_bucket("factory")
stop = datetime.now(timezone.utc)
start = stop - timedelta(hours=1)
when = {"$each_t": "30s", "&rms": {"$gt": 0.8}}

async for record in bucket.query(
"vibration", start=start, stop=stop, when=when
):
print(
record.timestamp,
record.labels["rms"],
record.labels["peak"],
record.size,
)


asyncio.run(main())

Expected output

1790786100000000 0.922 1.769 4000
1790786130000000 0.922 1.858 4000
1790786160000000 0.923 1.792 4000
...
Conditional queries →

Run SQL on PLC data

ReductSelect runs SQL on each stored CSV record. The temperature climbs during the fault.

SQL = """
SELECT round(avg(temperature), 1) AS temperature,
round(max(pressure), 2) AS pressure
FROM ENTRY()
"""

bucket = await client.get_bucket("factory")
when = {
"$each_t": "30s",
"#ext": {"select": {"sql": SQL, "export": {"format": "json"}}},
}

async for record in bucket.query(
"plc", start=start, stop=stop, when=when
):
print(record.timestamp, (await record.read_all()).decode())
Full script (factory_sql.py)
import asyncio
from datetime import datetime, timedelta, timezone

from reduct import Client

URL = "https://play.reduct.store"
TOKEN = "reductstore"

SQL = """
SELECT round(avg(temperature), 1) AS temperature,
round(max(pressure), 2) AS pressure
FROM ENTRY()
"""


async def main():
async with Client(URL, api_token=TOKEN) as client:
bucket = await client.get_bucket("factory")
stop = datetime.now(timezone.utc)
start = stop - timedelta(hours=1)
when = {
"$each_t": "30s",
"#ext": {"select": {"sql": SQL, "export": {"format": "json"}}},
}

async for record in bucket.query(
"plc", start=start, stop=stop, when=when
):
print(record.timestamp, (await record.read_all()).decode())


asyncio.run(main())

Expected output

1790784060000000 [{"temperature":60.9,"pressure":4.25}]
1790784090000000 [{"temperature":60.9,"pressure":4.27}]
...
1790786220000000 [{"temperature":64.5,"pressure":4.6}]
1790786280000000 [{"temperature":66.4,"pressure":4.58}]
...
ReductSelect →

Read MQTT topics

Each MQTT topic is its own entry, so topics are queried by name like any other stream.

TOPICS = [
"mqtt/line1/temperature",
"mqtt/line1/pressure",
"mqtt/line1/state",
]

bucket = await client.get_bucket("factory")

for topic in TOPICS:
async for record in bucket.query(
topic, start=start, stop=stop, when={"$each_t": "30s"}
):
print(
topic,
record.timestamp,
(await record.read_all()).decode(),
)
Full script (factory_mqtt.py)
import asyncio
from datetime import datetime, timedelta, timezone

from reduct import Client

URL = "https://play.reduct.store"
TOKEN = "reductstore"

TOPICS = [
"mqtt/line1/temperature",
"mqtt/line1/pressure",
"mqtt/line1/state",
]


async def main():
async with Client(URL, api_token=TOKEN) as client:
bucket = await client.get_bucket("factory")
stop = datetime.now(timezone.utc)
start = stop - timedelta(hours=1)

for topic in TOPICS:
async for record in bucket.query(
topic, start=start, stop=stop, when={"$each_t": "30s"}
):
print(
topic,
record.timestamp,
(await record.read_all()).decode(),
)


asyncio.run(main())

Expected output

mqtt/line1/temperature 1790784065000000 {"value": 60.86, "unit": "C"}
mqtt/line1/temperature 1790784095000000 {"value": 61.24, "unit": "C"}
...
mqtt/line1/pressure 1790784065000000 {"value": 4.177, "unit": "bar"}
...
mqtt/line1/state 1790784065000000 {"value": "running"}
...
MQTT data storage →

Datasets datasets

Select labeled training images

MNIST handwritten digits labeled with their value. The query picks five images of the digit 7.

bucket = await client.get_bucket("datasets")
when = {"&digit": {"$eq": 7}, "$limit": 5}

async for record in bucket.query("mnist_training", when=when):
with open(f"digit-{record.timestamp}.png", "wb") as file:
file.write(await record.read_all())
print(record.timestamp, record.labels["digit"], record.size)
Full script (dataset_labels.py)
import asyncio

from reduct import Client

URL = "https://play.reduct.store"
TOKEN = "reductstore"


async def main():
async with Client(URL, api_token=TOKEN) as client:
bucket = await client.get_bucket("datasets")
when = {"&digit": {"$eq": 7}, "$limit": 5}

async for record in bucket.query("mnist_training", when=when):
with open(f"digit-{record.timestamp}.png", "wb") as file:
file.write(await record.read_all())
print(record.timestamp, record.labels["digit"], record.size)


asyncio.run(main())

Expected output

23957 7 208
23958 7 207
23959 7 218
...
Handwritten digit 7, MNIST record 23957Handwritten digit 7, MNIST record 23958Handwritten digit 7, MNIST record 23959
Stream a dataset into PyTorch →

Use it in your tools

The factory line is simulated. Need another kind of data? Tell us.