scieee AI-readable full text Open interactive document viewer

Data centric workflows: The Maestro Middleware in the Destination Earth Twin Engine

Haus, Utz-Uwe

Full text

1 Utz-Uwe Haus, Head of HPE EMEA Research Lab 2025 -09 -02, REX -IO Workshop @ CLUSTER25, Edinburgh Data centric workflows: The Maestro Middleware in the Destination Earth Twin Engine 01 Maestro background and fundamentals 02 Destination Earth Climate Twins 03 Emergency Checkpointing 04 Maestro v0.5 feature update 05 Fault tolerance design (WIP) 06 Raw performance Agenda 2 Maestro https://gitlab.com/maestro -data/maestro -core 3 2020 motivating use case: Operational Weather Prediction Workflow Data acquisition Initial construction of the atmosphere High -resolution forecast Ensemble forecast Products generation (57 millions/day) Fields DB ~150 TB/day R/W Archive Current data movement Today’s bottleneck •Data movement between forecast stages and product generation •Archiving via I/O aggregator nodes into PFS •Each product generation job is reading from PFS Vision •Speed up data -movement for the Pgen step (the 42 TiB) •Exploit multiple storage technologies •More flexible dependencies Time critical: 1h 4 Application coupling 5Confidential | Authorized `````` `` Maestro Maestro -enabled Traditional C entral data repository (PFS, Database) or tightly integrated coupling framework (MPMD or split Comm) Data object ‘Marketplace’, peer -to-peer data transfer, cross -process 6 Overview/Architecture: Maestro in a nutshell APPLICATION 1 APPLICATION 3 APPLICATION 2 0010110100110 Core Data Object API: declare, offer/withdraw or require/demand, dispose CDO offer require+demand give require+demand Maestro Data Management CDO POOL Scope Object Maestro Data Transformation Unified Memory-storage API: mamba library Maestro System Model Sys Cost Mem Mem Mem Mem Mem CDO Resources DRAM HBM NVRAM SSD PFS CDO CDO CDO CDO CDO CDO GPU FPGA CDO CDO (Core Data Object) It is at the heart of Maestro’s design and is used to encapsulate data and metadata. Supports dependencies. OFFER+WITHDRAW Applications OFFER CDOs to the management pool. Maestro manages the data, until WITHDRAW occurs. REQUIRE+DEMAND When an application REQUIREs a CDO, Maestro makes data available. At DEMAND is hands over resources containing the data and relinquishes all control it. SCOPE OBJECT Captures information about scope, size, access relations and schedules of the data to enable efficient movement and/or transformation MAESTRO SYSTEM MODEL Computes the cost of moving, transforming or copying data a CDO SYS Interface to every memory level, enabling core functionality of that memory via mamba library. Scope Objec t Sys cdo = mstro_cdo_declare(“name”) —same name = same object mstro_cdo_attribute_add(cdo,key,val) —Important: size, layout, (distribution), data reference —Optional: user -defined attributes mstro_cdo_offer(cdo) —At this point all other workflow participants can access cdo mstro_cdo_withdraw(cdo) —may block, async variant available mstro_cdo_dispose(cdo) Producer side cdo = mstro_cdo_declare(“name”) —same name = same object mstro_cdo_attribute_add(cdo,key,val) —Important: size, layout, (distribution), data reference —Optional: user -defined attributes mstro_cdo_require(cdo) —At this point reference to a suitable source for CDO will be established mstro_cdo_demand(cdo) —may block, async variant available mstro_cdo_dispose(cdo) Consumer side 7 Low intrusiveness •Batch up CDOs in a CDO Group (for batched OFFERs) Alternatives •Subscribe to pool events (like offer, require, withdraw) and act on them •Create CDO Group based on attributes (think: SQL SELECT) and iterate on them Alternatives Data - and Memory -aware workflows with maestro Applications coupling bypassing filesystem intermediary. Pool events allow the implementation of useful workflow components . No programming paradigm or memory management layer enforced, but utilizes and can take advantage of Mamba memory management library (https://gitlab.com/cerl/mamba ) 8 9 Transport: Peer -to-Peer The Pool Manager is just the messenger Climate DT 16 Workflow Data Notification Integration (WIP) https:// doi.org /10.1016/j.jemets.2025.100015 Emergency Checkpointing Application Work in progress 18 •KAUST ACC project •Quickly move data from application to safety when receiving a SIGTERM •Avoid losing progress/important data on abnormal termination of an application •Data needs to be pre -registered for backup •Leverage Maestro distributed CDOs to backup/restore distributed data (no serialization) •Maestro librarian component archives data to a storage and stage them on request. Emergency Checkpointing Applications 19 Maestro Pool Manager Wrapper API M × N distribution Register_backup() Restore() Maestro librarian Worker job events/commands SIGTERM A. Esposito, C. Haine and A. Mohammed, "Emergency Backup for Scientific Applications," 2022 IEEE/ACM Third International Symposium on Checkpointing for Supercomputing ( SuperCheck ), Dallas, TX, USA, 2022, pp. 1 -8, doi : 10.1109/SuperCheck56652.2022.00008. Local memory Storage Benchmarks 20 0 5 10 15 20 25 0 0.05 0.1 0.15 0.2 Bandwidth (GB/s) Total data size (GB) MEB Model OSU Slingshot 10 Slingshot 11 Model (bandwidth)= 𝑠𝑏 𝜆𝑏+𝑠 s = data size b = 12.5 GB/s (nominal bandwidth) 𝜆 = 45 𝑢𝑠 (𝑙𝑎𝑡𝑒𝑛𝑐𝑦) s = data size b = 25 GB/s (nominal bandwidth) 𝜆 = 7 𝑢𝑠 (𝑙𝑎𝑡𝑒𝑛𝑐𝑦 𝑖𝑛𝑐𝑢𝑟𝑟𝑒𝑑 𝑏𝑦 7 𝑚𝑠𝑔𝑠) Feature updates 21 Major changes: ✓Revamp maestro core threading model ✓Support maestro core thread pinning for operations, transport, and fabric threads ✓Memory and bug fixes ✓Read the docs documentation Maestro 0.4 (Sept 2023) Major changes: ✓Python interface for maestro -core ✓Retire CentOS in CI and use rocky instead ✓OFI threads build their own private endpoints, i.e. less locking between threads ✓OFI and operation threads are NUMA -aware ✓Update to mamba 0.2.1 Maestro v0.5 -27-gcbb6000e (August 2025) Status overview 22 720 commits, + 41,810 /-7,389 , 14 new tests and examples since v0.3 (End of Horizon Europe project) New features: ✓A visualization tool for maestro logs ✓Support for transport of large CDOs with fragmentation ✓CXI/Slingshot 11 support ✓Inline transport for small sized CDOs ✓ New features: ✓New librarian component as an example ✓Support for OFI multi -recv ✓OpenFAM transport method ✓GPU memory support New Features As of 0.5 -27 23 •SWIG Interface (maestro -py.i ) •Type mappings: Convert between Python and C data types •Exception handling: Transform C status codes into Python exceptions •Memory management: Handle allocation/deallocation across language boundaries •Object wrapping: Create Python objects for opaque C handles •Python Module ( maestro_core ) •Direct access to all public Maestro C API functions •Pythonic error handling through exceptions •Automatic memory management •Type -safe attribute handling Python interface 24 import maestro_core as M import numpy as np import _mamba as mamba def data_producer(pm_info): M.mstro_init("numpy_workflow", "producer", 0) M.mstro_pm_attach(pm_info) # Create a 2D numpy array data_array = np.array([[0,1,2,3], [4,5,6,7], [8,9,10,11], [12,13,14,15]], dtype='double') # Wrap numpy array in a Mamba array mamba_array = mamba.new_array(data_array) mamba_array.describe() # Create CDO and attach the Mamba array cdo = M.mstro_cdo_declare("scientific_data", None) M.mstro_cdo_attribute_set(cdo, M.MSTRO_ATTR_CORE_CDO_MAMBA_ARRAY, mamba_array) # Make data available M.mstro_cdo_seal(cdo) M.mstro_cdo_offer(cdo) print("Producer: Offered CDO with numpy array data") M.mstro_cdo_withdraw(cdo) M.mstro_cdo_dispose(cdo) M.mstro_finalize() •Thread Teams : •Generic infrastructure for managing groups of worker threads •Each thread has its own FIFO queue for work items •OFI Thread Team : Specialized thread team for OpenFabrics Interface operations •Pool Operations Thread Team: Specialized thread team for maestro pool operations •Pool Manager (PM) : Server -side component managing resource pools •Pool Client (PC) : Client -side component for pool operations •NUMA Awareness : Round -robin work distribution with NUMA -aware scheduling Multithreading model 25 OFI thread PM Pool OP thread Transport thread Pool Manager OFI thread PC Pool OP thread App thread Client Synthetic Performance benchmarks 32 Handling CDOs (declare/offer) Single processing thread Two processing threads A: Number of attributes S: Size of attributes in bytes #nodes: Clients talking to the PM Almost 120k CDO -ops/s with 2 threads of a single Maestro pool manager. Synthetic Performance benchmarks 33 Aggregated bandwidth: 300 GB/s at 110k msg/s processed by a single thread of the pool manager. Wire speed up to 8 consumer nodes. Model bandwidth based on #msg/s on the PM: (3 msg/CDO)×CDO_size Transport Bandwidth © 2025 Hewlett Packard Enterprise Development LP [email protected] Thank You https://gitlab.com/maestro -data/maestro -core backup Multithreading model 36 OP queue processing CQ processing Discover and build EPs OP queue fi_read fi_send fi_mr_reg … Mem pool msg envelope msg context send buffers recv buffers OFI thread ID/locality/handle OP maker •Specialized thread teams for OpenFabrics Interface operations •Each thread manages its own set of OFI endpoints •Operation queues for network operations (send, receive, RDMA) •Memory pools for efficient buffer management OFI Threads Push Pool OPs Multithreading model 37 •Handle server/client -side pool operations •Operation Queues : Per -thread FIFO queues for work items •Pool Operation Engine : •State machine for multi -step operations •Steps could block/resume, wait for state change/event, or OFI op completion •May skip steps when needed •Event Domain : Asynchronous completion handling Pool OP Threads Pool OP engine OP queue join leave declare … Handle declare stack notify_event handle_acks handle_join send_welcome … OP thread ID/locality/handle fetch process step (OFI OP?) next step? would block? complete? Handle leave stack notify_event handle_acks handle_join send_welcome … Handle join stack notify_event handle_acks handle_join send_welcome … Push OFI OPs Multithreading model 38 OP queue processing CQ processing Discover and build EPs OP queue fi_read fi_send fi_mr_reg … Mem pool msg envelope msg context send buffers recv buffers OFI thread ID/locality/handle OP maker notify_event handle_acks handle_join send_welcome … Handle join stack Pool OP engine OP queue join leave declare … Handle declare stack notify_event handle_acks handle_join send_welcome … OP thread ID/locality/handle fetch process step (OFI OP?) next step? would block? complete? Handle leave stack notify_event handle_acks handle_join send_welcome … Handle join stack notify_event handle_acks handle_join send_welcome … Push to same locality Push to same locality Pool OP engine OP queue join leave declare … Handle declare stack notify_event handle_acks handle_join send_welcome … OP thread ID/locality/handle fetch process step (OFI OP?) next step? would block? complete? Pool OP engine OP queue join leave declare … Handle declare stack notify_event handle_acks handle_join send_welcome … OP thread ID/locality/handle fetch process step (OFI OP?) next step? would block? complete? OFI thread team Pool OP thread team Handle leave stack notify_event handle_acks handle_join send_welcome … Handle join stack notify_event handle_acks handle_join send_welcome …