Skip to content

[SPARK-60101][SQL] Support native data sources implemented in Rust or C++ - #59313

Open
HyukjinKwon wants to merge 11 commits into
apache:masterfrom
HyukjinKwon:SPARK-60101
Open

HyukjinKwon wants to merge 11 commits into
apache:masterfrom
HyukjinKwon:SPARK-60101

Conversation

@HyukjinKwon

@HyukjinKwon HyukjinKwon commented Oct 9, 2026 •

Copy link
Copy Markdown
Member

What changes were proposed in this pull request?

This PR adds native data sources: data sources implemented in native code, such as Rust or C++, without any JVM code. A native data source is a shared library that implements a binary interface. Spark finds it automatically when it is installed like other native libraries, or in a package file added to a session, and loads it when a query uses it. Spark adapts it to Data Source V2, with batch and micro-batch streaming reads, batch and streaming writes, and predicate, column and limit pushdown. Data is exchanged only as columnar batches, through the Arrow C data interface.

Public API: a single class, org.apache.spark.sql.datasource.NativeBridge (@Evolving, since 4.4.0). A library implements its native methods as JNI functions, such as Java_org_apache_spark_sql_datasource_NativeBridge_createDataSource, and its Javadoc is the specification of the interface: the lifecycle, the handles, the ownership of the Arrow structs, errors (Java exceptions, reported as NATIVE_DATA_SOURCE_ERROR), and the JSON format of the pushed predicates, with the Spark SQL semantics that a predicate evaluated by the library must follow, for example for NaN. Spark loads each library into its own copy of the class, with a separate class loader, so any number of libraries exporting the same functions can be loaded side by side. Spark itself ships no native code for native data sources, besides the JNI library that the new arrow-c-data dependency bundles. The interface uses JNI rather than the FFM API, which is final only from Java 22 while Spark supports Java 17 and 21; the interface is versioned, so that a later version can be a plain C ABI called through FFM.

Operation Required functions Optional functions
All abiVersion, createDataSource, closeDataSource schema
Batch read createReader, partitions, read, closeReader pushPredicates, pushLimit, pruneColumns, serializeReader
Streaming read createStreamReader, initialOffset, latestOffset, streamPartitions, read, closeStreamReader serializeStreamReader, commitOffset
Batch write createWriter, commit, abort, closeWriter, createDataWriter, write, commitDataWriter, abortDataWriter serializeWriter
Streaming write createStreamWriter, commit, abort, closeWriter, createDataWriter, write, commitDataWriter, abortDataWriter serializeWriter

A library only exports the functions of the operations it supports. Spark reports a missing required function as DATA_SOURCE_*_NOT_SUPPORTED.

Finding native data sources: when no Java or Python data source has the requested name, Spark looks for a native data source with that name, in two places:

  • Installed libraries, automatically: like the Python data sources installed in the Python path, a library installed like other native libraries is found without any configuration, by its file name. The library of the data source rust_range is spark_datasource_rust_range: libspark_datasource_rust_range.so on Linux, libspark_datasource_rust_range.dylib on macOS and spark_datasource_rust_range.dll on Windows. Spark checks whether the file exists in the native library path, which consists of:

    1. the directories of java.library.path, where the JVM itself finds native libraries. They include LD_LIBRARY_PATH on Linux, DYLD_LIBRARY_PATH on macOS and PATH on Windows, and so spark.driver.extraLibraryPath and spark.executor.extraLibraryPath;
    2. the lib directory of each installation prefix in PATH: <prefix>/lib for each <prefix>/bin.

    So libraries are found where native libraries are usually installed: /usr/local/lib (make install, CMake, cargo-c, Homebrew on Intel macs), /opt/homebrew/lib, $CONDA_PREFIX/lib and $VIRTUAL_ENV/lib in an activated environment, and /usr/lib. A library that implements several data sources is installed under each of their names, for example with symbolic links. The executors load the library from the same path as the driver, or else find it in their own native library path.

  • Packages of a session: a native data source package is a zip file with the extension .sparkpkg, with a manifest, spark-native-datasource.json, that lists the ABI version, the names of the data sources, and the library for each platform such as linux-x86_64 or osx-aarch64. A session adds packages with spark.addArtifact (Scala and Python, classic and Spark Connect: .sparkpkg files are added as file artifacts), or with spark.sql.dataSource.native.paths. They take precedence over the installed libraries, and the executors get them through the artifacts of the session, except in local mode, where they read them in place: a configured package is copied there under a name with its checksum, so it can be replaced in place. Spark ignores, and logs, a package that it cannot read or whose ABI version it does not support, unless its manifest lists the requested data source and no other package or installed library provides it.

Short names: no JVM class is needed for each data source. After the Java data sources (DataSourceRegister) and the Python ones, DataSource.lookupDataSource looks the name up in the packages of the session and in the native library path, case-insensitively, and resolves it to a single provider class, NativeDataSourceV2. The caller then gives it the requested name, the way PythonDataSourceV2 already works: a small NamedTableProvider trait generalizes how PythonDataSourceV2 gets its short name, so the native provider shares the same call sites. The provider reuses the result of the lookup that found its class, instead of looking the name up again. A Python data source can still be registered with the name of a native data source, and takes precedence over it. Lookups only read manifests and check whether files exist; Spark loads a library when a query uses it. Native data sources work with spark.read/write, readStream/writeStream, and SQL tables (CREATE TABLE ... USING rust_range OPTIONS (...), SELECT and INSERT INTO).

