Add working CI tests to validate TPCH across a matrix of python / ray versions (#74)

5 files changed
tree: 0eea8e9597aa3b09cc969caf06d21b3650c5ea94
  1. .cargo/
  2. .github/
  3. datafusion_ray/
  4. dev/
  5. docs/
  6. examples/
  7. src/
  8. testdata/
  9. tests/
  10. tpch/
  11. .asf.yaml
  12. .dockerignore
  13. .gitignore
  14. .pre-commit-config.yaml
  15. build.rs
  16. Cargo.lock
  17. Cargo.toml
  18. LICENSE
  19. NOTICE
  20. pyproject.toml
  21. README.md
README.md

Python Tests

DataFusion on Ray

Overview

DataFusion Ray is a distributed execution framework that enables DataFusion DataFrame and SQL queries to run on a Ray cluster. This integration allows users to leverage Ray's dynamic scheduling capabilities while executing queries in a distributed fashion.

Execution Modes

DataFusion Ray supports two execution modes:

Streaming Execution

This mode mimics the default execution strategy of DataFusion. Each operator in the query plan starts executing as soon as its inputs are available, leading to a more pipelined execution model.

Batch Execution

Note: Batch Execution is not implemented yet. Tracking issue: https://github.com/apache/datafusion-ray/issues/69

In this mode, execution follows a staged model similar to Apache Spark. Each query stage runs to completion, producing intermediate shuffle files that are persisted and used as input for the next stage.

Getting Started

See the contributor guide for instructions on building DataFusion Ray.

Once installed, you can run queries using DataFusion's familiar API while leveraging the distributed execution capabilities of Ray.

import ray
from datafusion_ray import DFRayContext, df_ray_runtime_env

ray.init(runtime_env=df_ray_runtime_env)
session = DFRayContext()
df = session.sql("SELECT * FROM my_table WHERE value > 100")
df.show()

Contributing

Contributions are welcome! Please open an issue or submit a pull request if you would like to contribute. See the contributor guide for more information.

License

DataFusion Ray is licensed under Apache 2.0.