Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion external/duckdb
Submodule duckdb updated 1398 files
4 changes: 4 additions & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -327,6 +327,10 @@ build = [
"cmake>=3.29.0",
"ninja>=1.10",
"nanobind>=3.0",
# C API headers for the editable build. Unpinned because the test group forces numpy<2 on
# Python 3.11 for tensorflow, and one universal resolution must satisfy both. Release wheels
# build isolated against the numpy>=2.0 in [build-system].
"numpy",
"scikit_build_core>=0.11.4",
]
dev = [ # tooling like uv will install this automatically when syncing the environment
Expand Down
82 changes: 39 additions & 43 deletions src/arrow/arrow_array_stream.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ nb::object PythonTableArrowArrayStreamFactory::ProduceScanner(nb::object &arrow_
const ClientProperties &client_properties) {
D_ASSERT(!nb::isinstance<nb::capsule>(arrow_obj_handle));
ArrowSchemaWrapper schema;
PythonTableArrowArrayStreamFactory::GetSchemaInternal(arrow_obj_handle, schema);
PythonTableArrowArrayStreamFactory::GetSchemaInternal(arrow_obj_handle, schema.arrow_schema);
ArrowTableSchema arrow_table;
ArrowTableFunction::PopulateArrowTableSchema(*client_properties.client_context.get_mutable(), arrow_table,
schema.arrow_schema);
Expand All @@ -59,13 +59,12 @@ nb::object PythonTableArrowArrayStreamFactory::ProduceScanner(nb::object &arrow_
return arrow_scanner(arrow_obj_handle, **kwargs);
}

unique_ptr<ArrowArrayStreamWrapper> PythonTableArrowArrayStreamFactory::Produce(uintptr_t factory_ptr,
ArrowStreamParameters &parameters) {
unique_ptr<ArrowArrayStreamWrapper>
PythonTableArrowArrayStreamFactory::ProduceStream(ArrowStreamParameters &parameters) {
nb::gil_scoped_acquire acquire;
auto factory = static_cast<PythonTableArrowArrayStreamFactory *>(reinterpret_cast<void *>(factory_ptr)); // NOLINT
D_ASSERT(factory->arrow_object);
nb::handle arrow_obj_handle(factory->arrow_object);
auto arrow_object_type = factory->cached_arrow_type;
D_ASSERT(arrow_object.obj.ptr());
nb::handle arrow_obj_handle(arrow_object.obj);
auto arrow_object_type = cached_arrow_type;

if (arrow_object_type == PyArrowObjectType::PolarsLazyFrame) {
nb::object lf = nb::borrow<nb::object>(arrow_obj_handle);
Expand All @@ -81,9 +80,9 @@ unique_ptr<ArrowArrayStreamWrapper> PythonTableArrowArrayStreamFactory::Produce(
// rather than silently returning unfiltered rows — the arrow scan does not
// re-apply pushed filters. Mirrors the pyarrow ProduceScanner path.
if (filters && filters->HasFilters()) {
auto filter_expr = PolarsFilterPushdown::TransformFilter(
*filters, parameters.projected_columns.projection_map, parameters.projected_columns.filter_to_col,
factory->client_properties);
auto filter_expr =
PolarsFilterPushdown::TransformFilter(*filters, parameters.projected_columns.projection_map,
parameters.projected_columns.filter_to_col, client_properties);
if (!filter_expr.is(nb::none())) {
lf = lf.attr("filter")(filter_expr);
filters_pushed = true;
Expand All @@ -93,13 +92,13 @@ unique_ptr<ArrowArrayStreamWrapper> PythonTableArrowArrayStreamFactory::Produce(
// If no filters were pushed and we have a cached Arrow table, reuse it. This avoids re-reading from source and
// re-converting on repeated unfiltered scans.
nb::object arrow_table;
if (!filters_pushed && factory->cached_arrow_table.ptr() != nullptr) {
arrow_table = factory->cached_arrow_table;
if (!filters_pushed && cached_arrow_table.obj.ptr() != nullptr) {
arrow_table = cached_arrow_table.obj;
} else {
arrow_table = lf.attr("collect")().attr("to_arrow")();
// Cache only unfiltered results (filtered results are partial)
if (!filters_pushed) {
factory->cached_arrow_table = arrow_table;
cached_arrow_table.obj = arrow_table;
}
}

Expand Down Expand Up @@ -141,7 +140,7 @@ unique_ptr<ArrowArrayStreamWrapper> PythonTableArrowArrayStreamFactory::Produce(
auto &import_cache = *DuckDBPyConnection::ImportCache();
nb::object arrow_batch_scanner = import_cache.pyarrow.dataset.Scanner().attr("from_batches");
nb::handle reader_handle = reader;
auto scanner = ProduceScanner(arrow_batch_scanner, reader_handle, parameters, factory->client_properties);
auto scanner = ProduceScanner(arrow_batch_scanner, reader_handle, parameters, client_properties);
auto record_batches = scanner.attr("to_reader")();
auto res = make_uniq<ArrowArrayStreamWrapper>();
auto export_to_c = record_batches.attr("_export_to_c");
Expand Down Expand Up @@ -179,12 +178,12 @@ unique_ptr<ArrowArrayStreamWrapper> PythonTableArrowArrayStreamFactory::Produce(
// If it's a scanner we have to turn it to a record batch reader, and then a scanner again since we can't stack
// scanners on arrow Otherwise pushed-down projections and filters will disappear like tears in the rain
auto record_batches = arrow_obj_handle.attr("to_reader")();
scanner = ProduceScanner(arrow_batch_scanner, record_batches, parameters, factory->client_properties);
scanner = ProduceScanner(arrow_batch_scanner, record_batches, parameters, client_properties);
break;
}
case PyArrowObjectType::Dataset: {
nb::object arrow_scanner = arrow_obj_handle.attr("__class__").attr("scanner");
scanner = ProduceScanner(arrow_scanner, arrow_obj_handle, parameters, factory->client_properties);
scanner = ProduceScanner(arrow_scanner, arrow_obj_handle, parameters, client_properties);
break;
}
default: {
Expand All @@ -201,15 +200,15 @@ unique_ptr<ArrowArrayStreamWrapper> PythonTableArrowArrayStreamFactory::Produce(
return res;
}

void PythonTableArrowArrayStreamFactory::GetSchemaInternal(nb::handle arrow_obj_handle, ArrowSchemaWrapper &schema) {
void PythonTableArrowArrayStreamFactory::GetSchemaInternal(nb::handle arrow_obj_handle, ArrowSchema &schema) {
// PyCapsule (from bare capsule Produce path)
if (nb::isinstance<nb::capsule>(arrow_obj_handle)) {
auto capsule = nb::borrow<nb::capsule>(arrow_obj_handle);
auto stream = reinterpret_cast<ArrowArrayStream *>(capsule.data());
if (!stream->release) {
throw InvalidInputException("This ArrowArrayStream has already been consumed and cannot be scanned again.");
}
if (stream->get_schema(stream, &schema.arrow_schema)) {
if (stream->get_schema(stream, &schema)) {
throw InvalidInputException("Failed to get Arrow schema from stream: %s",
stream->get_last_error ? stream->get_last_error(stream) : "unknown error");
}
Expand All @@ -221,40 +220,37 @@ void PythonTableArrowArrayStreamFactory::GetSchemaInternal(nb::handle arrow_obj_
auto &import_cache = *DuckDBPyConnection::ImportCache();
if (duckdb::PyUtil::IsInstance(arrow_obj_handle, import_cache.pyarrow.dataset.Scanner())) {
auto obj_schema = arrow_obj_handle.attr("projected_schema");
obj_schema.attr("_export_to_c")(reinterpret_cast<uint64_t>(&schema.arrow_schema));
obj_schema.attr("_export_to_c")(reinterpret_cast<uint64_t>(&schema));
} else {
auto obj_schema = arrow_obj_handle.attr("schema");
obj_schema.attr("_export_to_c")(reinterpret_cast<uint64_t>(&schema.arrow_schema));
obj_schema.attr("_export_to_c")(reinterpret_cast<uint64_t>(&schema));
}
}

void PythonTableArrowArrayStreamFactory::GetSchema(uintptr_t factory_ptr, ArrowSchemaWrapper &schema) {
auto factory = static_cast<PythonTableArrowArrayStreamFactory *>(reinterpret_cast<void *>(factory_ptr)); // NOLINT

// Fast path: return cached schema without GIL or Python calls
if (factory->schema_cached) {
schema.arrow_schema = factory->cached_schema; // struct copy
schema.arrow_schema.release = nullptr; // non-owning copy
void PythonTableArrowArrayStreamFactory::GetSchema(ArrowSchema &schema) {
if (schema_cached.load(std::memory_order_acquire)) {
schema = cached_schema;
schema.release = nullptr;
return;
}

nb::gil_scoped_acquire acquire;
D_ASSERT(factory->arrow_object);
nb::handle arrow_obj_handle(factory->arrow_object);
D_ASSERT(arrow_object.obj.ptr());
nb::handle arrow_obj_handle(arrow_object.obj);

auto type = factory->cached_arrow_type;
auto type = cached_arrow_type;
if (type == PyArrowObjectType::PolarsLazyFrame) {
// head(0).collect().to_arrow() gives the Arrow-exported schema (e.g. large_string) without materializing data.
// collect_schema() would give Polars-native types (e.g. string_view) that don't match the actual export.
const auto empty_arrow = arrow_obj_handle.attr("head")(0).attr("collect")().attr("to_arrow")();
const auto schema_capsule = empty_arrow.attr("schema").attr("__arrow_c_schema__")();
const auto capsule = nb::borrow<nb::capsule>(schema_capsule);
const auto arrow_schema = reinterpret_cast<ArrowSchema *>(capsule.data());
factory->cached_schema = *arrow_schema;
cached_schema = *arrow_schema;
arrow_schema->release = nullptr;
factory->schema_cached = true;
schema.arrow_schema = factory->cached_schema;
schema.arrow_schema.release = nullptr;
schema_cached.store(true, std::memory_order_release);
schema = cached_schema;
schema.release = nullptr;
return;
}
if (type == PyArrowObjectType::PyCapsuleInterface || type == PyArrowObjectType::Table) {
Expand All @@ -263,26 +259,26 @@ void PythonTableArrowArrayStreamFactory::GetSchema(uintptr_t factory_ptr, ArrowS
auto schema_capsule = arrow_obj_handle.attr("__arrow_c_schema__")();
auto capsule = nb::borrow<nb::capsule>(schema_capsule);
auto arrow_schema = reinterpret_cast<ArrowSchema *>(capsule.data());
factory->cached_schema = *arrow_schema; // factory takes ownership
cached_schema = *arrow_schema;
arrow_schema->release = nullptr;
factory->schema_cached = true;
schema.arrow_schema = factory->cached_schema; // non-owning copy
schema.arrow_schema.release = nullptr;
schema_cached.store(true, std::memory_order_release);
schema = cached_schema;
schema.release = nullptr;
return;
}
// Otherwise try to use .schema with _export_to_c
if (nb::hasattr(arrow_obj_handle, "schema")) {
auto obj_schema = arrow_obj_handle.attr("schema");
if (nb::hasattr(obj_schema, "_export_to_c")) {
obj_schema.attr("_export_to_c")(reinterpret_cast<uint64_t>(&schema.arrow_schema));
obj_schema.attr("_export_to_c")(reinterpret_cast<uint64_t>(&schema));
return;
}
}
// Fallback: create a temporary stream just for the schema (consumes single-use streams!)
auto stream_capsule = arrow_obj_handle.attr("__arrow_c_stream__")();
auto capsule = nb::borrow<nb::capsule>(stream_capsule);
auto stream = reinterpret_cast<ArrowArrayStream *>(capsule.data());
if (stream->get_schema(stream, &schema.arrow_schema)) {
if (stream->get_schema(stream, &schema)) {
throw InvalidInputException("Failed to get Arrow schema from stream: %s",
stream->get_last_error ? stream->get_last_error(stream) : "unknown error");
}
Expand All @@ -292,9 +288,9 @@ void PythonTableArrowArrayStreamFactory::GetSchema(uintptr_t factory_ptr, ArrowS

// Cache for Table and Dataset (immutable schema)
if (type == PyArrowObjectType::Table || type == PyArrowObjectType::Dataset) {
factory->cached_schema = schema.arrow_schema; // factory takes ownership
schema.arrow_schema.release = nullptr; // caller gets non-owning copy
factory->schema_cached = true;
cached_schema = schema;
schema.release = nullptr;
schema_cached.store(true, std::memory_order_release);
}
}

Expand Down
36 changes: 18 additions & 18 deletions src/include/duckdb_python/arrow/arrow_array_stream.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
#include "duckdb/main/client_config.hpp"
#include "duckdb/main/config.hpp"
#include "duckdb_python/nb/casters.hpp"
#include "duckdb_python/registered_py_object.hpp"

#include "duckdb/common/string.hpp"
#include "duckdb/common/vector.hpp"
Expand Down Expand Up @@ -66,44 +67,43 @@ void TransformDuckToArrowChunk(nb::object pyarrow_schema, ArrowArray &data, nb::

PyArrowObjectType GetArrowType(const nb::handle &obj);

class PythonTableArrowArrayStreamFactory {
class PythonTableArrowArrayStreamFactory : public ArrowScanFactory {
public:
explicit PythonTableArrowArrayStreamFactory(PyObject *arrow_table, const ClientProperties &client_properties_p,
PyArrowObjectType arrow_type_p)
: arrow_object(arrow_table), client_properties(client_properties_p), cached_arrow_type(arrow_type_p) {
//! Must be constructed while holding the GIL.
PythonTableArrowArrayStreamFactory(nb::object arrow_object_p, const ClientProperties &client_properties_p,
PyArrowObjectType arrow_type_p)
: arrow_object(std::move(arrow_object_p)), client_properties(client_properties_p),
cached_arrow_type(arrow_type_p) {
cached_schema.release = nullptr;
}

~PythonTableArrowArrayStreamFactory() {
if (cached_arrow_table.ptr() != nullptr) {
nb::gil_scoped_acquire acquire;
cached_arrow_table = nb::object();
}
~PythonTableArrowArrayStreamFactory() override {
// The release callback of a schema taken from a Python producer may itself be Python.
if (cached_schema.release) {
cached_schema.release(&cached_schema);
if (nb::detail::cleanup_guard guard {}) {
cached_schema.release(&cached_schema);
}
}
}

//! Produces an Arrow Scanner, should be only called once when initializing Scan States
static unique_ptr<ArrowArrayStreamWrapper> Produce(uintptr_t factory, ArrowStreamParameters &parameters);
void GetSchema(ArrowSchema &schema) override;
unique_ptr<ArrowArrayStreamWrapper> ProduceStream(ArrowStreamParameters &parameters) override;

//! Get the schema of the arrow object
static void GetSchemaInternal(nb::handle arrow_object, ArrowSchemaWrapper &schema);
static void GetSchema(uintptr_t factory_ptr, ArrowSchemaWrapper &schema);
static void GetSchemaInternal(nb::handle arrow_object, ArrowSchema &schema);

//! Arrow Object (i.e., Scanner, Record Batch Reader, Table, Dataset)
PyObject *arrow_object;
PyObjectHolder arrow_object;

const ClientProperties client_properties;
const PyArrowObjectType cached_arrow_type;

//! Cached Arrow table from an unfiltered .collect().to_arrow() on a LazyFrame.
//! Avoids re-reading from source and re-converting on repeated scans without filters.
nb::object cached_arrow_table;
PyObjectHolder cached_arrow_table;

private:
ArrowSchema cached_schema;
bool schema_cached = false;
atomic<bool> schema_cached {false};

static nb::object ProduceScanner(nb::object &arrow_scanner, nb::handle &arrow_obj_handle,
ArrowStreamParameters &parameters, const ClientProperties &client_properties);
Expand Down
6 changes: 4 additions & 2 deletions src/include/duckdb_python/filesystem_object.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,10 @@ class FileSystemObject : public RegisteredObject {
: RegisteredObject(std::move(fs)), filenames(std::move(filenames_p)) {
}
~FileSystemObject() override {
nb::gil_scoped_acquire acquire;
// Assert that the 'obj' is a filesystem
nb::detail::cleanup_guard guard {};
if (!guard) {
return;
}
D_ASSERT(duckdb::PyUtil::IsInstance(
obj, DuckDBPyConnection::ImportCache()->duckdb.filesystem.ModifiedMemoryFileSystem()));
// Destructors are implicitly noexcept: a Python exception escaping here (fsspec `_rm` raises
Expand Down
10 changes: 10 additions & 0 deletions src/include/duckdb_python/map.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -10,11 +10,21 @@

#include "duckdb.hpp"
#include "duckdb_python/nb/casters.hpp"
#include "duckdb_python/registered_py_object.hpp"
#include "duckdb/parser/parsed_data/create_table_function_info.hpp"
#include "duckdb/execution/execution_context.hpp"

namespace duckdb {

//! Carried through the bind input so that SQL text can never forge a reference to the callable
struct MapFunctionInfo : public TableFunctionInfo {
MapFunctionInfo(nb::object function_p, nb::object schema_p)
: function(std::move(function_p)), schema(std::move(schema_p)) {
}
PyObjectHolder function;
PyObjectHolder schema;
};

struct MapFunction : public TableFunction {

public:
Expand Down
8 changes: 8 additions & 0 deletions src/include/duckdb_python/pandas/pandas_scan.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -13,9 +13,17 @@
#include "duckdb_python/pandas/pandas_bind.hpp"

#include "duckdb_python/nb/casters.hpp"
#include "duckdb_python/registered_py_object.hpp"

namespace duckdb {

//! Carried through the bind input so that SQL text can never name the dataframe
struct PandasScanInfo : public TableFunctionInfo {
explicit PandasScanInfo(nb::object df_p) : df(std::move(df_p)) {
}
PyObjectHolder df;
};

struct PandasScanFunction : public TableFunction {
public:
static constexpr idx_t PANDAS_PARTITION_COUNT = 50 * STANDARD_VECTOR_SIZE;
Expand Down
19 changes: 7 additions & 12 deletions src/include/duckdb_python/pyconnection/pyconnection.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -31,14 +31,6 @@ enum class PythonEnvironmentType { NORMAL, INTERACTIVE, JUPYTER };

struct DuckDBPyRelation;

class RegisteredArrow : public RegisteredObject {

public:
RegisteredArrow(unique_ptr<PythonTableArrowArrayStreamFactory> arrow_factory_p, nb::object obj_p)
: RegisteredObject(std::move(obj_p)), arrow_factory(std::move(arrow_factory_p)) {};
unique_ptr<PythonTableArrowArrayStreamFactory> arrow_factory;
};

struct DefaultConnectionHolder {
public:
DefaultConnectionHolder() {
Expand Down Expand Up @@ -185,7 +177,7 @@ struct DuckDBPyConnection : public std::enable_shared_from_this<DuckDBPyConnecti
// Recursive so that the outer lock taken at the top of execute/fetch
// methods (while still holding the GIL) does not deadlock against the
// inner lock taken by PrepareQuery / ExecuteInternal /
// PrepareAndExecuteInternal (after releasing the GIL). Serialises every
// PrepareAndSubmitInternal (after releasing the GIL). Serialises every
// path that touches `con.result` so concurrent calls on a single
// DuckDBPyConnection cannot dereference an already-freed result — see
// duckdb-python#435.
Expand Down Expand Up @@ -261,8 +253,9 @@ struct DuckDBPyConnection : public std::enable_shared_from_this<DuckDBPyConnecti
void ExecuteImmediately(vector<unique_ptr<SQLStatement>> statements);
unique_ptr<PreparedStatement> PrepareQuery(unique_ptr<SQLStatement> statement);
unique_ptr<QueryResult> ExecuteInternal(PreparedStatement &prep, nb::object params = nb::list());
unique_ptr<QueryResult> PrepareAndExecuteInternal(unique_ptr<SQLStatement> statement,
nb::object params = nb::list());
//! Binds the parameters and submits the statement. The handle is returned undriven.
unique_ptr<QueryResult> PrepareAndSubmitInternal(unique_ptr<SQLStatement> statement,
nb::object params = nb::list());

std::shared_ptr<DuckDBPyConnection> Execute(const nb::object &query, nb::object params = nb::list());
std::shared_ptr<DuckDBPyConnection> ExecuteFromString(const string &query);
Expand Down Expand Up @@ -367,7 +360,9 @@ struct DuckDBPyConnection : public std::enable_shared_from_this<DuckDBPyConnecti
static bool IsAcceptedArrowObject(const nb::object &object);
static NumpyObjectType IsAcceptedNumpyObject(const nb::object &object);

static unique_ptr<QueryResult> CompletePendingQuery(PendingQueryResult &pending_query);
//! Runs a submitted query to a retained, ended result on the calling thread, checking for Python
//! signals between tasks. Throws the query's error, and leaves the handle holding its collection.
static void CompleteQuery(QueryResult &result);

private:
std::unique_ptr<DuckDBPyRelation> CreateRelation(shared_ptr<Relation> rel);
Expand Down
Loading
Loading