Internals:

  • execution.datasources.v2.columnar: an internal columnar data source interface, adapted to Data Source V2. Writes convert rows to Arrow batches that end at spark.sql.execution.arrow.maxRecordsPerBatch rows or once they hold spark.sql.execution.arrow.maxBytesPerBatch bytes, whichever is reached first, as for Python UDFs. The writer of a write is created once the stages it depends on ran, right before its tasks start, so that each writer is committed or aborted. A micro-batch streaming write gets a new stream writer for each micro-batch, and the previous one of the query is closed when the next one is created.
  • execution.datasources.v2.ffi: the discovery of packages and installed libraries, package parsing, the distribution of packages to the executors, library loading, the implementation of the internal interface on top of NativeBridge with arrow-c-data, which reads the Arrow streams array by array to reject arrays with a non-zero offset (Arrow Java ignores offsets), and the JSON encoding of pushed predicates.

Other changes:

  • Configurations: spark.sql.dataSource.native.enabled (static, default true) and spark.sql.dataSource.native.paths.
  • Error conditions: NATIVE_DATA_SOURCE_ERROR, INVALID_NATIVE_DATA_SOURCE_LIBRARY (CANNOT_LOAD, UNSUPPORTED_ABI_VERSION), INVALID_NATIVE_DATA_SOURCE_PACKAGE (INVALID_MANIFEST, INVALID_LIBRARY, UNSUPPORTED_ABI_VERSION, UNSUPPORTED_PLATFORM), NATIVE_DATA_SOURCE_LIBRARY_NOT_FOUND, NATIVE_DATA_SOURCE_PACKAGE_CONFLICT and NATIVE_DATA_SOURCE_PACKAGE_NOT_FOUND.
  • New dependency: org.apache.arrow:arrow-c-data, with the same version as arrow-vector. It bundles its own JNI library for Linux, macOS and Windows.
  • Documentation: a new page, "Native Data Sources", with the examples below, linked from the Data Source V2 page.
  • Existing classes:
    • DataSource.lookupDataSource looks up the native data sources after the Python ones, and DataSourceRegistration skips them when it registers a Python data source.
    • ArtifactManager lists the native data source packages of a session, and forgets the cached packages read from the artifacts of a session when it cleans them up.
    • The classic DataFrameWriter builds its commands, and DataStreamWriter resolves its sink, with the session of the Dataset active, as the analysis of a query does: the lookups of the Python and native data sources use the active session.
    • PySpark addArtifacts adds .sparkpkg files to the session whatever the flags: in classic mode through the JVM session, and with Spark Connect as file artifacts.
    • CheckConnectJvmClientCompatibility excludes the new org.apache.spark.sql.datasource package, which the Spark Connect client does not have.

Examples

Both examples implement a read-only data source that returns the numbers [0, end) as a column id, in two partitions, with the option end (10 by default). They only export the 8 functions needed for batch reads.

Rust (with the jni crate and arrow-rs)

Cargo.toml:

[package]
name = "rust_range"
version = "0.1.0"
edition = "2021"

[lib]
name = "spark_datasource_rust_range"
crate-type = ["cdylib"]

[dependencies]
arrow-array = { version = "58", features = ["ffi"] }
arrow-schema = { version = "58", features = ["ffi"] }
jni = "0.21"

src/lib.rs:

//! A native data source for Apache Spark, written in Rust. It reads the numbers [0, end) as a
//! column `id`, in two partitions. The option `end` defaults to 10.

use std::panic::{catch_unwind, AssertUnwindSafe};
use std::sync::Arc;

use arrow_array::ffi_stream::FFI_ArrowArrayStream;
use arrow_array::{Int64Array, RecordBatch, RecordBatchIterator};
use arrow_schema::ffi::FFI_ArrowSchema;
use arrow_schema::{DataType, Field, Schema, SchemaRef};
use jni::objects::{JByteArray, JClass, JObject, JObjectArray, JString};
use jni::sys::{jint, jlong, jobjectArray};
use jni::JNIEnv;

type Result<T> = std::result::Result<T, Box<dyn std::error::Error>>;

struct RangeSource {
    end: i64,
}

struct RangeReader {
    end: i64,
}

fn schema() -> SchemaRef {
    Arc::new(Schema::new(vec![Field::new("id", DataType::Int64, false)]))
}

/// Runs `f`, and reports its error or panic as a Java exception: neither may cross into the JVM.
fn call<T>(env: &mut JNIEnv, default: T, f: impl FnOnce(&mut JNIEnv) -> Result<T>) -> T {
    let message = match catch_unwind(AssertUnwindSafe(|| f(env))) {
        Ok(Ok(value)) => return value,
        Ok(Err(e)) => e.to_string(),
        Err(_) => "the data source panicked".to_string(),
    };
    if !env.exception_check().unwrap_or(true) {
        let _ = env.throw_new("java/lang/RuntimeException", message);
    }
    default
}

fn strings(env: &mut JNIEnv, array: &JObjectArray) -> Result<Vec<String>> {
    let mut result = Vec::new();
    for i in 0..env.get_array_length(array)? {
        let value = JString::from(env.get_object_array_element(array, i)?);
        result.push(env.get_string(&value)?.into());
    }
    Ok(result)
}

#[no_mangle]
pub extern "system" fn Java_org_apache_spark_sql_datasource_NativeBridge_abiVersion(
    _env: JNIEnv,
    _class: JClass,
) -> jint {
    1
}

