HFlow is een open-source SDK ontwikkeld door Hebbian Robotics om de flessenhals in dataverwerking binnen de robotica op te lossen. Het framework biedt tools voor het beheren van multimodale datapijplijnen, waardoor teams van elke omvang data kunnen verwerken, controleren en cureren.
Kernfunctionaliteiten:
- Data-orchestratie: HFlow ondersteunt het MCAP-formaat en LeRobot Dataset v3, en gebruikt Airflow 3 DAG's voor de visuele weergave en uitvoering van pijplijnen.
- Provenancie en Catalogisering: Elke dataset krijgt herkomstinformatie (provenance). Metadata en kwaliteitsmetingen worden opgeslagen in een Parquet-catalogus, waardoor data bevraagd kan worden via DuckDB SQL zonder de zware bronbestanden te laden.
- Flexibele Integratie: Gebruikers kunnen hun eigen Python-transformaties en kwaliteitscontroles implementeren via adapters, waardoor bestaande code behouden blijft.
- Levenscyclus: Het proces volgt een vaste route van collectie naar ingestie, curatie en uiteindelijke levering voor training.
Het systeem is ontworpen voor schaalbaarheid en kan worden gedraaid via Docker Compose of in bestaande Airflow-omgevingen, met ondersteuning voor cloud-storage zoals S3, GS en Azure.
HFlow: Open-source SDK voor schaalbare multimodale datapijplijnen in robotica en fysieke AI
Het probleem: de bottleneck in dataverwerking
Het verwerken van data is een grote flessenhals in de robotica. Een corpus kan video's, toestanden (states), acties, tijdstempels en metadata van diverse opnamesystemen combineren. Teams merken dit probleem vaak als eerste bij de kwaliteitscontrole: het bepalen of camera's zijn vastgelopen, streams uit sync zijn gelopen, vereiste onderwerpen (topics) zijn verdwenen of er dubbele opnames in het corpus zijn terechtgekomen. Naarmate het corpus groeit, maken gefragmenteerde scripts het moeilijk om te weten wat er is uitgevoerd, de resultaten te controleren of een dataset te reproduceren.
Teams kunnen beginnen met de ingebouwde controles van HFlow, nieuwe transformaties, controles, labels en verrijkingen schrijven, of hun reeds bestaande verwerkingscode koppelen. HFlow regelt de orchestratie, opslag, versiering en curatie rondom deze stappen.
HFlow voorziet elke verwerkte episode van herkomstinformatie (provenance), rendert de pijplijn als een grafiek en legt metadata en kwaliteitsbewijzen vast in een doorzoekbare catalogus. Hierdoor kun je herleiden hoe outputs zijn geproduceerd, elke fase monitoren en een corpus onderzoeken zonder de onderliggende opnames te hoeven laden.
De grenzen van HFlow
Input
- Directe ondersteuning voor standaard MCAP-episodes.
- LeRobot Dataset v3-repositories via
hflow import lerobot.
Verwerking
- Jouw eigen Python-transformaties, controles, labels en verrijkingen.
Uitvoering
- In-process voor ontwikkeling.
- Gegenereerde Airflow 3 DAG's voor geplande runs.
Duurzame output
- Canonieke MCAP-episodes, provenance, artefacten en een Parquet-catalogus.
Curatie
- DuckDB SQL die een versie-gepinde manifest schrijft.
De levenscyclus van data
Menselijke en robotdata doorlopen een levenscyclus van vier fasen:
collectie (landing bucket) --> ingestie (transform -> QC gate -> enrich, als catalogus / Airflow DAG) --> curatie (SQL over episode) --> levering (gecureerde MCAP + manifest; conversie voor training)
Belangrijke kenmerken:
- Behoud van eigen code: Je verwerkingscode blijft van jou. Transformaties, kwaliteitscontroles, labels en verrijkingen zijn eenvoudige Python-functies in je eigen omgeving. Bestaande code wordt aangesloten via kleine adapters in plaats van dat het herschreven moet worden voor een propriëtair framework.
- MCAP-formaat: Episodes zijn MCAP, de container die ROS 2 standaard opneemt en die Foxglove/Rerun direct kunnen openen. Er wordt gebruikgemaakt van twee optimalisaties (zoals beschreven in het artikel van Dyna): in-band H.264 met een GOP-lengte die is afgestemd op de leeswijze, en topic-group chunking (cameraproducten en statusstreams delen nooit een chunk, waardoor een trainingssample slechts één read per groep kost in plaats van één per topic).
- Herkomst (Provenance): Verwerkte episodes bevatten hun herkomst. Het bestand zelf registreert het schema, de pijplijn en de tool-versies die het hebben geproduceerd, plus de bron-URI indien beschikbaar. Catalogusrecords koppelen metingen en resultaten aan stap-versies, waardoor het gemakkelijker is om een slecht resultaat terug te herleiden naar de oorsprong.
- Visuele pijplijn: De pijplijn is zichtbaar als een grafiek. HFlow rendert Airflow DAG's, zodat je kunt zien hoe fasen verbonden zijn en je de status van taken, logs, retries en reruns kunt monitoren.
- Herbruikbaar bewijs: Kwaliteitscontroles produceren herbruikbaar bewijs. Accessors extraheren de inputs die bestaande verwerkingscode verwacht (numpy arrays, MP4-paden, JPEG-frames), en resultaten worden opgeslagen als doorzoekbare metingen in plaats van hardgecodeerde oordelen. Verschillende datasets kunnen verschillende drempelwaarden toepassen zonder de media opnieuw te hoeven verwerken.
- Querying zonder laden: Je kunt het corpus bevragen zonder de opnames te laden. Metadata, kwaliteitsmetingen, tags, versie-stempels en locaties van artefacten bevinden zich in de Parquet-catalogus. DuckDB kan vragen over het hele corpus beantwoorden en manifests bouwen zonder de onderliggende MCAP-bestanden te openen.
Je kunt op elk moment de DuckDB-browser over de catalogus openen, zelfs voordat de eerste run start: hflow catalog ui
Hosting en schaalbaarheid
De open-source implementatie is ontworpen om eenvoudig te beheren: draai één single-tenant workspace met de inbegrepen Docker Compose-runtime, of implementeer de gegenereerde DAG-bundel in een Airflow 3-omgeving die je al beheert. Het bevat geen gebruikersaccounts, RBAC of een multi-tenant control plane.
De data plane is gescheiden van account- en control-plane zaken, zodat dezelfde engine kan worden geschaald als meerdere geïsoleerde workspaces (bijvoorbeeld één per team of klant) achter een externe control plane. Dit is het beoogde pad naar een toekomstige hosted versie, maar de hosted control plane is niet geïmplementeerd in deze repository en is geen commitment voor de pre-v1 release.
Voor details over het data-plane contract, de workspace-unit, de trust model en huidige limieten kan worden verwezen naar docs/HOSTING.md.
Installatie en gebruik
Installatie
Installeer de SDK via PyPI met uv: uv add hflow
Let op: Het Hebbian Robotics-project start bij versie 0.2.0. Eerdere 0.1.x releases onder dezelfde PyPI-naam behoorden tot een niet-gerelateerd, inactief project voordat de naam werd overgedragen.
Quickstart
Om de gebundelde quickstart van de repository uit te voeren:
git clone https://github.com/Hebbian-Robotics/hflow.git
cd hflow
uv sync --locked
uv run python examples/quickstart.py
De quickstart synthetiseert een kleine multimodale episode met camera- en statusstreams wanneer er geen inputbestand is gegeven, voert de pijplijn in-process uit en schrijft de outputs naar de (gitignored) data/ directory. Er is geen Docker of Airflow voor nodig. Om je eigen opname te gebruiken: uv run python examples/quickstart.py path/to/episode.mcap
Gebruik uv run hflow --help voor de CLI.
LeRobot import
Om een LeRobot Dataset v3 episode te importeren naar de canonieke MCAP-boundary:
uv run hflow import lerobot \
--repo lerobot/pusht --revision main \
--camera observation.image --episode-index 0 \
--output-dir ./data/lerobot_pusht
Codevoorbeeld
Je kunt beginnen in zes regels code. Dit uitgebreidere voorbeeld gebruikt een robot-teleoperation episode, maar dezelfde interface is van toepassing op egocentrische video's en andere physical-AI opnames.
import hflow
from hflow.checks import camera_frame_stats
from your_existing_qc import check_joint_smoothness # gebruik je bestaande controles
app = hflow.App("kitchen-pipeline") # data root: $HFLOW_DATA_ROOT, hflow.toml, of ./data
@app.check(version="1")
def joint_smoothness(ep: hflow.Episode) -> hflow.CheckResult:
joints = ep.channel("/joint_states").to_numpy() # extractie
result = check_joint_smoothness(joints, rate_hz=100) # bestaande logica
return hflow.CheckResult(measurements=result) # registratie
@app.check(version="1", critical=True)
def camera_blackout(ep: hflow.Episode) -> hflow.CheckResult:
camera_topic = next(topic for topic in ep.cameras if "wrist_cam" in topic)
evidence = camera_frame_stats(ep, cameras=[camera_topic])
black_frame_percent = evidence.measurements[f"{camera_topic}/black_frame_pct"]
assert isinstance(black_frame_percent, float)
return hflow.CheckResult(
measurements={"black_pct": black_frame_percent},
verdict=black_frame_percent < 50.0, # jouw drempelwaarde
)
if __name__ == "__main__":
app.test("episode_0001.mcap") # volledige pijplijn, in-process, zonder infra
Elke check, verrijking en afgeleid kanaal declareert een versie. HFlow slaat deze waarde exact op; behoud deze voor refactors die het gedrag behouden, en verhoog deze wanneer oude en nieuwe resultaten niet langer als vergelijkbaar beschouwd moeten worden.
Curatie gebeurt achteraf via hflow.curate(data_root / "catalog", sql, output="manifest.parquet") of via de command line: hflow curate "<sql>".
Voorbeeld SQL voor curatie:
SELECT episode_id, uri FROM episodes
WHERE task = 'fold_napkin'
AND status = 'ok'
AND black_pct < 1.0 -- percentage, drempelwaarde van gebruiker
AND pipeline_version = 'a41c9f27b3d8' # pin één specifieke verwerkingsgeneratie
Ontwerpprincipes
- Democratiseer de architectuur, stel optimalisaties uit: Behoud de nuttige workflow en standaard interfaces op kleine schaal. Label elk mechanisme voor productieschaal eerlijk als geïmplementeerd, vereenvoudigd, uitgesteld of buiten scope.
- Bewijs, geen oordelen: Controles registreren metingen met dekking; het pass/fail-beleid behoort tot de consument tijdens de curatie. Kwaliteitstags routeren episodes; ze verwijderen nooit data.
- Standaardformaten bij elke grens: MCAP-episodes, Parquet-catalogi, Airflow DAG's. De code bestaat alleen waar het formaat een brug vereist of waar een valkuil echt niet voor de hand ligt.
- Jouw code blijft jouw code: Bestaande transformaties, controles en verrijkingen worden aangesloten via kleine adapters in plaats van dat ze herschreven moeten worden.
Systeemvereisten
- Python: $\ge$ 3.11
- Runtime: Docker (voor de pijplijn runtime;
app.test() heeft dit niet nodig), of een eigen Airflow-implementatie (Astronomer, MWAA, Cloud Composer, self-managed).
- Opslag: De eerste
hflow up downloadt ongeveer 2 GB aan container-images en bouwt de task venv.
- Cloud: Native
s3://, gs:// en Azure data roots gebruiken de optionele bucket-backend (uv sync --extra bucket); lokale paden hebben dit niet nodig.
- Video: Op Linux x8664/aarch64 downloadt de eerste video-operatie een checksum-geverifieerde, gepinde ffmpeg/ffprobe build in de user cache. Gebruik
HFLOWFFMPEG en HFLOW_FFPROBE om eigen binaries te gebruiken.
- Windows: Ondersteund via WSL2 (Airflow draait niet native op Windows).
Documentatie en Referenties
Documentatiebronnen
- Home: Start per taak, gevolgd door tutorials, how-to gidsen, referenties of uitleg.
- FAQ: Informatie over formaten, infrastructuur, schaal, projectscope en huidige releasestatus.
- Data Stack: Hoe HFlow past binnen de robotics data stack (MCAP, Airflow, Foxglove, Rerun, DuckDB, object storage, training formats).
- Voorbeelden: Uitvoerbare commando's, vereisten en links naar gidsen.
- Architectuur: Matrix van geïmplementeerde, vereenvoudigde, uitgestelde en buiten-scope functies.
- OpenAI Vision: Specifieke gids voor het aanroepen van OpenAI vision vanuit een stap.
- Bijdragen: Ontwikkelsetup, validatiecommando's en PR-verwachtingen.
Referenties
- Dyna Robotics: Training Dyna-2 at million-hour scale, repeatably.
- MCAP-specificatie en Python-libraries (Foxglove).
foxglove.CompressedVideo schema: in-band H.264/H.265/VP9/AV1 video in MCAP.
- Apache Airflow.
- DuckDB.
- Foxglove en Rerun.
- FFmpeg.
- Pareto: Het robotics data curatieplatform van Hebbian Robotics.
HFlow: Open-source SDK voor schaalbare multimodale datapijplijnen in robotica en fysieke AI
Het probleem: de bottleneck in dataverwerking
Het verwerken van data is een grote flessenhals in de robotica. Een corpus kan video's, toestanden (states), acties, tijdstempels en metadata van diverse opnamesystemen combineren. Teams merken dit probleem vaak als eerste bij de kwaliteitscontrole: het bepalen of camera's zijn vastgelopen, streams uit sync zijn gelopen, vereiste onderwerpen (topics) zijn verdwenen of er dubbele opnames in het corpus zijn terechtgekomen. Naarmate het corpus groeit, maken gefragmenteerde scripts het moeilijk om te weten wat er is uitgevoerd, de resultaten te controleren of een dataset te reproduceren.
Teams kunnen beginnen met de ingebouwde controles van HFlow, nieuwe transformaties, controles, labels en verrijkingen schrijven, of hun reeds bestaande verwerkingscode koppelen. HFlow regelt de orchestratie, opslag, versiering en curatie rondom deze stappen.
HFlow voorziet elke verwerkte episode van herkomstinformatie (provenance), rendert de pijplijn als een grafiek en legt metadata en kwaliteitsbewijzen vast in een doorzoekbare catalogus. Hierdoor kun je herleiden hoe outputs zijn geproduceerd, elke fase monitoren en een corpus onderzoeken zonder de onderliggende opnames te hoeven laden.
De grenzen van HFlow
Input
- Directe ondersteuning voor standaard MCAP-episodes.
- LeRobot Dataset v3-repositories via
hflow import lerobot.
Verwerking
- Jouw eigen Python-transformaties, controles, labels en verrijkingen.
Uitvoering
- In-process voor ontwikkeling.
- Gegenereerde Airflow 3 DAG's voor geplande runs.
Duurzame output
- Canonieke MCAP-episodes, provenance, artefacten en een Parquet-catalogus.
Curatie
- DuckDB SQL die een versie-gepinde manifest schrijft.
De levenscyclus van data
Menselijke en robotdata doorlopen een levenscyclus van vier fasen:
collectie (landing bucket) --> ingestie (transform -> QC gate -> enrich, als catalogus / Airflow DAG) --> curatie (SQL over episode) --> levering (gecureerde MCAP + manifest; conversie voor training)
Belangrijke kenmerken:
- Behoud van eigen code: Je verwerkingscode blijft van jou. Transformaties, kwaliteitscontroles, labels en verrijkingen zijn eenvoudige Python-functies in je eigen omgeving. Bestaande code wordt aangesloten via kleine adapters in plaats van dat het herschreven moet worden voor een propriëtair framework.
- MCAP-formaat: Episodes zijn MCAP, de container die ROS 2 standaard opneemt en die Foxglove/Rerun direct kunnen openen. Er wordt gebruikgemaakt van twee optimalisaties (zoals beschreven in het artikel van Dyna): in-band H.264 met een GOP-lengte die is afgestemd op de leeswijze, en topic-group chunking (cameraproducten en statusstreams delen nooit een chunk, waardoor een trainingssample slechts één read per groep kost in plaats van één per topic).
- Herkomst (Provenance): Verwerkte episodes bevatten hun herkomst. Het bestand zelf registreert het schema, de pijplijn en de tool-versies die het hebben geproduceerd, plus de bron-URI indien beschikbaar. Catalogusrecords koppelen metingen en resultaten aan stap-versies, waardoor het gemakkelijker is om een slecht resultaat terug te herleiden naar de oorsprong.
- Visuele pijplijn: De pijplijn is zichtbaar als een grafiek. HFlow rendert Airflow DAG's, zodat je kunt zien hoe fasen verbonden zijn en je de status van taken, logs, retries en reruns kunt monitoren.
- Herbruikbaar bewijs: Kwaliteitscontroles produceren herbruikbaar bewijs. Accessors extraheren de inputs die bestaande verwerkingscode verwacht (numpy arrays, MP4-paden, JPEG-frames), en resultaten worden opgeslagen als doorzoekbare metingen in plaats van hardgecodeerde oordelen. Verschillende datasets kunnen verschillende drempelwaarden toepassen zonder de media opnieuw te hoeven verwerken.
- Querying zonder laden: Je kunt het corpus bevragen zonder de opnames te laden. Metadata, kwaliteitsmetingen, tags, versie-stempels en locaties van artefacten bevinden zich in de Parquet-catalogus. DuckDB kan vragen over het hele corpus beantwoorden en manifests bouwen zonder de onderliggende MCAP-bestanden te openen.
Je kunt op elk moment de DuckDB-browser over de catalogus openen, zelfs voordat de eerste run start: hflow catalog ui
Hosting en schaalbaarheid
De open-source implementatie is ontworpen om eenvoudig te beheren: draai één single-tenant workspace met de inbegrepen Docker Compose-runtime, of implementeer de gegenereerde DAG-bundel in een Airflow 3-omgeving die je al beheert. Het bevat geen gebruikersaccounts, RBAC of een multi-tenant control plane.
De data plane is gescheiden van account- en control-plane zaken, zodat dezelfde engine kan worden geschaald als meerdere geïsoleerde workspaces (bijvoorbeeld één per team of klant) achter een externe control plane. Dit is het beoogde pad naar een toekomstige hosted versie, maar de hosted control plane is niet geïmplementeerd in deze repository en is geen commitment voor de pre-v1 release.
Voor details over het data-plane contract, de workspace-unit, de trust model en huidige limieten kan worden verwezen naar docs/HOSTING.md.
Installatie en gebruik
Installatie
Installeer de SDK via PyPI met uv: uv add hflow
Let op: Het Hebbian Robotics-project start bij versie 0.2.0. Eerdere 0.1.x releases onder dezelfde PyPI-naam behoorden tot een niet-gerelateerd, inactief project voordat de naam werd overgedragen.
Quickstart
Om de gebundelde quickstart van de repository uit te voeren:
git clone https://github.com/Hebbian-Robotics/hflow.git
cd hflow
uv sync --locked
uv run python examples/quickstart.py
De quickstart synthetiseert een kleine multimodale episode met camera- en statusstreams wanneer er geen inputbestand is gegeven, voert de pijplijn in-process uit en schrijft de outputs naar de (gitignored) data/ directory. Er is geen Docker of Airflow voor nodig. Om je eigen opname te gebruiken: uv run python examples/quickstart.py path/to/episode.mcap
Gebruik uv run hflow --help voor de CLI.
LeRobot import
Om een LeRobot Dataset v3 episode te importeren naar de canonieke MCAP-boundary:
uv run hflow import lerobot \
--repo lerobot/pusht --revision main \
--camera observation.image --episode-index 0 \
--output-dir ./data/lerobot_pusht
Codevoorbeeld
Je kunt beginnen in zes regels code. Dit uitgebreidere voorbeeld gebruikt een robot-teleoperation episode, maar dezelfde interface is van toepassing op egocentrische video's en andere physical-AI opnames.
import hflow
from hflow.checks import camera_frame_stats
from your_existing_qc import check_joint_smoothness # gebruik je bestaande controles
app = hflow.App("kitchen-pipeline") # data root: $HFLOW_DATA_ROOT, hflow.toml, of ./data
@app.check(version="1")
def joint_smoothness(ep: hflow.Episode) -> hflow.CheckResult:
joints = ep.channel("/joint_states").to_numpy() # extractie
result = check_joint_smoothness(joints, rate_hz=100) # bestaande logica
return hflow.CheckResult(measurements=result) # registratie
@app.check(version="1", critical=True)
def camera_blackout(ep: hflow.Episode) -> hflow.CheckResult:
camera_topic = next(topic for topic in ep.cameras if "wrist_cam" in topic)
evidence = camera_frame_stats(ep, cameras=[camera_topic])
black_frame_percent = evidence.measurements[f"{camera_topic}/black_frame_pct"]
assert isinstance(black_frame_percent, float)
return hflow.CheckResult(
measurements={"black_pct": black_frame_percent},
verdict=black_frame_percent < 50.0, # jouw drempelwaarde
)
if __name__ == "__main__":
app.test("episode_0001.mcap") # volledige pijplijn, in-process, zonder infra
Elke check, verrijking en afgeleid kanaal declareert een versie. HFlow slaat deze waarde exact op; behoud deze voor refactors die het gedrag behouden, en verhoog deze wanneer oude en nieuwe resultaten niet langer als vergelijkbaar beschouwd moeten worden.
Curatie gebeurt achteraf via hflow.curate(data_root / "catalog", sql, output="manifest.parquet") of via de command line: hflow curate "<sql>".
Voorbeeld SQL voor curatie:
SELECT episode_id, uri FROM episodes
WHERE task = 'fold_napkin'
AND status = 'ok'
AND black_pct < 1.0 -- percentage, drempelwaarde van gebruiker
AND pipeline_version = 'a41c9f27b3d8' # pin één specifieke verwerkingsgeneratie
Ontwerpprincipes
- Democratiseer de architectuur, stel optimalisaties uit: Behoud de nuttige workflow en standaard interfaces op kleine schaal. Label elk mechanisme voor productieschaal eerlijk als geïmplementeerd, vereenvoudigd, uitgesteld of buiten scope.
- Bewijs, geen oordelen: Controles registreren metingen met dekking; het pass/fail-beleid behoort tot de consument tijdens de curatie. Kwaliteitstags routeren episodes; ze verwijderen nooit data.
- Standaardformaten bij elke grens: MCAP-episodes, Parquet-catalogi, Airflow DAG's. De code bestaat alleen waar het formaat een brug vereist of waar een valkuil echt niet voor de hand ligt.
- Jouw code blijft jouw code: Bestaande transformaties, controles en verrijkingen worden aangesloten via kleine adapters in plaats van dat ze herschreven moeten worden.
Systeemvereisten
- Python: $\ge$ 3.11
- Runtime: Docker (voor de pijplijn runtime;
app.test() heeft dit niet nodig), of een eigen Airflow-implementatie (Astronomer, MWAA, Cloud Composer, self-managed).
- Opslag: De eerste
hflow up downloadt ongeveer 2 GB aan container-images en bouwt de task venv.
- Cloud: Native
s3://, gs:// en Azure data roots gebruiken de optionele bucket-backend (uv sync --extra bucket); lokale paden hebben dit niet nodig.
- Video: Op Linux x8664/aarch64 downloadt de eerste video-operatie een checksum-geverifieerde, gepinde ffmpeg/ffprobe build in de user cache. Gebruik
HFLOWFFMPEG en HFLOW_FFPROBE om eigen binaries te gebruiken.
- Windows: Ondersteund via WSL2 (Airflow draait niet native op Windows).
Documentatie en Referenties
Documentatiebronnen
- Home: Start per taak, gevolgd door tutorials, how-to gidsen, referenties of uitleg.
- FAQ: Informatie over formaten, infrastructuur, schaal, projectscope en huidige releasestatus.
- Data Stack: Hoe HFlow past binnen de robotics data stack (MCAP, Airflow, Foxglove, Rerun, DuckDB, object storage, training formats).
- Voorbeelden: Uitvoerbare commando's, vereisten en links naar gidsen.
- Architectuur: Matrix van geïmplementeerde, vereenvoudigde, uitgestelde en buiten-scope functies.
- OpenAI Vision: Specifieke gids voor het aanroepen van OpenAI vision vanuit een stap.
- Bijdragen: Ontwikkelsetup, validatiecommando's en PR-verwachtingen.
Referenties
- Dyna Robotics: Training Dyna-2 at million-hour scale, repeatably.
- MCAP-specificatie en Python-libraries (Foxglove).
foxglove.CompressedVideo schema: in-band H.264/H.265/VP9/AV1 video in MCAP.
- Apache Airflow.
- DuckDB.
- Foxglove en Rerun.
- FFmpeg.
- Pareto: Het robotics data curatieplatform van Hebbian Robotics.