From b53a8433969ac127126c9653f4f22357ffbd651b Mon Sep 17 00:00:00 2001 From: "logan.riggs@gmail.com" Date: Wed, 30 Sep 2026 18:01:32 +0000 Subject: [PATCH 1/2] I originally identified [this issue in arrow-java](https://github.com/apache/arrow-java/issues/601) and fixed it on the arrow-java side with a coarse grained lock on the Projector and Filter Make methods. That fixed the issue but adds some undesirable overhead in certain multi threaded workflows. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit This issue captures the deeper problem that is occurring in the C++ side of Gandiva and was not touched by my previous fix though it is prevented from happening. I have a proposed solution for this and believe it is a better overall solution. Projector::Make() and Filter::Make() read the shared expression cache once to decide the is_cached status and then call SetLLVMObjectCache(). That performs its own second, unsynchronized read of the same key before pre-loading a cached object into the LLJIT. Those two reads could disagree. If another thread compiling the identical (schema, expressions, selection vector mode, configuration) tuple inserted between them, the first thread would: take the is_cached == false path, generating expr_0_0 into its IR module, and also see a hit on the second read and addObjectFile() a cached object that defines expr_0_0 as well. Both then call JITDylib::define for the same symbol in the same JITDylib, and ORC's duplicate-symbol detection fires: CodeGenError in Gandiva: Failed to add IR module to LLJIT: In gdv_module_..., duplicate definition of symbol 'expr_0_0' C++, Gandiva Thread the single cache-lookup result (prev_cached_obj) that already determined is_cached directly into SetLLVMObjectCache(), instead of letting it perform an independent second lookup. Engine::SetLLVMObjectCache and LLVMGenerator::SetLLVMObjectCache now take the resolved shared_ptr rather than a GandivaObjectCache&. No new locking — this removes the window rather than serializing around it. New cpp/src/gandiva/tests/concurrent_make_test.cc. Verified tests fail with the duplicate-symbol error without this fix. Verified internally with our product stress tests after removing the arrow-java side of the fix. Once this change is in C++ I can make the java change to remove the synchronization. --- cpp/src/gandiva/engine.cc | 11 +- cpp/src/gandiva/engine.h | 8 +- cpp/src/gandiva/filter.cc | 7 +- cpp/src/gandiva/llvm_generator.cc | 5 +- cpp/src/gandiva/llvm_generator.h | 7 +- cpp/src/gandiva/projector.cc | 7 +- cpp/src/gandiva/tests/CMakeLists.txt | 1 + cpp/src/gandiva/tests/concurrent_make_test.cc | 205 ++++++++++++++++++ 8 files changed, 238 insertions(+), 13 deletions(-) create mode 100644 cpp/src/gandiva/tests/concurrent_make_test.cc diff --git a/cpp/src/gandiva/engine.cc b/cpp/src/gandiva/engine.cc index 205d2f866b2..56ec569ef41 100644 --- a/cpp/src/gandiva/engine.cc +++ b/cpp/src/gandiva/engine.cc @@ -304,9 +304,14 @@ void RemoveBuildTargetAttributes(llvm::Module& module) { } // namespace -Status Engine::SetLLVMObjectCache(GandivaObjectCache& object_cache) { - auto cached_buffer = object_cache.getObject(nullptr); - if (cached_buffer) { +Status Engine::SetLLVMObjectCache( + const std::shared_ptr& prev_cached_obj) { + if (prev_cached_obj) { + // Copy rather than hand over the cache's own buffer -- addObjectFile takes + // ownership and may free it, but the shared cache map holds a reference to the + // same object too. + auto cached_buffer = prev_cached_obj->getMemBufferCopy( + prev_cached_obj->getBuffer(), prev_cached_obj->getBufferIdentifier()); auto error = lljit_->addObjectFile(std::move(cached_buffer)); if (error) { return Status::CodeGenError("Failed to add cached object file to LLJIT: ", diff --git a/cpp/src/gandiva/engine.h b/cpp/src/gandiva/engine.h index 20165787cb6..23fc366fe15 100644 --- a/cpp/src/gandiva/engine.h +++ b/cpp/src/gandiva/engine.h @@ -74,8 +74,12 @@ class GANDIVA_EXPORT Engine { /// Optimise and compile the module. Status FinalizeModule(); - /// Set LLVM ObjectCache. - Status SetLLVMObjectCache(GandivaObjectCache& object_cache); + /// Pre-load the LLJIT with an already-compiled object for the current module, if + /// one was found by the caller's own cache lookup. Takes the lookup result directly + /// rather than re-querying the cache, so this can never disagree with the \p cached + /// flag passed to \p Make -- avoiding the load being based on a newer cache state + /// than the one that decided whether Make() would still build (and load) fresh IR. + Status SetLLVMObjectCache(const std::shared_ptr& prev_cached_obj); /// Get the compiled function corresponding to the irfunction. Result CompiledFunction(const std::string& function); diff --git a/cpp/src/gandiva/filter.cc b/cpp/src/gandiva/filter.cc index 8a270cfdc06..be67c12254f 100644 --- a/cpp/src/gandiva/filter.cc +++ b/cpp/src/gandiva/filter.cc @@ -76,8 +76,11 @@ Status Filter::Make(SchemaPtr schema, ConditionPtr condition, ARROW_RETURN_NOT_OK(expr_validator.Validate(condition)); } - // Set the object cache for LLVM - ARROW_RETURN_NOT_OK(llvm_gen->SetLLVMObjectCache(obj_cache)); + // Pre-load the already-compiled object from the cache lookup above, if there was + // one. Reuses prev_cached_obj instead of letting SetLLVMObjectCache re-query the + // shared cache, so this can never disagree with the is_cached snapshot that + // decided whether Build() below still compiles fresh IR (see GH-601). + ARROW_RETURN_NOT_OK(llvm_gen->SetLLVMObjectCache(prev_cached_obj)); ARROW_RETURN_NOT_OK(llvm_gen->Build({condition}, SelectionVector::Mode::MODE_NONE)); diff --git a/cpp/src/gandiva/llvm_generator.cc b/cpp/src/gandiva/llvm_generator.cc index 7caafa58369..7aa91c4dd6b 100644 --- a/cpp/src/gandiva/llvm_generator.cc +++ b/cpp/src/gandiva/llvm_generator.cc @@ -64,8 +64,9 @@ LLVMGenerator::GetCache() { return shared_cache; } -Status LLVMGenerator::SetLLVMObjectCache(GandivaObjectCache& object_cache) { - return engine_->SetLLVMObjectCache(object_cache); +Status LLVMGenerator::SetLLVMObjectCache( + const std::shared_ptr& prev_cached_obj) { + return engine_->SetLLVMObjectCache(prev_cached_obj); } Status LLVMGenerator::Add(const ExpressionPtr expr, const FieldDescriptorPtr output) { diff --git a/cpp/src/gandiva/llvm_generator.h b/cpp/src/gandiva/llvm_generator.h index c6e4e821dfe..0c796b4b2a2 100644 --- a/cpp/src/gandiva/llvm_generator.h +++ b/cpp/src/gandiva/llvm_generator.h @@ -58,8 +58,11 @@ class GANDIVA_EXPORT LLVMGenerator { static std::shared_ptr>> GetCache(); - /// \brief Set LLVM ObjectCache. - Status SetLLVMObjectCache(GandivaObjectCache& object_cache); + /// \brief Pre-load the engine with an already-compiled object for the current + /// module, if the caller's own cache lookup found one. See the equivalent Engine + /// method for why this takes the lookup result directly instead of re-querying + /// the cache. + Status SetLLVMObjectCache(const std::shared_ptr& prev_cached_obj); /// \brief Build the code for the expression trees for default mode with a LLVM /// ObjectCache. Each element in the vector represents an expression tree diff --git a/cpp/src/gandiva/projector.cc b/cpp/src/gandiva/projector.cc index ec0302146ff..eb52477ff6d 100644 --- a/cpp/src/gandiva/projector.cc +++ b/cpp/src/gandiva/projector.cc @@ -94,8 +94,11 @@ Status Projector::Make(SchemaPtr schema, const ExpressionVector& exprs, } } - // Set the object cache for LLVM - ARROW_RETURN_NOT_OK(llvm_gen->SetLLVMObjectCache(obj_cache)); + // Pre-load the already-compiled object from the cache lookup above, if there was + // one. Reuses prev_cached_obj instead of letting SetLLVMObjectCache re-query the + // shared cache, so this can never disagree with the is_cached snapshot that + // decided whether Build() below still compiles fresh IR (see GH-601). + ARROW_RETURN_NOT_OK(llvm_gen->SetLLVMObjectCache(prev_cached_obj)); ARROW_RETURN_NOT_OK(llvm_gen->Build(exprs, selection_vector_mode)); diff --git a/cpp/src/gandiva/tests/CMakeLists.txt b/cpp/src/gandiva/tests/CMakeLists.txt index be635e16d3a..52d32cb7f01 100644 --- a/cpp/src/gandiva/tests/CMakeLists.txt +++ b/cpp/src/gandiva/tests/CMakeLists.txt @@ -19,6 +19,7 @@ add_gandiva_test(projector-test SOURCES binary_test.cc boolean_expr_test.cc + concurrent_make_test.cc date_time_test.cc decimal_alignment_test.cc decimal_single_test.cc diff --git a/cpp/src/gandiva/tests/concurrent_make_test.cc b/cpp/src/gandiva/tests/concurrent_make_test.cc new file mode 100644 index 00000000000..f6bb7e80534 --- /dev/null +++ b/cpp/src/gandiva/tests/concurrent_make_test.cc @@ -0,0 +1,205 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +// Regression tests for GH-601: concurrent Projector::Make()/Filter::Make() calls for the +// *same* expression cache key. +// +// Before the fix, Make() read the shared object cache once to decide `is_cached`, then +// SetLLVMObjectCache() performed a second, independent read of the same key. A thread that +// saw a miss on the first read and a hit on the second one would both pre-load the cached +// object (defining e.g. "expr_0_0" in its JITDylib) *and* add its own freshly compiled IR +// module defining the same symbol, so LLVM ORC's duplicate-symbol detection fired and +// Make() returned "Failed to add IR module to LLJIT: Duplicate definition of symbol". +// +// Note that this surfaces as a returned Status rather than a crash, so every Make() below +// must be asserted on -- that is precisely what the pre-existing Java-side repro +// (ProjectorTest#testMakeProjectorParallel) failed to do. + +#include + +#include +#include +#include +#include +#include +#include + +#include "arrow/memory_pool.h" +#include "gandiva/filter.h" +#include "gandiva/projector.h" +#include "gandiva/tests/test_util.h" +#include "gandiva/tree_expr_builder.h" + +namespace gandiva { + +using arrow::int32; + +namespace { + +// Releases every waiting thread at once, so the racing Make() calls actually overlap +// instead of trickling in as each thread is spawned. Hand-rolled rather than std::latch +// to avoid depending on the standard library's C++20 support. +class StartGate { + public: + void Wait() { + std::unique_lock lock(mutex_); + cv_.wait(lock, [this] { return open_; }); + } + + void Open() { + { + std::lock_guard lock(mutex_); + open_ = true; + } + cv_.notify_all(); + } + + private: + std::mutex mutex_; + std::condition_variable cv_; + bool open_ = false; +}; + +int NumThreads() { + auto hw = static_cast(std::thread::hardware_concurrency()); + return std::max(8, 2 * std::max(1, hw)); +} + +// Each iteration needs a cache key that has never been built before, so that every thread +// in that iteration genuinely starts from a miss. Note that TestConfiguration() returns +// the *same* shared Configuration instance every call, and ExpressionCacheKey compares +// configurations by pointer identity, so varying the configuration would not produce a +// miss -- vary the schema instead. +constexpr int kIterations = 25; + +} // namespace + +class TestConcurrentMake : public ::testing::Test { + public: + void SetUp() { pool_ = arrow::default_memory_pool(); } + + protected: + arrow::MemoryPool* pool_; +}; + +TEST_F(TestConcurrentMake, TestProjectorMakeSameCacheKey) { + const int num_threads = NumThreads(); + + for (int iter = 0; iter < kIterations; ++iter) { + // Unique per iteration => guaranteed cache miss for the first thread to get there. + auto field0 = field("f0_" + std::to_string(iter), int32()); + auto field1 = field("f1_" + std::to_string(iter), int32()); + auto schema = arrow::schema({field0, field1}); + auto field_sum = field("add_" + std::to_string(iter), int32()); + auto sum_expr = TreeExprBuilder::MakeExpression("add", {field0, field1}, field_sum); + auto configuration = TestConfiguration(); + + StartGate gate; + std::vector threads; + std::vector statuses(num_threads); + std::vector> projectors(num_threads); + + threads.reserve(num_threads); + for (int i = 0; i < num_threads; ++i) { + threads.emplace_back([&, i] { + gate.Wait(); + statuses[i] = Projector::Make(schema, {sum_expr}, configuration, &projectors[i]); + }); + } + gate.Open(); + for (auto& thread : threads) { + thread.join(); + } + + // Create a row-batch with some sample data. + int num_records = 4; + auto array0 = MakeArrowArrayInt32({1, 2, 3, 4}, {true, true, true, true}); + auto array1 = MakeArrowArrayInt32({11, 13, 15, 17}, {true, true, true, true}); + auto exp_sum = MakeArrowArrayInt32({12, 15, 18, 21}, {true, true, true, true}); + auto in_batch = arrow::RecordBatch::Make(schema, num_records, {array0, array1}); + + for (int i = 0; i < num_threads; ++i) { + ASSERT_OK(statuses[i]) << "iteration " << iter << ", thread " << i; + ASSERT_NE(projectors[i], nullptr) << "iteration " << iter << ", thread " << i; + + // Every projector must evaluate correctly, including the ones built from the cache + // hit path. No pre-existing test evaluates a cache-hit projector, so a + // corrupt-but-non-null cached module would otherwise go unnoticed. + arrow::ArrayVector outputs; + ASSERT_OK(projectors[i]->Evaluate(*in_batch, pool_, &outputs)) + << "iteration " << iter << ", thread " << i; + EXPECT_ARROW_ARRAY_EQUALS(exp_sum, outputs.at(0)); + } + } +} + +TEST_F(TestConcurrentMake, TestFilterMakeSameCacheKey) { + const int num_threads = NumThreads(); + + for (int iter = 0; iter < kIterations; ++iter) { + auto field0 = field("g0_" + std::to_string(iter), int32()); + auto field1 = field("g1_" + std::to_string(iter), int32()); + auto schema = arrow::schema({field0, field1}); + + // Condition: f0 + f1 < 10 + auto node_f0 = TreeExprBuilder::MakeField(field0); + auto node_f1 = TreeExprBuilder::MakeField(field1); + auto sum_func = + TreeExprBuilder::MakeFunction("add", {node_f0, node_f1}, arrow::int32()); + auto literal_10 = TreeExprBuilder::MakeLiteral(static_cast(10)); + auto less_than_10 = TreeExprBuilder::MakeFunction("less_than", {sum_func, literal_10}, + arrow::boolean()); + auto condition = TreeExprBuilder::MakeCondition(less_than_10); + auto configuration = TestConfiguration(); + + StartGate gate; + std::vector threads; + std::vector statuses(num_threads); + std::vector> filters(num_threads); + + threads.reserve(num_threads); + for (int i = 0; i < num_threads; ++i) { + threads.emplace_back([&, i] { + gate.Wait(); + statuses[i] = Filter::Make(schema, condition, configuration, &filters[i]); + }); + } + gate.Open(); + for (auto& thread : threads) { + thread.join(); + } + + int num_records = 5; + auto array0 = MakeArrowArrayInt32({1, 2, 3, 4, 6}, {true, true, true, false, true}); + auto array1 = MakeArrowArrayInt32({5, 9, 6, 17, 3}, {true, true, false, true, true}); + auto exp = MakeArrowArrayUint16({0, 4}); + auto in_batch = arrow::RecordBatch::Make(schema, num_records, {array0, array1}); + + for (int i = 0; i < num_threads; ++i) { + ASSERT_OK(statuses[i]) << "iteration " << iter << ", thread " << i; + ASSERT_NE(filters[i], nullptr) << "iteration " << iter << ", thread " << i; + + std::shared_ptr selection_vector; + ASSERT_OK(SelectionVector::MakeInt16(num_records, pool_, &selection_vector)); + ASSERT_OK(filters[i]->Evaluate(*in_batch, selection_vector)) + << "iteration " << iter << ", thread " << i; + EXPECT_ARROW_ARRAY_EQUALS(exp, selection_vector->ToArray()); + } + } +} + +} // namespace gandiva From 7894cd99e87eb365edc0c6398e4411621e0f84b0 Mon Sep 17 00:00:00 2001 From: "logan.riggs@gmail.com" Date: Wed, 30 Sep 2026 18:10:41 +0000 Subject: [PATCH 2/2] lint --- cpp/src/gandiva/tests/concurrent_make_test.cc | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/cpp/src/gandiva/tests/concurrent_make_test.cc b/cpp/src/gandiva/tests/concurrent_make_test.cc index f6bb7e80534..af5ea310942 100644 --- a/cpp/src/gandiva/tests/concurrent_make_test.cc +++ b/cpp/src/gandiva/tests/concurrent_make_test.cc @@ -19,11 +19,12 @@ // *same* expression cache key. // // Before the fix, Make() read the shared object cache once to decide `is_cached`, then -// SetLLVMObjectCache() performed a second, independent read of the same key. A thread that -// saw a miss on the first read and a hit on the second one would both pre-load the cached -// object (defining e.g. "expr_0_0" in its JITDylib) *and* add its own freshly compiled IR -// module defining the same symbol, so LLVM ORC's duplicate-symbol detection fired and -// Make() returned "Failed to add IR module to LLJIT: Duplicate definition of symbol". +// SetLLVMObjectCache() performed a second, independent read of the same key. A thread +// that saw a miss on the first read and a hit on the second one would both pre-load the +// cached object (defining e.g. "expr_0_0" in its JITDylib) *and* add its own freshly +// compiled IR module defining the same symbol, so LLVM ORC's duplicate-symbol detection +// fired and Make() returned "Failed to add IR module to LLJIT: Duplicate definition of +// symbol". // // Note that this surfaces as a returned Status rather than a crash, so every Make() below // must be asserted on -- that is precisely what the pre-existing Java-side repro