#[no_mangle]
pub extern "system" fn Java_org_apache_spark_sql_datasource_NativeBridge_createDataSource(
    mut env: JNIEnv,
    _class: JClass,
    _name: JString,
    keys: JObjectArray,
    values: JObjectArray,
) -> jlong {
    call(&mut env, 0, |env| {
        let mut source = RangeSource { end: 10 };
        for (key, value) in strings(env, &keys)?.into_iter().zip(strings(env, &values)?) {
            if key == "end" {
                source.end = value.parse()?;
            }
        }
        Ok(Box::into_raw(Box::new(source)) as jlong)
    })
}

#[no_mangle]
pub extern "system" fn Java_org_apache_spark_sql_datasource_NativeBridge_schema(
    mut env: JNIEnv,
    _class: JClass,
    _source: jlong,
    schema_address: jlong,
) {
    call(&mut env, (), |_| {
        let schema = FFI_ArrowSchema::try_from(schema().as_ref())?;
        unsafe { std::ptr::write(schema_address as *mut FFI_ArrowSchema, schema) };
        Ok(())
    })
}

#[no_mangle]
pub extern "system" fn Java_org_apache_spark_sql_datasource_NativeBridge_createReader(
    mut env: JNIEnv,
    _class: JClass,
    source: jlong,
    _schema_address: jlong,
) -> jlong {
    call(&mut env, 0, |_| {
        let source = unsafe { &*(source as *const RangeSource) };
        Ok(Box::into_raw(Box::new(RangeReader { end: source.end })) as jlong)
    })
}

#[no_mangle]
pub extern "system" fn Java_org_apache_spark_sql_datasource_NativeBridge_closeDataSource(
    _env: JNIEnv,
    _class: JClass,
    source: jlong,
) {
    drop(unsafe { Box::from_raw(source as *mut RangeSource) });
}

#[no_mangle]
pub extern "system" fn Java_org_apache_spark_sql_datasource_NativeBridge_partitions(
    mut env: JNIEnv,
    _class: JClass,
    reader: jlong,
) -> jobjectArray {
    call(&mut env, std::ptr::null_mut(), |env| {
        let end = unsafe { &*(reader as *const RangeReader) }.end;
        let bounds = [(0, end / 2), (end / 2, end)];
        let partitions = env.new_object_array(bounds.len() as i32, "[B", JObject::null())?;
        for (i, (start, end)) in bounds.into_iter().enumerate() {
            let bytes = [start.to_le_bytes(), end.to_le_bytes()].concat();
            let partition = env.byte_array_from_slice(&bytes)?;
            env.set_object_array_element(&partitions, i as i32, partition)?;
        }
        Ok(partitions.into_raw())
    })
}

#[no_mangle]
pub extern "system" fn Java_org_apache_spark_sql_datasource_NativeBridge_closeReader(
    _env: JNIEnv,
    _class: JClass,
    reader: jlong,
) {
    drop(unsafe { Box::from_raw(reader as *mut RangeReader) });
}

#[no_mangle]
pub extern "system" fn Java_org_apache_spark_sql_datasource_NativeBridge_read(
    mut env: JNIEnv,
    _class: JClass,
    _reader_state: JByteArray,
    partition: JByteArray,
    stream_address: jlong,
) {
    call(&mut env, (), |env| {
        let bytes = env.convert_byte_array(&partition)?;
        let start = i64::from_le_bytes(bytes[0..8].try_into()?);
        let end = i64::from_le_bytes(bytes[8..16].try_into()?);
        let ids = Int64Array::from_iter_values(start..end);
        let batch = RecordBatch::try_new(schema(), vec![Arc::new(ids)])?;
        let batches = RecordBatchIterator::new(vec![Ok(batch)], schema());
        let stream = FFI_ArrowArrayStream::new(Box::new(batches));
        unsafe { std::ptr::write(stream_address as *mut FFI_ArrowArrayStream, stream) };
        Ok(())
    })
}

Build and install, here on macOS:

cargo build --release
cp target/release/libspark_datasource_rust_range.dylib /usr/local/lib/

Or package, to add it to a session with spark.addArtifact:

mkdir -p package/osx-aarch64
cp target/release/libspark_datasource_rust_range.dylib package/osx-aarch64/
cat > package/spark-native-datasource.json <<EOF
{
  "abiVersion": 1,
  "dataSources": ["rust_range"],
  "libraries": {"osx-aarch64": "osx-aarch64/libspark_datasource_rust_range.dylib"}
}
EOF
(cd package && zip -r ../rust_range.sparkpkg .)
C++ (with Arrow C++)

cpp_range.cc:

// A native data source for Apache Spark, written in C++ with Arrow C++. It reads the numbers
// [0, end) as a column `id`, in two partitions. The option `end` defaults to 10.

#include <jni.h>

#include <cstdint>
#include <memory>
#include <stdexcept>
#include <string>
#include <type_traits>

#include <arrow/api.h>
#include <arrow/c/bridge.h>

#define SPARK_JNI(ret, method) \
  extern "C" JNIEXPORT ret JNICALL Java_org_apache_spark_sql_datasource_NativeBridge_##method

