From dac49dc595ef1988a78ee52270a47c322150f4c2 Mon Sep 17 00:00:00 2001 From: Bill Slacum Date: Fri, 10 Jul 2026 13:32:49 -0400 Subject: [PATCH] Make parallel scan thread count configurable Replace the hardcoded single-thread BatchScanner in ParallelSnapshotScanner with a configurable value via the new fluo.scan.num.threads property. Add setNumParallelScanThreads and getNumParallelScanThreads to FluoConfiguration (default 1 thread) with accompanying Javadocs. # Conflicts: # modules/core/src/main/java/org/apache/fluo/core/impl/ParallelSnapshotScanner.java --- .../fluo/api/config/FluoConfiguration.java | 26 +++++++++++++++++++ .../core/impl/ParallelSnapshotScanner.java | 6 ++--- 2 files changed, 29 insertions(+), 3 deletions(-) diff --git a/modules/api/src/main/java/org/apache/fluo/api/config/FluoConfiguration.java b/modules/api/src/main/java/org/apache/fluo/api/config/FluoConfiguration.java index b18b21c3..490709a5 100644 --- a/modules/api/src/main/java/org/apache/fluo/api/config/FluoConfiguration.java +++ b/modules/api/src/main/java/org/apache/fluo/api/config/FluoConfiguration.java @@ -188,6 +188,9 @@ public class FluoConfiguration extends SimpleConfiguration { public static final String WORKER_NUM_THREADS_PROP = WORKER_PREFIX + ".num.threads"; public static final int WORKER_NUM_THREADS_DEFAULT = 10; + public static final String PARALLEL_SCAN_THREADS_PROP = FLUO_PREFIX + ".scan.num.threads"; + public static final int NUM_PARALLEL_SCAN_THREADS_DEFAULT = 1; + // Loader properties private static final String LOADER_PREFIX = FLUO_PREFIX + ".loader"; public static final String LOADER_NUM_THREADS_PROP = LOADER_PREFIX + ".num.threads"; @@ -634,6 +637,29 @@ public int getWorkerThreads() { return getPositiveInt(WORKER_NUM_THREADS_PROP, WORKER_NUM_THREADS_DEFAULT); } + /** + * Sets the number of threads used to execute parallel scans, must be positive. The default is + * {@value #NUM_PARALLEL_SCAN_THREADS_DEFAULT} thread. Sets this value in the property + * {@value #PARALLEL_SCAN_THREADS_PROP} + * + * @param numThreads The number of threads to use, must be positive + * @since 2.0.0 + */ + public FluoConfiguration setNumParallelScanThreads(int numThreads) { + return setPositiveInt(PARALLEL_SCAN_THREADS_PROP, numThreads); + } + + /** + * Gets the value of the property {@value #PARALLEL_SCAN_THREADS_PROP} if set otherwise returns + * {@value #NUM_PARALLEL_SCAN_THREADS_DEFAULT} + * + * @return The number of threads used for parallel scans. + * @since 2.0.0 + */ + public int getNumParallelScanThreads() { + return getPositiveInt(PARALLEL_SCAN_THREADS_PROP, NUM_PARALLEL_SCAN_THREADS_DEFAULT); + } + /** * @deprecated since 1.1.0. Replaced by {@link #setObserverProvider(String)} and * {@link #getObserverProvider()} diff --git a/modules/core/src/main/java/org/apache/fluo/core/impl/ParallelSnapshotScanner.java b/modules/core/src/main/java/org/apache/fluo/core/impl/ParallelSnapshotScanner.java index 17f4620f..d11d776e 100644 --- a/modules/core/src/main/java/org/apache/fluo/core/impl/ParallelSnapshotScanner.java +++ b/modules/core/src/main/java/org/apache/fluo/core/impl/ParallelSnapshotScanner.java @@ -116,9 +116,9 @@ private BatchScanner setupBatchScanner() { BatchScanner scanner; try { - // TODO hardcoded number of threads! - // one thread is probably good.. going for throughput - scanner = env.getAccumuloClient().createBatchScanner(env.getTable(), this.authorizations, 1); + scanner = + env.getAccumuloClient().createBatchScanner(env.getTable(), env.getAuthorizations(), + env.getConfiguration().getNumParallelScanThreads()); } catch (TableNotFoundException e) { throw new RuntimeException(e); }