ray-project/raydp
61.8
Adequate · 20 September 2026
10.4k
lines of production code
Scala
with Python, Java
1
measurement over time
What this system is
RayDP is a hybrid runtime system that enables Apache Spark workloads to execute on top of the Ray distributed computing framework. It provides a Python API for managing Spark clusters within Ray, including utilities for converting data between Spark DataFrames and Ray Datasets. The system also integrates machine learning libraries such as PyTorch, TensorFlow, and XGBoost, offering distributed training estimators that leverage Ray's execution engine.
How it got here
2020 — Initial release and ML estimator integration
9 changes.
This period established the RayDP 1.7.0 foundation, introducing the core Spark-on-Ray integration and comprehensive test suites. It expanded the library's capabilities by adding distributed machine learning support through TFEstimator for TensorFlow and TorchEstimator for PyTorch, alongside practical usage examples.
2021–2022 — Spark-on-Ray runtime and MPI support
8 changes.
This period focused on establishing the core runtime infrastructure for executing Spark workloads on Ray, including the introduction of version-agnostic shim abstractions and specific support for Spark 4.2. It also expanded the library's capabilities by adding MPI job execution support and simplifying deployment through new Docker, Kubernetes, and submission tooling.
2023–2026 — Spark 3.4/4.x support and logging improvements
5 changes.
This period focused on extending RayDP compatibility to Spark 3.4 and 4.x through new shim layers, alongside introducing a Java agent to improve log management and job ID tracking. It also added a new XGBoostEstimator API for distributed training and expanded test coverage for version parsing and actor state management.
Features
Add Docker support for RayDP on Kubernetes
Users can now build and deploy RayDP on Kubernetes using the provided Dockerfile, build script, and Kubernetes configuration. The Dockerfile installs Java 8 and RayDP on top of the Ray base image, while the included \legacy.yaml\ file provides a Kubernetes cluster configuration for running Ray head and worker nodes with appropriate resource limits and shared memory settings.
docker · high confidence
Add MPI job execution support on Ray
Introduces a new \raydp.mpi\ module that enables running MPI jobs on Ray clusters. This includes a \create\_mpi\_job\ API to launch jobs with support for OpenMPI, Intel MPI, and MPICH, along with the underlying worker and driver services for process coordination and function execution.
python/raydp/mpi · high confidence
Add NYC Taxi Fare Prediction examples with data processing and training scripts
The examples directory now includes a complete NYC Taxi Fare Prediction workflow, featuring a random dataset generator (\random\_nyctaxi.py\) and a preprocessing script (\data\_process.py\) that handles feature engineering and JDK 17 compatibility options. Users can now run distributed training using PyTorch via \TorchEstimator\ (\pytorch\_nyctaxi.py\), Horovod (\horovod\_nyctaxi.py\), or the Ray Train API (\raytrain\_nyctaxi.py\), alongside a DLRM implementation notebook and a submission helper script.
examples · high confidence
Add Spark 3.4.0–3.4.4 shim layer
This change introduces a new shim implementation for Spark versions 3.4.0 through 3.4.4, enabling RayDP to run on these specific Spark releases. The shim provides version-specific adapters for core operations, including DataFrame creation from Arrow batches, Arrow schema conversion, command-line argument handling via SparkSubmitUtils, and the instantiation of custom executor backends (RayCoarseGrainedExecutorBackend) through a dedicated factory. This allows the platform to correctly bridge RayDP's internal APIs with the Spark 3.4.x runtime environment.
core/shims/spark340 · high confidence
Add Spark 4.2 shim support
This change adds a new shim module for Apache Spark 4.2 (version 4.2.0) to enable RayDP compatibility with this specific Spark minor line. The implementation includes the necessary service provider configuration, a shim provider descriptor, and concrete shim classes that bridge RayDP internals to Spark 4.2 APIs. Specifically, it provides utilities for converting DataFrames to and from Arrow formats, creates a custom executor backend factory for Spark 4.2, and handles command-line submission arguments, allowing users to run RayDP workloads on Spark 4.2 clusters.
(repo-wide) · high confidence
Add support for Spark 4.0 and 4.1 shims
New shim implementations for Spark 4.0 and 4.1 have been added to the core/shims/spark400 and core/shims/spark410 directories. These changes introduce version-specific providers, SQL utilities, executor backend factories, and task context utilities, enabling RayDP to run on Spark 4.x patch releases.
core/shims/spark400, core/shims/spark410 · high confidence
Added MPI network communication layer with gRPC services
The \python/raydp/mpi/network\ module now includes the protocol definition and generated Python code for a gRPC-based communication layer between RayDP driver and worker processes. This adds the \DriverService\ (handling worker registration and function result reporting) and \WorkerService\ (handling function execution and shutdown) definitions, enabling the underlying MPI job execution infrastructure to coordinate tasks across distributed nodes.
python/raydp/mpi/network · high confidence
Initial project scaffolding and documentation
The repository is initialized with core project files including an Apache 2.0 LICENSE, a SECURITY.md policy pointing to the Intel Security Center, and a CONTRIBUTING.md guide for setting up the mixed Scala/Python development environment (requiring JDK 8/17, Python 3.10+, and Maven 3.6+). A build.sh script is added to compile the core module and Python wheel, alongside a README.md detailing the RayDP architecture, installation, and usage for running Spark on Ray and integrating with AI libraries.
(repo-wide) · high confidence
Initial release of RayDP 1.7.0.dev0 with Spark-on-Ray integration
This change introduces the RayDP library (version 1.7.0.dev0), providing a Python API to run Apache Spark on Ray. Users can now initialize a Spark cluster within a Ray environment using \raydp.init\_spark()\, which supports configuration for executors, memory, and placement groups. The release includes a context manager for session lifecycle management, an abstract \EstimatorInterface\ for machine learning workflows, and utilities for parsing memory sizes and handling Spark DataFrames. It also introduces version-aware configuration selection for Spark's logging system (log4j vs log4j2) based on the installed PySpark version.
python/raydp · high confidence
Introduce Spark-on-Ray runtime core components and configuration
This change adds the foundational Java and Scala classes required to run Spark workloads on Ray, including the RayAppMaster and RayDPExecutor actors, their creation utilities, and the SparkOnRayConfigs configuration constants. It registers the RayClusterManager as the external cluster manager via the META-INF service file and provides the AppMasterEntryPoint to bootstrap the Python-to-Java gateway. These components establish the runtime wiring for executor lifecycle, resource allocation, and shuffle service management within the Ray cluster.
core/raydp-main/src/main · high confidence
Introduce TFEstimator for distributed TensorFlow Keras training
Users can now train TensorFlow Keras models in a distributed manner using a scikit-learn-like API via the new TFEstimator class. This component leverages Ray Train's TensorflowTrainer to handle distributed execution, allowing users to specify model architectures, optimizers, losses, and metrics while automatically managing data sharding and worker resources.
python/raydp/tf · high confidence
New XGBoostEstimator API for distributed training
Users can now train XGBoost models using a new \XGBoostEstimator\ class in the \raydp.xgboost\ module. This estimator wraps Ray's \XGBoostTrainer\ to provide a unified interface for distributed training, supporting both direct Ray Dataset inputs and Spark DataFrames via \fit\_on\_spark\. It allows configuration of worker counts, per-worker resources, data shuffling, and checkpointing, and exposes the trained model via \get\_model\.
python/raydp/xgboost · high confidence
New raydp-submit script for launching Spark on Ray
A new \bin/raydp-submit\ shell script has been added to simplify submitting Spark applications to a Ray cluster. This script automatically detects \SPARK\_HOME\, \RAY\_HOME\, and the driver node IP, and injects necessary JVM arguments (including JDK 17+ module access flags) and Spark configurations (such as the Ray driver host and log4j settings) before invoking \SparkSubmit\.
bin · high confidence
New tutorials for PyTorch and Ray Train integration
Added two new Colab-compatible Jupyter notebooks demonstrating how to use RayDP with PyTorch and Ray Train. The PyTorch example (\pytorch\_example.ipynb\) shows an end-to-end pipeline using Spark for data preprocessing and a distributed estimator for training, including a fix to add the \torchmetrics\ dependency. The Ray Train example (\raytrain\_example.ipynb\) illustrates using the new Ray Train API for distributed training and evaluation. Both tutorials use the healthcare stroke prediction dataset (\healthcare-dataset-stroke-data.csv\) to predict stroke risk based on patient data.
tutorials · high confidence
Architecture
Introduce common Spark shim abstractions for version-agnostic integration
The core/shims/common module now provides the foundational traits and loader logic that allow the library to adapt to different Spark versions at runtime. This includes a new SparkShimLoader that dynamically resolves the correct shim provider based on the running Spark version, a SparkShimProvider interface (with a helper for minor-line matching), and core traits like SparkShims, CommandLineUtilsBridge, and RayDPExecutorBackendFactory. These abstractions decouple version-specific implementation details from the core codebase, enabling support for multiple Spark versions (including 3.x and 4.x) without requiring separate code paths in the main application logic.
core/shims/common · high confidence
Behavioural changes
Introduce TorchEstimator with Ray Train integration and Intel oneCCL support
The \python/raydp/torch\ module now provides a \TorchEstimator\ that leverages the Ray Train API for distributed PyTorch training, replacing previous internal implementations. This change introduces a new \TorchConfig\ that enables the Intel oneCCL backend for distributed process group initialization, requiring \intel\_extension\_for\_pytorch\ and \torch-ccl\. Users can now specify evaluation metrics via \torchmetrics\, configure placement strategies for Ray actors, and utilize the \TorchMLDataset\ wrapper for efficient data loading from Ray's MLDataset.
python/raydp/torch · high confidence
Introduce recoverable Spark-to-Ray Dataset conversion and SparkCluster configuration
The \python/raydp/spark\ module now exposes a recoverable pipeline for converting Spark DataFrames to Ray Datasets via \from\_spark\_recoverable\ and \spark\_dataframe\_to\_ray\_dataset\, which caches Arrow bytes in Spark and builds Ray Dataset blocks using Ray task lineage to survive executor loss. The \SparkCluster\ class has been updated to automatically set the Spark driver host to the Ray node IP address and to support custom Spark locations by picking up the \$SPARK\_HOME\ environment variable.
python/raydp/spark · high confidence
Introduces Java agent to redirect SLF4J logs and prepend job IDs to worker logs
Adds a new Java agent (Agent.java) and SLF4J binding (StaticLoggerBinder.java) to the RayDP core. The agent intercepts system output/error streams during startup to redirect SLF4J binding warnings and logs into a dedicated file (slf4j-\<pid\>.log) instead of cluttering the Spark shell console, then restores the original streams. Additionally, for Ray versions 2.3.0–2.3.1, it prepends the job ID in the format ':job\_id:\<jobid\>' to the beginning of the java-worker log file to comply with Ray's log parsing requirements, using explicit UTF-8 encoding for file paths and content.
core/agent · high confidence
Test coverage
Added test suite for RayDP core, estimators, and utilities; Added tests for Spark version parsing and application actor-slot state management.
Dependencies
Python package initialization and dependency configuration
The Python package is initialized with a new MANIFEST.in and pylintrc configuration. The setup.py defines the package metadata, specifying a base version of 1.7.0 and enforcing Python 3.10+ compatibility. It establishes strict dependency constraints, requiring Ray \>= 2.37.0, PySpark \>= 4.0.0, PyArrow \>= 4.0.1, and Protobuf \> 3.19.5, while also including optional TensorFlow support.
python · high confidence
RayDP 1.7.0 dependency and shim updates
This change updates the RayDP build to version 1.7.0-SNAPSHOT, raising the default Spark target to 4.0.0 and adding dedicated shim modules for Spark 4.0.0, 4.1.0, and 4.2. It upgrades the Ray dependency to 2.1.0, bumps Jackson to 2.18.2, Rhino to 1.7.14.1, and commons-lang3 to 3.18.0, while introducing a new core/agent module with log4j 2.25.3 and slf4j 1.7.32.
(dependencies) · high confidence
Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.
How this codebase got here
Baseline
- First survey — no prior run to compare against. CAI 62.
Lenses
- Code Health 91
- Architecture 100
- Maturity 64
- Readiness 53
- Security 62
Changes since last survey
- 269 commits — 216 feature/other, 53 fixes
By area
- python/raydp — 88 commits
- (root) — 33 commits
- .github/workflows — 30 commits
- core/pom.xml — 20 commits
- core/src — 20 commits
- core/shims — 17 commits
- core/raydp-main — 15 commits
- python/setup.py — 8 commits
- core/agent — 7 commits
- doc/spark_on_ray.md — 6 commits
- python/spark_on_ray — 5 commits
- bin/raydp-submit — 3 commits
- docker/Dockerfile — 3 commits
- python/spark — 2 commits
- .github/ISSUE_TEMPLATE — 1 commit
- .github/pull_request_template.md — 1 commit
- core/javastyle.xml — 1 commit
- doc/mpi.md — 1 commit
- examples/NYC_Taxi.ipynb — 1 commit
- examples/NYC_Taxi_Pytorch.ipynb — 1 commit
Notable commits
- fix: A quick fix for issue #94 (#96)
- fix: Add and fixes UT (#41)
- fix: Fix CI (#252)
- fix: Fix CI with latest ray (#403)
- fix: Fix CVE snappy issue (#384)
- fix: Fix Jackson and Scala vulnerabilities (#373)
- fix: Fix RayDPDriverAgent Rpc bind (#467)
- fix: Fix Spark 3.5.x support + support Spark 3.4.3, deprecate python 3.7+support python 3.10 (#411)
- fix: Fix Spark download Url in tests (#255)
- fix: Fix TorchEstimator timeout when fit on large dataset (#20)
- fix: Fix connection error in CI (#224)
- fix: Fix document since PR#313 (#329)
- fix: Fix executor class path issue (#145)
- fix: Fix maven dependency and prepare for release (#356)
- fix: Fix the bug when use BlockSetSampler in TorchTrainer evaluation (#21)
- fix: Fix the pyspark driver log level (#269)
- fix: Fix unbound error in _SparkContext for int memory (#355)
- fix: Fix vulnerabilities listed in issue #117 (#118)
- fix: Fix: Ray nightly now uses its own schema type (#351)
- fix: Fixes multiple SLF4J bindings warning (#93)
- …and 249 more
Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.
Survey your own repository
ray-project/raydp was measured the same way every project in this corpus was: the same rubric, at a pinned commit, with the result published in full. Point a surveyor at a repository you know and see whether you agree with it.
About this page
- The score is its most recent published measurement, taken on 20 September 2026 at a pinned commit. It is not a live figure and does not change until the project is measured again.
- Measured at commit f1868429df7f46100127b5e259a4ca105cb654fa — the exact code this score is about.
- Scored under rubric-2026.09.15 — the same rubric and the same method as every other entry in this index.
- Measured by watchdog.canine.dev using codehealth-analyzer preprod-b51f968c9b10.