namespace {

struct RangeSource {
  int64_t end = 10;
};

struct RangeReader {
  int64_t end;
};

std::shared_ptr<arrow::Schema> Schema() {
  return arrow::schema({arrow::field("id", arrow::int64(), /*nullable=*/false)});
}

// Runs `f`, and reports its error as a Java exception: no C++ exception may cross into the JVM.
template <typename F>
auto Call(JNIEnv* env, F&& f) -> decltype(f()) {
  try {
    return f();
  } catch (const std::exception& e) {
    env->ThrowNew(env->FindClass("java/lang/RuntimeException"), e.what());
    if constexpr (!std::is_void_v<decltype(f())>) return {};
  }
}

void Check(const arrow::Status& status) {
  if (!status.ok()) throw std::runtime_error(status.ToString());
}

std::string ToString(JNIEnv* env, jobject value) {
  const char* chars = env->GetStringUTFChars(static_cast<jstring>(value), nullptr);
  std::string result(chars);
  env->ReleaseStringUTFChars(static_cast<jstring>(value), chars);
  return result;
}

}  // namespace

SPARK_JNI(jint, abiVersion)(JNIEnv*, jclass) { return 1; }

SPARK_JNI(jlong, createDataSource)(JNIEnv* env, jclass, jstring, jobjectArray keys,
                                   jobjectArray values) {
  return Call(env, [&] {
    auto source = std::make_unique<RangeSource>();
    for (jsize i = 0; i < env->GetArrayLength(keys); i++) {
      if (ToString(env, env->GetObjectArrayElement(keys, i)) == "end") {
        source->end = std::stoll(ToString(env, env->GetObjectArrayElement(values, i)));
      }
    }
    return reinterpret_cast<jlong>(source.release());
  });
}

SPARK_JNI(void, schema)(JNIEnv* env, jclass, jlong, jlong schema_address) {
  Call(env, [&] {
    Check(arrow::ExportSchema(*Schema(), reinterpret_cast<ArrowSchema*>(schema_address)));
  });
}

SPARK_JNI(jlong, createReader)(JNIEnv* env, jclass, jlong source, jlong) {
  return Call(env, [&] {
    auto* reader = new RangeReader{reinterpret_cast<RangeSource*>(source)->end};
    return reinterpret_cast<jlong>(reader);
  });
}

SPARK_JNI(void, closeDataSource)(JNIEnv*, jclass, jlong source) {
  delete reinterpret_cast<RangeSource*>(source);
}

SPARK_JNI(jobjectArray, partitions)(JNIEnv* env, jclass, jlong reader) {
  return Call(env, [&] {
    int64_t end = reinterpret_cast<RangeReader*>(reader)->end;
    int64_t bounds[2][2] = {{0, end / 2}, {end / 2, end}};
    jobjectArray partitions = env->NewObjectArray(2, env->FindClass("[B"), nullptr);
    for (jsize i = 0; i < 2; i++) {
      jbyteArray partition = env->NewByteArray(sizeof(bounds[i]));
      env->SetByteArrayRegion(partition, 0, sizeof(bounds[i]),
                              reinterpret_cast<const jbyte*>(bounds[i]));
      env->SetObjectArrayElement(partitions, i, partition);
    }
    return partitions;
  });
}

SPARK_JNI(void, closeReader)(JNIEnv*, jclass, jlong reader) {
  delete reinterpret_cast<RangeReader*>(reader);
}

SPARK_JNI(void, read)(JNIEnv* env, jclass, jbyteArray, jbyteArray partition,
                      jlong stream_address) {
  Call(env, [&] {
    int64_t bounds[2];
    env->GetByteArrayRegion(partition, 0, sizeof(bounds), reinterpret_cast<jbyte*>(bounds));
    arrow::Int64Builder ids;
    for (int64_t id = bounds[0]; id < bounds[1]; id++) Check(ids.Append(id));
    std::shared_ptr<arrow::Array> array;
    Check(ids.Finish(&array));
    auto batch = arrow::RecordBatch::Make(Schema(), array->length(), {array});
    auto batches = arrow::RecordBatchReader::Make({batch}, Schema()).ValueOrDie();
    Check(arrow::ExportRecordBatchReader(batches,
                                         reinterpret_cast<ArrowArrayStream*>(stream_address)));
  });
}

Build (Arrow C++ 24 requires C++20) and install, or package as above:

c++ -std=c++20 -O2 -shared -fPIC \
  -I"$JAVA_HOME/include" -I"$JAVA_HOME/include/darwin" \
  -I"$ARROW_HOME/include" -L"$ARROW_HOME/lib" -larrow \
  -o libspark_datasource_cpp_range.dylib cpp_range.cc
cp libspark_datasource_cpp_range.dylib /usr/local/lib/

Usage, in classic PySpark or with Spark Connect. Once installed, the data sources are found automatically, without spark.addArtifact or any configuration:

>>> spark.read.format("rust_range").option("end", 5).load().show()
+---+
| id|
+---+
|  0|
|  1|
|  2|
|  3|
|  4|
+---+
>>> spark.read.format("cpp_range").option("end", 5).load().filter("id >= 3").collect()
[Row(id=3), Row(id=4)]
>>> spark.sql("CREATE TABLE t USING rust_range OPTIONS (end 3)")
>>> spark.sql("SELECT * FROM t").collect()
[Row(id=0), Row(id=1), Row(id=2)]

Or added to a session as a package:

>>> spark.addArtifact("rust_range/rust_range.sparkpkg")
>>> spark.read.format("rust_range").option("end", 5).load().count()
5

Why are the changes needed?

Many data systems have their best client libraries in native code, such as arrow-rs, DataFusion, delta-rs, iceberg-rust and Lance in Rust, or C++ libraries. Using them from Spark today requires writing a JVM Data Source V2 connector with its own JNI glue. With this change, a native data source implements a stable binary interface once, is installed like any other native library or distributed as a single file, and works with every Spark API, including Spark Connect clients in any language, without any JVM code.

