EXTRACT-TASKA: A Distributed Workflow Orchestrator for Radio Astronomy Data Processing
Abstract
We present the distributed workflow orchestrator framework developed within the EXTRACT project, for the TASKA (Transients Astrophysics with and SKA pathfinder). The orchestrator is built on existing software developed and maintained by the partner of the EXTRACT project. We demonstrate that a simple workflow for radio astronomy interferometric imaging can be run with this framework, using various cloud infrastructures. It has been tested with a private cloud provider (OVH), as well as on academic cloud infrastructures (EGI/CESNET, EOSC EU Node, ObsParis local OKD cluster). The demonstration can run from a jupyter notebook (e.g., on the EOSC EU Node), and the processing is run remotely on the selected cloud infrastructure.
Full text
The$EXTRACT$Project$has$received$funding$from$the$European$Union’s$ Horizon$Europe$programme$under$grant$agreement$number$101093110. TASKA Transient$Astrophysics$ with$a$Square$Kilometre$Array$pathfinder Baptiste Cecconi, Stéphane Aicardi and the EXTRACT collaboration extract-project.eu
2 EU project https://extract-project.eu/ EXTRACT aims to create a data-mining software platform for extreme data ! across the compute continuum e.g. NenuFAR SKA ?
Data Infrastructure Object Storage Data Catalog Semantic Data Staging Data Mining Framework Machine Learning Workflow Description Data-driven Orchestrator Big Data Scheduling Deployment Monitoring Edge Frameworks Cloud Frameworks HPC Frameworks Interoperability Abstraction Layer Compute Continuum Data security Cyber security Full-stack Security Layer ML model security 3
EXTRACT)Use-Cases,)sharing)the)same)platform Transient Astrophysics with the Square Kilometre Array pathfinder (TASKA) NenuFAR generating high-volume and highvelocity data TASKA Personalised Evacuation Route (PER) in the City of Venice based on an Urban Digital Twin and an AI engine PER 4
NA~2000 antennes NenuFAR Pathfinder de SKA (LOW) , Infrastructure de recherche New extension in Nançay Upgrading loFAR F = 20-80 MHz Fonctionnement en mode réseau phasé et interféromètre 5
Extreme)data)with)NenuFAR ➢~2000)antennas)(96$mini-arrays$of$19$ antennas$+$6$remotes$mini-arrays)) ➢Beamform)data:$$ raw$data$rate$=$1.2$GB/s$(~35$PB/yr)$ ➢Imaging)data:$$ raw$data$rate$=$8.6$TB/hr$(~74$PB/yr)$ ➢Local)data)storage:$3.5$PB) ➢Science)teams)are)reducing)data) down$to$about$1$PB/yr$ (>1/100$of$raw$data$rate))) ➢Reduced)data)transferred)to) distributed)datacenter)) (currently:$3$PB$in$Meudon$and$5$PB$ in$Orléans) 6 Compact core (400 m) : 96 MA x 2 polar. (or 96 antennas x 2 polar.) Distant MA (≤3 km) : 6 MA x 2 polar. antenna beam MA beam (analog-phased) : 8°-64° • 3 Gbit/s link to LOFAR correlator • Channelization down to <1 kHz • Correlation with signals from LBA stations (~1 s) • ADC 200 MHz (x 96 x 2) • PFB → 200 kHz SB [+TBB] • Phasing+Summation per SB = beamlet (244 simultaneous ~ 48 MHz total bandwidth) • SST, BST, XST → Calibration LOFAR FR606 backend • ADC 200 MHz (x 96 x 2) • PFB → 200 kHz SB [+TBB] • Phasing+Summation per SB = beamlet (768 simultaneous ~ 150 MHz total bandwidth) • Incoherent summation • SST, BST, XST → Calibration LaNewBa • Channelization down to <1 kHz • MA coherent sum (by 1, 2, 4) • MA (or MA sum) correlation (~1 s) Local correlator UnDySPuTeD • Channelization down to δf<1 kHz • Dedispersion (predefined, parametric) • Integration over δt ≥5 µs • ADC 200 MHz (x 6 x 2, at MA foot) • PFB → 200 kHz subbands NenuFAR Standalone Radio Imager NenuFAR Standalone Beamformer LOFAR Super Station NenuFAR analog signals MA = 19 antennas mini-array SB = ~200 kHz sub-band PFB = polyphase filter bank TBB = Transient Buffer Boards SST/BST/XST = Subband/Beamlet/Crosslet STatistics analog signals High-resolution imaging (~0.1" , IQUV) [+TBB] 40 kB /s /channel 9.6 MB /s /SB • Dynamic spectra (δt x δf) • Time series (δt) (IQUV, multibeam) [+TBB] 16/δt bytes /s /channel NenuFAR-Backend-0 (= UnDySPuTeD offline on <<768 SB) • Channelization down to δf<1 kHz • Dedispersion (predefined, parametric) • Integration over δt ≥5 µs 150 MB /s /SB • Instantaneous (~1 s) coarse imaging (~1°) • Slow (6-8h) multi-λ mid-resolution (8’) imaging → confusion ~ thermal noise (IQUV) 1.6 MB /s /SB 40 kB /s /channel NenuFAR : 3 instruments in 1 signal path/data flow/receivers/data products analog signals 1.6 MB /s /SB 1.6 MB /s /SB ➢Processing)of)imaging)data)on) distributed)datacenter:$ data$can’t$be$moved$easily,$need$to$ process$where$the$data$is$located
“Edge”$=$Nançay$facility$(NenuFAR$backend$+$Nançay$Data$Centre)$ “Cloud”$=$“Datalake”$(NenuFAR$Data$Centre$+$partners) NenuFAR)digital)infrastructure 7
➢Use)Case)A:$Early$detection$and$selective$resolution$data$recording$(space$optimality)$ ➢Use)Case)C:$Workflow$orchestration$of$interferometric$data$processing$with$a$focus$on$ improving$the$processing$speed,$accuracy$and$automation$on$large$datasets$ ➢Use)Case)D:$Prototype$development$for$“dynamic”$imaging$of$the$variable$Universe$ ➢Use)Case)E:$Advanced$data$reduction$workflows$for$multi-dimensional$real-time$ analysis$and$inference$(joining$A$and$C$together) TASKA)Use)Case)Overview) 8 (DL transient imaging)
Starting dataset: Visibilities (Measurement Sets (MS) Format) time/freq (Flagging, Rebining, Calibration, Imaging) Final product: Image cubes Data transfer data processing Optimal dataset distribution ? (Multiple sites) (Multiple tools) from combination of known analytic “bricks” Flexible reduction? Analytics Multiple tools (scientific quality, fidelity) Use-Case)C:)Workflow)orchestration)for)radiointerferometry 9
Implementing)the)chaining)of)tasks)and)decision)making)&)replay)capabilities x x √ TASKA)-)MVP)“Automated”)Workflow On-going development on EOSC, EGI/CESNET, OVH, (Soon Nançay/Obs) 16 x (provenance needed)
EXTRACT)-)TASKA)-)Summary 17 TASKA-C) ●We$have$developed$a$framework)for)distributed)data)computing$on$cloud$clusters$ ● Currently$validating$ $-$unsupervised/automated$workflow$$ $-$running$a$step$on$an$HPC$resource$ $-$running$on$a$multi-cluster$scale$(data$distributed$in$several$data$centers)$$$ ● Application$on$NenuFAR$(SKA$pathfinder)$ ● Clear$huge$potential$for$SRCNet$ Other$work$not$addressed$in$this$presentation:$$$ TASKA-A) ● Real$time$detection$(possibly$with$AI)$on$high$resolution$data$stream$(dynamic$spectra):$ implemented$on$NenuFAR$beamformer$backend$ TASKA-D) ●On$going$work$on$new$imager$for$dynamical$sources$with$AI-based$video$reconstruction$$
Julien Girard Astro Radio - 14/11/2024
Data)Infrastructure) •Object$Storage$$ •Data$Staging$Engine$ •Data$Catalog$ •Semantic$engine$ Data)Mining)Framework) •Complex$workflows$ description$ •Serverless$approach$ •Support$to$task-$and$databased$parallelism Data Infrastructure Object Storage Data Catalog Semantic Data Staging Data Mining Framework Machine Learning Workflow Description Data-driven Orchestrator Big Data Scheduling Deployment Monitoring Edge Frameworks Cloud Frameworks HPC Frameworks Interoperability Abstraction Layer Compute Continuum Data security Cyber security Full-stack Security Layer ML model security Components S3/InfluxDB Dataplug Nuvla Virtuoso COMPSs Lithops Pytorch Kserve Components 19
Components Data-driven)Orchestrator)) •Select$computing$ resources$for$workflows$ based$on$monitoring$ Compute)Continuum) •Unified$computing$ abstraction$layer$based$ on$containers$ •Programming$ paradigms$optimized$for$ edge,$cloud$and$HPC Data Infrastructure Object Storage Data Catalog Semantic Data Staging Data Mining Framework Machine Learning Workflow Description Data-driven Orchestrator Big Data Scheduling Deployment Monitoring Edge Frameworks Cloud Frameworks HPC Frameworks Interoperability Abstraction Layer Compute Continuum Data security Cyber security Full-stack Security Layer ML model security Components COMPSs Prometheus Kubernetes CUDA SkyStore 20
Components Cybersecurity)Capabilities) •Data$protection,$privacy$ and$confidentiality$ •AI)models)protection$ •Authenticity$and$ security$for$computing$ nodes$$ •Trustworthiness$and$ verifiability$of$routines$ and$libraries Data Infrastructure Object Storage Data Catalog Semantic Data Staging Data Mining Framework Machine Learning Workflow Description Data-driven Orchestrator Big Data Scheduling Deployment Monitoring Edge Frameworks Cloud Frameworks HPC Frameworks Interoperability Abstraction Layer Compute Continuum Data security Cyber security Full-stack Security Layer ML model security Components Trivy MultiParty$Computing Homomorphic$ Encryption 21
22 Data)Catalog:)Nuvla)data)management)(DM) Elements:) •S3%infrastructure%service% •endpoint$and$credentials$ •data%object% •reference$to$S3$object$on$S3%infrastructure%service$ •operations:$create,$obtain$S3$pre-signed$URL$for$$ upload,$download,$and$delete$ •data%record% •metadata$about$the$data%object$ •contains$reference$to$the$data%object$it$describes$ •data%set% •query$against$data%records$ Third-party)app)integration)with)DM)workflow:) •Nuvla$API$(JSON$over$HTTPS)$for$data$management$ •Nuvla$API$python-library$with$data$management$examples$ •Notifications$to$MQTT$on$data$objects$creation https://nuvla.io/ui/
SkyStore •Joint$project$UC$Berkeley$-$IBM$ •Provides$a$virtual)global)object)store$namespace$across$clouds$/$premises$ •Each$client$connects$to$local$S3-proxy,$which$is$connected$to$a$nearby$S3)cluster$and$to$a$central$location) DB)server •Data$replication$and$consistency$controlled$via$policies$ •Specifically,$remote$objects$can$be$automatically$cached$using$closer/local$object$storage$ •SkyStore$work$in$Phase$2$of$EXTRACT:$ •Matured$base$prototype$ •Joint$paper$(re-submission)$ (Slide from Erez Hadad, IBM) 23
Dataplug)dynamic)data)staging •Dataplug:)extensible$framework$that$implements$on-the-fly$data$ partitioning$ •Hide$complexities$of$pre-processing$and$partitioning$unstructured$ scientific$data$ •Data-driven$and$dynamic,$efficient$parallel$access$to$data$ •Generate$arbitrary$data$partitions$without$modifying$existing$data$ •Extensible$to$multiple$data$formats$ •KPI)1.1:$faster$partitioning$(up$to$65.6%$less$pre-processing$time,$ and$3.7x$in$fetching$partitions)$and$an$important$reduction$of$data$ transfers$in$staging 24
Throughput per step with different VCPU configurations. Data)throughput)analysis Data volumes per step (input-output 25