diff --git a/gandiva/src/main/java/org/apache/arrow/gandiva/evaluator/Filter.java b/gandiva/src/main/java/org/apache/arrow/gandiva/evaluator/Filter.java index 6c8540cde8..06577a516e 100644 --- a/gandiva/src/main/java/org/apache/arrow/gandiva/evaluator/Filter.java +++ b/gandiva/src/main/java/org/apache/arrow/gandiva/evaluator/Filter.java @@ -111,14 +111,18 @@ public static Filter make(Schema schema, Condition condition, boolean optimize) * @param configurationId Custom configuration created through config builder. * @return A native evaluator object that can be used to invoke these projections on a RecordBatch */ - public static synchronized Filter make(Schema schema, Condition condition, long configurationId) + public static Filter make(Schema schema, Condition condition, long configurationId) throws GandivaException { // Invoke the JNI layer to create the LLVM module representing the filter. GandivaTypes.Condition conditionBuf = condition.toProtobuf(); GandivaTypes.Schema schemaBuf = ArrowTypeHelper.arrowSchemaToProtobuf(schema); + byte[] schemaBytes = schemaBuf.toByteArray(); + byte[] conditionBytes = conditionBuf.toByteArray(); JniWrapper wrapper = JniLoader.getInstance().getWrapper(); - long moduleId = - wrapper.buildFilter(schemaBuf.toByteArray(), conditionBuf.toByteArray(), configurationId); + // No lock here, deliberately -- see the equivalent comment in Projector.make(). The + // duplicate-LLVM-symbol race from GH-601 is fixed in Gandiva's native Filter::Make(), so this + // no longer needs to be serialized. + long moduleId = wrapper.buildFilter(schemaBytes, conditionBytes, configurationId); logger.debug("Created module for the filter with id {}", moduleId); return new Filter(wrapper, moduleId, schema); } diff --git a/gandiva/src/main/java/org/apache/arrow/gandiva/evaluator/Projector.java b/gandiva/src/main/java/org/apache/arrow/gandiva/evaluator/Projector.java index 9c5b22d659..13fc040f34 100644 --- a/gandiva/src/main/java/org/apache/arrow/gandiva/evaluator/Projector.java +++ b/gandiva/src/main/java/org/apache/arrow/gandiva/evaluator/Projector.java @@ -188,7 +188,7 @@ public static Projector make( * @param configurationId Custom configuration created through config builder. * @return A native evaluator object that can be used to invoke these projections on a RecordBatch */ - public static synchronized Projector make( + public static Projector make( Schema schema, List exprs, SelectionVectorType selectionVectorType, @@ -202,13 +202,19 @@ public static synchronized Projector make( // Invoke the JNI layer to create the LLVM module representing the expressions GandivaTypes.Schema schemaBuf = ArrowTypeHelper.arrowSchemaToProtobuf(schema); + byte[] schemaBytes = schemaBuf.toByteArray(); + byte[] exprBytes = builder.build().toByteArray(); JniWrapper wrapper = JniLoader.getInstance().getWrapper(); + // No lock here, deliberately. This method used to be `static synchronized`, which serialized + // every LLVM compilation in the process, because concurrent builds for the same native + // expression-cache key could race and fail with a duplicate-LLVM-symbol error (GH-601). + // That race was a double read of the native object cache inside Projector::Make(), and it is + // fixed in Gandiva itself -- see "GH-601: Fix TOCTOU race in Gandiva's LLVM object cache read" + // in the arrow C++ tree. Concurrent buildProjector() calls are safe against a Gandiva that + // contains that fix; do not re-add a lock here without first checking the native side. long moduleId = wrapper.buildProjector( - schemaBuf.toByteArray(), - builder.build().toByteArray(), - selectionVectorType.getNumber(), - configurationId); + schemaBytes, exprBytes, selectionVectorType.getNumber(), configurationId); logger.debug("Created module for the projector with id {}", moduleId); return new Projector(wrapper, moduleId, schema, exprs.size()); } diff --git a/gandiva/src/test/java/org/apache/arrow/gandiva/evaluator/MakeConcurrencyBenchmark.java b/gandiva/src/test/java/org/apache/arrow/gandiva/evaluator/MakeConcurrencyBenchmark.java new file mode 100644 index 0000000000..397e850894 --- /dev/null +++ b/gandiva/src/test/java/org/apache/arrow/gandiva/evaluator/MakeConcurrencyBenchmark.java @@ -0,0 +1,277 @@ +/* + * 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. + */ +package org.apache.arrow.gandiva.evaluator; + +import com.google.common.collect.Lists; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicLong; +import org.apache.arrow.gandiva.expression.ExpressionTree; +import org.apache.arrow.gandiva.expression.TreeBuilder; +import org.apache.arrow.gandiva.expression.TreeNode; +import org.apache.arrow.vector.types.pojo.Field; +import org.apache.arrow.vector.types.pojo.Schema; +import org.junit.jupiter.api.Disabled; +import org.junit.jupiter.api.Test; + +/** + * Measures the throughput of concurrent {@link Projector#make} calls, to quantify what removing the + * process-wide lock on {@code make()} actually buys (GH-601). + * + *

{@code Projector.make()} used to be declared {@code static synchronized}, so every LLVM module + * build in the JVM serialized on one monitor regardless of whether the calls were related. That + * monitor was {@code Projector.class}, which means the old behaviour can be reproduced exactly by + * wrapping the call in {@code synchronized (Projector.class)}. This benchmark runs both arms in a + * single JVM over identical pre-generated inputs, so the comparison does not depend on two builds + * or two machine states. + * + *

Two workloads, because a cache hit is not free: + * + *

+ * + *

Do not expect linear scaling from the unlocked arm. Two process-wide native locks remain on + * this path — the expression cache's own mutex, taken on every lookup and insert, and the JNI + * module-id map's mutex, taken on every make and every close. Where the curve flattens is where + * those start to dominate. + * + *

Disabled by default; this is a benchmark, not a test. Run it with: + * + *

+ * mvn -Parrow-jni -pl gandiva -Dtest=MakeConcurrencyBenchmark -DfailIfNoTests=false \
+ *     -Djunit.jupiter.conditions.deactivate='*' \
+ *     -Darrow.cpp.build.dir=/path/to/dir/containing/libgandiva_jni.so test
+ * 
+ * + *

Keep {@code COLD_MAKES * THREAD_COUNTS.length} well under Gandiva's native cache capacity + * (default 5000, overridable only via the {@code GANDIVA_CACHE_SIZE} environment variable, which is + * read once per process) so that LRU eviction does not perturb the cold numbers. + */ +@Disabled("perf benchmark; run explicitly, see class javadoc") +public class MakeConcurrencyBenchmark extends BaseEvaluatorTest { + + private static final int[] THREAD_COUNTS = {1, 2, 4, 8, 16}; + + /** Distinct compiles per measured cold run. Each is tens of ms, so keep this modest. */ + private static final int COLD_MAKES = 256; + + /** Cache-hit makes per measured warm run. Each is ~ms, so this can be larger. */ + private static final int WARM_MAKES = 2048; + + /** Distinct expressions burned by warmup passes, reserved outside the measured key ranges. */ + private static final int WARMUP_MAKES = 64; + + /** Monotonic source of distinct cache keys, so no two runs in this JVM collide. */ + private static final AtomicLong KEY_SEQUENCE = new AtomicLong(); + + /** Whether the timed call emulates the old {@code static synchronized} declaration. */ + private enum Arm { + LOCKED("legacy (static synchronized)"), + UNLOCKED("current (no lock)"); + + private final String label; + + Arm(String label) { + this.label = label; + } + } + + /** + * Builds one expression whose native cache key is unique to {@code key}. + * + *

Varying an integer literal is the cheapest way to get a distinct key: Gandiva's cache key is + * derived from the printed form of the expression, and changing the literal leaves the generated + * IR the same size, so per-compile cost stays flat across the sweep. Note this deliberately + * avoids any {@code like()} function — Gandiva mixes the thread id into the cache key for those, + * which would make the hit/miss ratio depend on thread count. + */ + private List distinctExpr(long key) { + Field x = Field.nullable("x", int32); + TreeNode sum = + TreeBuilder.makeFunction( + "add", + Lists.newArrayList(TreeBuilder.makeField(x), TreeBuilder.makeLiteral((int) key)), + int32); + return Lists.newArrayList(TreeBuilder.makeExpression(sum, Field.nullable("out", int32))); + } + + private Schema schema() { + return new Schema(Lists.newArrayList(Field.nullable("x", int32))); + } + + /** + * Runs {@code numMakes} Projector.make() calls across {@code numThreads} threads and returns the + * wall-clock nanos for the whole batch. + * + * @param keys the cache key to use for each call; pre-generated so expression construction is not + * inside the timed region. Repeat a value to force a cache hit. + */ + private long timeMakes(Arm arm, int numThreads, List keys) throws Exception { + final Schema schema = schema(); + + // Build every expression up front: only make() should be inside the timed region. + List> exprs = new ArrayList<>(keys.size()); + for (Long key : keys) { + exprs.add(distinctExpr(key)); + } + + ExecutorService pool = Executors.newFixedThreadPool(numThreads); + List failures = Collections.synchronizedList(new ArrayList<>()); + // All threads meet here, so the measured window is contention and not pool ramp-up. + CyclicBarrier gate = new CyclicBarrier(numThreads + 1); + + try { + List> futures = new ArrayList<>(); + for (int t = 0; t < numThreads; t++) { + final int threadIndex = t; + futures.add( + pool.submit( + () -> { + try { + gate.await(60, TimeUnit.SECONDS); + // Stripe the work so every thread does an equal share. + for (int i = threadIndex; i < exprs.size(); i += numThreads) { + makeOnce(arm, schema, exprs.get(i)); + } + } catch (Throwable e) { + failures.add(e); + } + })); + } + + gate.await(60, TimeUnit.SECONDS); + long start = System.nanoTime(); + for (Future future : futures) { + future.get(10, TimeUnit.MINUTES); + } + long elapsed = System.nanoTime() - start; + + if (!failures.isEmpty()) { + AssertionError error = new AssertionError("make() failed during benchmark"); + error.initCause(failures.get(0)); + throw error; + } + return elapsed; + } finally { + pool.shutdownNow(); + } + } + + /** + * One make/close pair. The LOCKED arm holds the {@code Projector.class} monitor across the call, + * which is exactly what {@code public static synchronized Projector make(...)} did. Verified that + * this mechanism works the same as the original sycnhronized method call. + */ + private void makeOnce(Arm arm, Schema schema, List exprs) throws Exception { + Projector projector = null; + try { + if (arm == Arm.LOCKED) { + synchronized (Projector.class) { + projector = Projector.make(schema, exprs); + } + } else { + projector = Projector.make(schema, exprs); + } + } finally { + // Projector is not AutoCloseable and has no Cleaner: a dropped reference leaks a whole LLJIT, + // which at these iteration counts will exhaust memory. + if (projector != null) { + projector.close(); + } + } + } + + /** Distinct key per call: every make() is a real LLVM compile. */ + private List coldKeys(int count) { + List keys = new ArrayList<>(count); + for (int i = 0; i < count; i++) { + keys.add(KEY_SEQUENCE.incrementAndGet()); + } + return keys; + } + + /** One shared key: all but the first make() hits the native object cache. */ + private List warmKeys(int count) { + long key = KEY_SEQUENCE.incrementAndGet(); + List keys = new ArrayList<>(count); + for (int i = 0; i < count; i++) { + keys.add(key); + } + return keys; + } + + private void warmup() throws Exception { + // Absorbs JVM JIT, the one-time JniLoader native load, and LLVM's one-time target init. + timeMakes(Arm.UNLOCKED, 4, coldKeys(WARMUP_MAKES)); + } + + @Test + public void benchmarkConcurrentMake() throws Exception { + warmup(); + + StringBuilder table = new StringBuilder(); + table.append("\n### Projector.make() throughput (makes/sec, higher is better)\n\n"); + table.append("cold = distinct expression per call (real LLVM compile)\n"); + table.append("warm = one shared expression (native object-cache hits)\n\n"); + table.append("| threads | cold legacy | cold current | cold speedup "); + table.append("| warm legacy | warm current | warm speedup |\n"); + table.append("|--------:|------------:|-------------:|-------------:"); + table.append("|------------:|-------------:|-------------:|\n"); + + for (int numThreads : THREAD_COUNTS) { + double coldLocked = rate(COLD_MAKES, timeMakes(Arm.LOCKED, numThreads, coldKeys(COLD_MAKES))); + double coldFree = rate(COLD_MAKES, timeMakes(Arm.UNLOCKED, numThreads, coldKeys(COLD_MAKES))); + double warmLocked = rate(WARM_MAKES, timeMakes(Arm.LOCKED, numThreads, warmKeys(WARM_MAKES))); + double warmFree = rate(WARM_MAKES, timeMakes(Arm.UNLOCKED, numThreads, warmKeys(WARM_MAKES))); + + table.append( + String.format( + "| %7d | %11.1f | %12.1f | %11.2fx | %11.1f | %12.1f | %11.2fx |%n", + numThreads, + coldLocked, + coldFree, + coldFree / coldLocked, + warmLocked, + warmFree, + warmFree / warmLocked)); + } + + table.append("\nArms: legacy = ").append(Arm.LOCKED.label); + table.append(", current = ").append(Arm.UNLOCKED.label).append('\n'); + table.append("Cores available: ").append(Runtime.getRuntime().availableProcessors()); + table.append(" (speedup is bounded by this, not by thread count)\n"); + + System.out.println(table); + } + + private static double rate(int numMakes, long elapsedNanos) { + return numMakes * 1_000_000_000.0 / elapsedNanos; + } +} diff --git a/gandiva/src/test/java/org/apache/arrow/gandiva/evaluator/ProjectorTest.java b/gandiva/src/test/java/org/apache/arrow/gandiva/evaluator/ProjectorTest.java index f2590226b1..f2d626f7ab 100644 --- a/gandiva/src/test/java/org/apache/arrow/gandiva/evaluator/ProjectorTest.java +++ b/gandiva/src/test/java/org/apache/arrow/gandiva/evaluator/ProjectorTest.java @@ -19,6 +19,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotEquals; +import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; @@ -29,10 +30,13 @@ import java.time.Instant; import java.util.ArrayList; import java.util.Arrays; +import java.util.Collections; import java.util.List; import java.util.Set; +import java.util.concurrent.CyclicBarrier; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; +import java.util.concurrent.Future; import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.IntStream; import org.apache.arrow.gandiva.exceptions.GandivaException; @@ -99,12 +103,93 @@ List binaryBufs(String[] strings) { return varBufs(strings, utf16Charset); } + /** + * Builds projectors concurrently and fails if any of them errors out. + * + *

Failures here surface as a {@link GandivaException} carrying a native {@code CodeGenError}, + * not as a JVM crash -- see GH-601, where concurrent builds for the same native expression-cache + * key raced and produced "Duplicate definition of symbol 'expr_0_0'". Every exception must + * therefore be collected and re-thrown; swallowing them makes this test unable to observe the + * very bug it exists for. + * + * @param schemas schemas to pick from, one per task, round-robin + * @param exprs the expressions to compile + * @param configOptions custom configuration, or null for the default + */ + private void makeProjectorsConcurrently( + List schemas, + List exprs, + ConfigurationBuilder.ConfigOptions configOptions) + throws Exception { + final int numTasks = 1000; + ExecutorService executors = Executors.newFixedThreadPool(16); + List failures = Collections.synchronizedList(new ArrayList<>()); + + try { + IntStream.range(0, numTasks) + .forEach( + i -> { + Schema schema = schemas.get(i % schemas.size()); + executors.submit( + () -> { + Projector evaluator = null; + try { + evaluator = + configOptions == null + ? Projector.make(schema, exprs) + : Projector.make(schema, exprs, configOptions); + assertNotNull(evaluator); + } catch (Throwable t) { + failures.add(t); + } finally { + if (evaluator != null) { + try { + evaluator.close(); + } catch (Throwable t) { + failures.add(t); + } + } + } + }); + }); + executors.shutdown(); + assertTrue( + executors.awaitTermination(100, java.util.concurrent.TimeUnit.SECONDS), + "projector builds did not finish within the timeout"); + } finally { + executors.shutdownNow(); + } + + if (!failures.isEmpty()) { + AssertionError error = + new AssertionError( + failures.size() + + " of " + + numTasks + + " concurrent Projector.make() calls failed; first failure attached"); + error.initCause(failures.get(0)); + throw error; + } + } + + private List greaterThanExpr(Field a, Field b) { + TreeNode aNode = TreeBuilder.makeField(a); + TreeNode bNode = TreeBuilder.makeField(b); + List args = Lists.newArrayList(aNode, bNode); + + TreeNode cond = TreeBuilder.makeFunction("greater_than", args, boolType); + TreeNode ifNode = TreeBuilder.makeIf(cond, aNode, bNode, int64); + + ExpressionTree expr = TreeBuilder.makeExpression(ifNode, Field.nullable("c", int64)); + return Lists.newArrayList(expr); + } + private void testMakeProjectorParallel(ConfigurationBuilder.ConfigOptions configOptions) - throws InterruptedException { + throws Exception { List schemas = Lists.newArrayList(); Field a = Field.nullable("a", int64); Field b = Field.nullable("b", int64); - IntStream.range(0, 1000) + IntStream.range(0, 100) .forEach( i -> { Field c = Field.nullable("" + i, int64); @@ -112,40 +197,94 @@ private void testMakeProjectorParallel(ConfigurationBuilder.ConfigOptions config schemas.add(new Schema(cols)); }); - TreeNode aNode = TreeBuilder.makeField(a); - TreeNode bNode = TreeBuilder.makeField(b); - List args = Lists.newArrayList(aNode, bNode); + // Build projectors in parallel across many distinct schemas. This is the throughput case: + // the cache keys mostly differ, so these compilations should proceed independently. + makeProjectorsConcurrently(schemas, greaterThanExpr(a, b), configOptions); + } - TreeNode cond = TreeBuilder.makeFunction("greater_than", args, boolType); - TreeNode ifNode = TreeBuilder.makeIf(cond, aNode, bNode, int64); + /** + * The actual GH-601 shape: every thread races to build the *same* native expression-cache entry, + * starting from a genuine cache miss. + * + *

Two details matter, and both were missing from the older parallel test. First, the threads + * are released from a {@link CyclicBarrier} rather than trickling in as the executor ramps up, so + * the builds genuinely overlap. Second, each round uses a fresh schema: once a key is in the + * native cache every later build takes the cached path and the race window is gone, so a single + * key gives at most one chance to observe it. + * + *

Be aware of what this test is and is not. Against a Gandiva without the native fix it does + * reproduce the duplicate-symbol failure, but only at roughly one round in 300 -- at the round + * count below it will usually pass even on affected builds. Treat it as a smoke test that + * concurrent same-key builds succeed and produce usable projectors. The reliable regression guard + * for GH-601 is the native test (TestConcurrentMake in + * cpp/src/gandiva/tests/concurrent_make_test.cc), which reproduces within a handful of iterations + * because it races the cache read directly without the protobuf and JNI round trip in between. + */ + private void testMakeProjectorParallelSameKey(ConfigurationBuilder.ConfigOptions configOptions) + throws Exception { + final int numThreads = 16; + final int numRounds = 50; - ExpressionTree expr = TreeBuilder.makeExpression(ifNode, Field.nullable("c", int64)); - List exprs = Lists.newArrayList(expr); + Field a = Field.nullable("a", int64); + Field b = Field.nullable("b", int64); + List exprs = greaterThanExpr(a, b); - // build projectors in parallel choosing schema at random - // this should hit the same cache entry thus exposing - // any threading issues. - ExecutorService executors = Executors.newFixedThreadPool(16); + ExecutorService executors = Executors.newFixedThreadPool(numThreads); + List failures = Collections.synchronizedList(new ArrayList<>()); - IntStream.range(0, 1000) - .forEach( - i -> { + try { + for (int round = 0; round < numRounds; round++) { + // Unique per round => the first thread to arrive genuinely misses the native cache. + Field c = Field.nullable("c" + round, int64); + final Schema schema = new Schema(Lists.newArrayList(a, b, c)); + final CyclicBarrier gate = new CyclicBarrier(numThreads); + + List> futures = new ArrayList<>(); + for (int i = 0; i < numThreads; i++) { + futures.add( executors.submit( () -> { + Projector evaluator = null; try { - Projector evaluator = + gate.await(30, java.util.concurrent.TimeUnit.SECONDS); + evaluator = configOptions == null - ? Projector.make(schemas.get((int) (Math.random() * 100)), exprs) - : Projector.make( - schemas.get((int) (Math.random() * 100)), exprs, configOptions); - evaluator.close(); - } catch (GandivaException e) { - e.printStackTrace(); + ? Projector.make(schema, exprs) + : Projector.make(schema, exprs, configOptions); + assertNotNull(evaluator); + } catch (Throwable t) { + failures.add(t); + } finally { + if (evaluator != null) { + try { + evaluator.close(); + } catch (Throwable t) { + failures.add(t); + } + } } - }); - }); - executors.shutdown(); - executors.awaitTermination(100, java.util.concurrent.TimeUnit.SECONDS); + })); + } + for (Future future : futures) { + future.get(120, java.util.concurrent.TimeUnit.SECONDS); + } + } + } finally { + executors.shutdownNow(); + } + + if (!failures.isEmpty()) { + AssertionError error = + new AssertionError( + failures.size() + + " concurrent same-cache-key Projector.make() calls failed across " + + numRounds + + " rounds of " + + numThreads + + " threads; first failure attached"); + error.initCause(failures.get(0)); + throw error; + } } @Test @@ -156,6 +295,14 @@ public void testMakeProjectorParallel() throws Exception { new ConfigurationBuilder.ConfigOptions().withTargetCPU(false).withOptimize(false)); } + @Test + public void testMakeProjectorParallelSameKey() throws Exception { + // Only the default configuration. The configuration is part of the native cache key, so varying + // it just builds unrelated cache entries -- it adds runtime without exercising anything new in + // the same-key race this test is about. + testMakeProjectorParallelSameKey(null); + } + // Will be fixed by https://issues.apache.org/jira/browse/ARROW-4371 @Disabled @Test