Does this PR introduce any user-facing change?

Yes. It adds native data sources, with the NativeBridge interface, the configurations and the error conditions above. Spark finds the native data source libraries installed in the native library path, and spark.addArtifact now accepts native data source packages (.sparkpkg files). Also, in classic mode, writing with a session that is not the active one now finds the data sources of that session, such as its Python data sources, because DataFrameWriter and DataStreamWriter resolve the data source with that session active.

How was this patch tested?

  • New NativeDataSourceSuite, with a dependency-free C++ test library (sql/core/src/test/resources/native-datasource/test_native_datasource.cc) compiled by the tests with the C++ compiler of the machine; the tests that load it are skipped if there is no compiler. It covers inferred and user-specified schemas, partitions and batches, the supported types, predicate, column and limit pushdown, batch writes (append and overwrite), aborted writes, streaming reads and writes over several micro-batches with offset commits, SQL tables with OPTIONS (SELECT and INSERT INTO), errors reported by native code, missing functions, several libraries exporting the same functions side by side, spark.addArtifact, packages in directories, package conflicts, the precedence of Java data sources, and invalid packages (manifest, ABI version, platform and library). It also has unit tests for the manifest, the platform names, the native library path and the JSON encoding of predicates.
  • NativeDataSourceSuite also covers installed libraries: found automatically in the native library path with case-insensitive names, a library installed under several names, writes, the precedence of the packages of a session, the lookup by name on executors (NATIVE_DATA_SOURCE_LIBRARY_NOT_FOUND), and invalid installed libraries (CANNOT_LOAD, and UNSUPPORTED_ABI_VERSION, with the same error when the library is used again).
  • New NativeDataSourceConnectSuite, with a Spark Connect client and server: spark.addArtifact of a package, reads, writes, SQL tables, the isolation of the packages of different sessions, and installed libraries.
  • New ColumnarDataSourceSuite, which tests the internal Data Source V2 adapter with a data source implemented on the JVM.
  • New PySpark tests for adding .sparkpkg artifacts in classic mode, whatever the flags and in lists with other files, and with Spark Connect.
  • Tests added during the review, in NativeDataSourceSuite unless noted:
    • Packages: invalid and unsupported packages are ignored unless no other package or installed library provides the data source; a configured package replaced in place; the copies distributed to the executors are named by checksum and are not packages of the session; the cached packages of a session are forgotten when its artifacts are cleaned up; configured directories that cannot be read; distributing a package without an active session.
    • Lookups: a single lookup for each resolution, which the next lookup on the thread clears; a Python data source registered with the name of an installed native one; batch and streaming writes with a session that is not the active one.
    • Writes: a failed commitDataWriter aborts the data writer; a failed commit aborts the write, then closes the writer; the test library records when a writer is closed, which checks that a batch write closes its writer, and that each micro-batch of a streaming write closes the writer of the previous one; ColumnarDataSourceSuite checks that each micro-batch creates its stream writer with the ID of the query, and that batches end at spark.sql.execution.arrow.maxBytesPerBatch.
    • Reads: predicates on columns with collated strings are not pushed down, read schemas with collations, and the scan described by the name of the data source; sliced arrays and decimal256 are rejected with the name of the column, maps whose entry fields are named keys and values (as arrow-rs names them) are read, and limit(...).offset(...) asks the library for the limit once.
    • Writes whose upstream stages fail with adaptive execution create no writer (also in ColumnarDataSourceSuite).
  • The Rust and C++ examples above were built from scratch with the commands above, and run on macOS (osx-aarch64) with classic PySpark and with Spark Connect, with the output shown above: installed in the lib directory of an installation prefix in PATH without any configuration (and not found without it), found through java.library.path, added with spark.addArtifact, and as SQL tables.
  • The C++ test library also builds without warnings with GCC 15 and libstdc++ (-Wall -Wextra), and NativeDataSourceSuite runs on Linux with g++ in CI.
  • Existing suites: ResolvedDataSourceSuite, ArtifactManagerSuite, PythonDataSourceSuite, PythonStreamingDataSourceSuite, DataStreamReaderWriterSuite, DataFrameReaderWriterSuite, DataFrameWriterV2Suite, DataSourceSuite, DDLSourceLoadSuite, DataSourceV2Suite, StreamingDataSourceV2Suite, SQLConfSuite, SparkThrowableSuite and SparkConfigBindingPolicySuite.

Was this patch authored or co-authored using generative AI tooling?

Generated-by: Claude Code (Claude Opus 5.5)

This pull request and its description were written by Isaac.

HyukjinKwon and others added 2 commits October 9, 2026 10:09
…(Rust, C++) data sources

This adds the columnar data source API, a simpler way to plug a data source into Spark than
implementing the Data Source V2 interfaces directly. It follows the Python Data Source API, but
data is exchanged only as columnar batches: readers return ColumnarBatches, and writers receive
Arrow-backed ColumnarBatches. Spark adapts a data source to Data Source V2, with batch and
micro-batch reads, batch and streaming writes, and predicate, column and limit pushdown.

A data source is implemented either on the JVM, with the interfaces of
org.apache.spark.sql.datasource and a DataSourceProvider, or in native code such as Rust or C++,
as a shared library that implements the native methods of NativeBridge as JNI functions and
exchanges data through the Arrow C data interface. Spark loads each library into its own copy of
NativeBridge, so libraries exporting the same functions can be loaded side by side, and Spark
itself ships no native code.

A native data source is distributed as a native data source package, a zip file with the
extension .sparkpkg that contains a manifest and the library for each platform. Packages are added
with spark.addArtifact, in Scala and Python, classic and Spark Connect, or found under
spark.sql.dataSource.native.paths; executors get them through the artifacts of the session.

The change also adds the spark.sql.dataSource.native.enabled and spark.sql.dataSource.native.paths
configurations, the native data source error conditions, the arrow-c-data dependency, and a
documentation page with examples in Scala, Rust and C++.

Co-authored-by: Isaac <no-reply@databricks.com>
Downstream data sources are only native, so the JVM interfaces of the columnar data source API
(DataSource, DataSourceReader, DataSourceWriter, DataSourceStreamReader, DataSourceStreamWriter
and DataSourceProvider) are not public anymore. They become the internal ColumnarDataSource
traits that the Data Source V2 adapter and the native data sources share, and the Javadoc of
NativeBridge now specifies the whole interface on its own. The documentation page is renamed to
"Native Data Sources".

Co-authored-by: Isaac <no-reply@databricks.com>
@HyukjinKwon HyukjinKwon changed the title [SPARK-60101][SQL] Add a columnar data source API for JVM and native (Rust, C++) data sources [SPARK-60101][SQL] Support native data sources implemented in Rust or C++ Oct 9, 2026
HyukjinKwon and others added 5 commits October 9, 2026 10:41
A table created with CREATE TABLE ... USING ... OPTIONS (...) passes its options to
TableProvider.getTable, while its scans and writes only get the options of their query. The
columnar table now merges the options of the table under the options of each scan and write, so
that SELECT and INSERT INTO on such a table use them.

Co-authored-by: Isaac <no-reply@databricks.com>
…cally

Like the Python data sources installed in the Python path, the native data source packages in the
native-datasources directory of SPARK_HOME are now found without any configuration. The packages
of a session, added with spark.addArtifact or under spark.sql.dataSource.native.paths, take
precedence over the installed ones, like the Python data sources registered at runtime take
precedence over the installed ones.

Co-authored-by: Isaac <no-reply@databricks.com>
- Exclude the `org.apache.spark.sql.datasource` package, which has no
  counterpart in the Spark Connect client, from the client compatibility
  check, like the other packages specific to classic Spark.
- Escape the Rust and C++ examples of the native data sources page from
  Liquid, which failed to parse `{{0, end / 2}, ...}` in the C++ example.

Co-authored-by: Isaac <no-reply@databricks.com>
Instead of the native-datasources directory of SPARK_HOME, find the
library of a native data source where native libraries are usually
installed, by its file name: the library spark_datasource_<name>, such
as libspark_datasource_<name>.so on Linux, implements the data source
<name>. Spark looks for it in the native library path: the directories
of java.library.path, which include LD_LIBRARY_PATH on Linux, and then
the lib directory of each installation prefix in PATH, such as
/usr/local/lib for /usr/local/bin. The executors load the library from
the same path as the driver, or else find it in their own native
library path. The packages of a session still take precedence.

Also:
- Keep the error of a library that is loaded but rejected, because it
  cannot be loaded again in another class loader: using it again failed
  with "already loaded in another classloader".
- Add NativeDataSourceConnectSuite, which tests native data sources with
  a Spark Connect client and server, and move the helpers that build the
  test library to NativeDataSourceTestUtils.

Co-authored-by: Isaac <no-reply@databricks.com>
The Scala linter checks the format of the Spark Connect modules with
scalafmt.

Co-authored-by: Isaac <no-reply@databricks.com>
@pan3793

pan3793 commented Oct 9, 2026

Copy link
Copy Markdown
Member

Just out of curiosity, was the FFM API considered instead of JNI?

@HyukjinKwon

Copy link
Copy Markdown
Member Author

@pan3793 Good question. Yes, but FFM cannot be the only bridge yet because of the Java baseline. Spark compiles with --release 17 and supports Java 17, 21 and 25, while the FFM API is final only from Java 22 (JEP 454): it is a preview API in Java 21 and an incubator module in Java 17. So an FFM bridge would need a separate source set for Java 22+, and we would still need the JNI path for Java 17 and 21, with two interfaces for library authors to choose from.

FFM would fit this design well, though, and the design leaves room for it:

  • With FFM, a library could export a plain C ABI, without JNIEnv or the jni crate, and report errors with return codes instead of Java exceptions. Spark could also open each library with SymbolLookup.libraryLookup instead of defining a copy of NativeBridge in a class loader of its own for each library.
  • Performance does not decide it either way: the calls are coarse-grained (per scan, partition and batch), and the data crosses the boundary as Arrow C Data Interface structs, without copying.
  • The interface is versioned (NativeBridge.abiVersion, and abiVersion in the manifest of a package), and only NativeBridge is public: loading and calling the libraries is internal. Once the minimum Java version of Spark allows it, a later version of the interface can be a plain C ABI called through FFM, while Spark keeps loading the libraries of version 1.
  • Spark already runs with --enable-native-access=ALL-UNNAMED, without which recent JDKs warn about both JNI (JEP 472) and the restricted methods of FFM, so neither needs new flags.

@HyukjinKwon
HyukjinKwon marked this pull request as draft October 9, 2026 07:17
@pan3793

pan3793 commented Oct 9, 2026 •

Copy link
Copy Markdown
Member

... FFM cannot be the only bridge yet because of the Java baseline.

I suppose we can raise our JDK baseline in the upcoming 5.0.0 (early 2027)? And this can also be a 25-specific feature before that happens, with the decoupled Connect client-server architecture, upgrading server-side JDK version should be much easier.

With FFM, a library could export a plain C ABI, without JNIEnv or the jni crate ...

Yeah, this is the biggest advantage in my mind, non-JVM developers see a clean C header file without needing to learn JNI stuff.

@HyukjinKwon

Copy link
Copy Markdown
Member Author

I took a look about FFM (and did prototype). I think it's sort of good and bad. For example, we have to release header file together whereas JNI one does not require it. BTW, we will release 4.4.0 in Dec. and 4.5.0 around Mar. since we will do the quarterly release. I think we can keep JNI version, and do FFM one later around Spark 5.0.0.

@HyukjinKwon
HyukjinKwon marked this pull request as ready for review October 9, 2026 09:28
@pan3793

pan3793 commented Oct 9, 2026 •

Copy link
Copy Markdown
Member

we have to release header file together whereas JNI one does not require it.

for JNI, we define the interface via a Java class (NativeBridge), and the developer uses javap to generate the header file; for FFM, we define a header file as the interface directly.

I suppose, from the perspective of the native datasource connector developer, the latter should be more friendly (hygienic, and save one step to generate it from Java).

If we choose to define a header file as the interface directly, we can:

  • for Java 17 and 21: Spark provides a JNI bridge to call it
  • for Java 25 and later: Spark uses the FFM API to call it directly

and once we raise our JDK baseline to 25+, we can get rid of the JNI glue code, the ABI keeps no change from the beginning.

@HyukjinKwon

Copy link
Copy Markdown
Member Author

Yeah. Let's actually do that. We can deprecate this in 4.5.0, and drop it in 5.0.0 with replacing it to FFM. I am thinking about bringing more native code into Spark 5 in fact ..

@dongjoon-hyun dongjoon-hyun left a comment •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for this proposal, @HyukjinKwon. This is great! I reviewed the lookup, distribution and lifecycle of the native data sources, and left 10 inline comments, summarized below in the same order.

Main issues

  1. One invalid package breaks every lookup miss (NativeDataSourceRegistry.scala L65)
  2. Streaming tasks look up the package under the wrong session UUID (NativeDataSourceRegistry.scala L168)

Other issues

  1. Classic addArtifact(..., file=True) doesn't register the package (session.py L2299)
  2. Spark comparison semantics (NaN, collation) for pushed predicates (NativePredicates.scala L79)
  3. Artifact dedup by file name can ship a different package (NativeDataSourceRegistry.scala L152)
  4. The stream writer handle is never closed explicitly (NativeDataSource.scala L294)
  5. A failed commitDataWriter skips the task abort (NativeDataSource.scala L484)
  6. The registry lookup runs twice for each resolution (NativeDataSourceV2.scala L46)
  7. Installed native libraries now block Python data source registration (DataSource.scala L686)
  8. The scan description shows the class name instead of the data source name (ColumnarTable.scala L148)

The first two items seem to be the most important: one invalid package breaks unrelated lookups, and streaming queries on a cluster with isolated sessions would not find the package on the executors.

Comment thread python/pyspark/sql/session.py Outdated
…d lifecycle

- Skip and log invalid packages during the lookup, unless the manifest of the
  invalid package lists the requested data source.
- Distribute a package when its reader or writer is created, in the active
  session, which is the cloned session of a streaming query; executors also
  check the copies in the files of the other sessions by checksum.
- Route .sparkpkg files to the JVM in classic addArtifacts regardless of the
  flags, also in mixed lists, as Spark Connect does.
- Do not push predicates on columns with non-binary collated strings, and
  document the Spark SQL semantics (NaN, -0.0, nulls) of fully evaluated
  predicates in NativeBridge.
- Report a clear conflict when a different package with the same file name
  is already an artifact of the session.
- Close the previous stream writer of a streaming query when it creates a
  new one, and document the stream writer lifecycle.
- Keep the data writer handle when commitDataWriter fails, so that Spark
  aborts it with abortDataWriter.
- Reuse the lookup of DataSource.lookupDataSource in NativeDataSourceV2.
- Skip native data sources when registering a Python data source, which
  takes precedence over them.
- Describe the scan by the name of the data source.

Co-authored-by: Isaac <no-reply@databricks.com>

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thank you for addressing all the comments, @HyukjinKwon. I checked 9dd4f88 against each thread, and all 10 points are handled. I left 4 more minor, non-blocking comments:

  1. An invalid package that lists the name hides a valid one (NativeDataSourceRegistry.scala L88)
  2. The thread-local can keep a stale lookup on pooled threads (NativeDataSourceRegistry.scala L70)
  3. The package cache is never evicted (NativeDataSourcePackage.scala L102)
  4. No end-to-end test for closing the previous stream writer (NativeDataSource.scala L565)

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

+1, LGTM.

@sarutak sarutak left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I went over the lookup/distribution and the reader/writer lifecycle again on 9dd4f88. The fixes for the earlier round all check out in the code, and the overall design (per-library BridgeClassLoader for side-by-side loading, NativeHandle backed by AtomicLong plus Cleaner, Arrow struct release/close on every path, tryWithSafeFinally around the native calls) is solid.

I left four inline comments, none blocking:

  1. distribute collapses "no active session" and "local mode" into one branch, so a missing active session would silently skip distribution rather than fail clearly. I could not find a planning path that actually has no active session, so this is defensive hardening.
  2. The batch writer's commit closes the handle outside tryWithSafeFinally, unlike abort; if library.commit throws, the driver-side handle falls back to the Cleaner. One-line symmetry fix.
  3. spark.sql.dataSource.native.enabled defaults to true: a note on whether opt-in is preferable for the first release of a feature that loads native code (nit).
  4. NativeLibrary.load uses setAccessible across class loaders: a comment on why would help future readers (nit).

The last two are nits. Thanks for the very thorough design and docs on this.

…m writers

- Prefer a usable package or an installed library over an invalid package or a
  package of an unsupported interface version that lists the same data source,
  and only report those when nothing else provides it.
- Clear the lookup kept for the provider at the start of every
  DataSource.lookupDataSource, so a provider never takes an earlier lookup.
- Bound the cache of read packages, and forget the packages of a session when
  its artifacts are cleaned up.
- Record writer closes in the test library and check end to end that each
  micro-batch closes the stream writer of the previous one, and that each
  micro-batch creates its stream writer with the ID of the query.

Co-authored-by: Isaac <no-reply@databricks.com>

@LuciferYang LuciferYang left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I went through 3814147 again with a local build and the test library. The fixes from the earlier rounds hold up. Five more inline comments: four reproduced on 3814147, and the batch-size one is computed from the offset limit.

…batching

- Resolve the provider of classic DataFrameWriter commands and of
  DataStreamWriter sinks with the session of the Dataset active, so that the
  lookups of Python and native data sources use that session.
- Distribute a package that is not an artifact of the session as a copy named
  after its file name and checksum, which the lookups do not count as a
  package of the session: a configured package can be replaced in place, and
  packages with the same file name no longer collide. Remove the then
  unreachable NATIVE_DATA_SOURCE_PACKAGE_FILE_NAME_CONFLICT error.
- Log a warning when a package is distributed without an active session.
- Skip, and log, a configured directory that cannot be read.
- Compare the data with the read schema as it round-trips through Arrow, which
  does not carry collations, interval fields and user-defined types.
- Also end the write batches at spark.sql.execution.arrow.maxBytesPerBatch.
- Document that a failed batch commit is aborted, which closes the writer, why
  NativeBridge.load is called reflectively, and the security of discovering
  libraries by file name.

Co-authored-by: Isaac <no-reply@databricks.com>
@sarutak

sarutak commented Oct 10, 2026

Copy link
Copy Markdown
Member

One small thing before merge: could you update the PR description to mention spark.sql.execution.arrow.maxBytesPerBatch? The Internals section still says

Writes convert rows to Arrow batches of at most
spark.sql.execution.arrow.maxRecordsPerBatch rows.

Since 9ac5a49, batches are also bounded by maxBytesPerBatch, so a line like "...at most maxRecordsPerBatch rows or maxBytesPerBatch bytes, whichever is reached first" would match the code. The description already names maxRecordsPerBatch explicitly, so adding the byte limit keeps the two consistent. Not blocking, just so the description stays accurate for the record.

@HyukjinKwon

Copy link
Copy Markdown
Member Author

BTW, I plan to leave this open a while. I want others to take a look too :-).

@HyukjinKwon

Copy link
Copy Markdown
Member Author

@sarutak Thanks, updated. The Internals section now says writes convert rows to Arrow batches "that end at spark.sql.execution.arrow.maxRecordsPerBatch rows or once they hold spark.sql.execution.arrow.maxBytesPerBatch bytes, whichever is reached first, as for Python UDFs". I worded it this way because the byte limit is checked after each row, so a batch can go slightly over it, as in PythonArrowInput.

While there, I went through the rest of the description against 9ac5a49 and updated the other parts that changed during the review:

  • the per-micro-batch stream writers;
  • the precedence of usable packages over invalid or unsupported ones;
  • the checksum-named copies of configured packages;
  • the single lookup per resolution and the Python data source precedence on registration;
  • the changes to existing classes (DataFrameWriter/DataStreamWriter under withActive, ArtifactManager, DataSourceRegistration, PySpark addArtifacts, the Connect compatibility exclusion);
  • the user-facing fix for writes from a non-active session;
  • the tests added during the review.

@LuciferYang LuciferYang left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I went through 9ac5a49 again. The fixes from the earlier rounds hold up in the code. Four more inline comments, each checked on 9ac5a49 with a small test suite: Arrow Java for the first two, and a JVM ColumnarDataSource for the last two.

The first one is the one I'd like fixed before merge: when a library returns a sliced Arrow array, Spark reads the wrong rows without any error. The other three are places where checkSchema lets through a layout that ArrowColumnVector can't read, or where NativeBridge promises something Spark doesn't do (calling abort when a write fails, and calling pushLimit at most once).

…h limits once

- Read the Arrow stream of a native library array by array, and reject an array
  with a non-zero offset at any level, which Arrow Java would read from the
  start of its buffers, with a NATIVE_DATA_SOURCE_ERROR that names the column.
- Name the fields of the entries of maps `key` and `value` before creating the
  vectors, as ArrowColumnVector looks them up by name, so that maps built by
  arrow-rs can be read; reject decimal256, and report any other column that
  ArrowColumnVector cannot read with its name.
- Create the writer of a batch or streaming write with its writer factory, after
  the stages that the write depends on ran, so that Spark commits or aborts each
  writer it creates, also when such a stage fails with adaptive execution.
- Ask the reader for a limit only once, as NativeBridge documents, although
  Spark pushes the limit of LIMIT ... OFFSET twice.
- Document these in NativeBridge and the native data sources page.

Co-authored-by: Isaac <no-reply@databricks.com>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants