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
19 changes: 9 additions & 10 deletions core/src/main/java/org/apache/accumulo/core/conf/Property.java
Original file line number Diff line number Diff line change
Expand Up @@ -2065,16 +2065,15 @@ public static boolean isValidTablePropertyKey(String key) {
TSERV_WALOG_TOLERATED_CREATION_FAILURES, TSERV_WAL_TOLERATED_WAIT_INCREMENT,
TSERV_WALOG_TOLERATED_WAIT_INCREMENT, TSERV_WAL_TOLERATED_MAXIMUM_WAIT_DURATION,
TSERV_WALOG_TOLERATED_MAXIMUM_WAIT_DURATION, TSERV_MAX_IDLE, TSERV_SESSION_MAXIDLE,
TSERV_SCAN_RESULTS_MAX_TIMEOUT, TSERV_MAJC_DELAY, TSERV_COMPACTION_SERVICE_DEFAULT_MAX_OPEN,
TSERV_MAJC_THREAD_MAXOPEN, TSERV_MAJC_MAXCONCURRENT, TSERV_MAJC_THROUGHPUT,
TSERV_MINC_MAXCONCURRENT, TSERV_MONITOR_FS, TSERV_HEALTH_CHECK_FREQ, TSERV_THREADCHECK,
TSERV_LOG_BUSY_TABLETS_COUNT, TSERV_LOG_BUSY_TABLETS_INTERVAL, TSERV_WAL_SORT_MAX_CONCURRENT,
TSERV_RECOVERY_MAX_CONCURRENT, TSERV_WORKQ_THREADS, TSERV_SLOW_FILEPERMIT_MILLIS,
TSERV_READ_AHEAD_MAXCONCURRENT, TSERV_METADATA_READ_AHEAD_MAXCONCURRENT, TSERV_WAL_BLOCKSIZE,
TSERV_CLIENTPORT, TSERV_PORTSEARCH, TSERV_MAX_MESSAGE_SIZE, TSERV_CACHE_MANAGER_IMPL,
TSERV_DATACACHE_SIZE, TSERV_INDEXCACHE_SIZE, TSERV_SUMMARYCACHE_SIZE, TSERV_DEFAULT_BLOCKSIZE,
TSERV_MINTHREADS, TSERV_MINTHREADS_TIMEOUT, TSERV_NATIVEMAP_ENABLED, TSERV_MAXMEM,
TSERV_SCAN_MAX_OPENFILES,
TSERV_MAJC_DELAY, TSERV_COMPACTION_SERVICE_DEFAULT_MAX_OPEN, TSERV_MAJC_THREAD_MAXOPEN,
TSERV_MAJC_MAXCONCURRENT, TSERV_MAJC_THROUGHPUT, TSERV_MINC_MAXCONCURRENT, TSERV_MONITOR_FS,
TSERV_HEALTH_CHECK_FREQ, TSERV_THREADCHECK, TSERV_LOG_BUSY_TABLETS_COUNT,
TSERV_LOG_BUSY_TABLETS_INTERVAL, TSERV_WAL_SORT_MAX_CONCURRENT, TSERV_RECOVERY_MAX_CONCURRENT,
TSERV_WORKQ_THREADS, TSERV_SLOW_FILEPERMIT_MILLIS, TSERV_READ_AHEAD_MAXCONCURRENT,
TSERV_METADATA_READ_AHEAD_MAXCONCURRENT, TSERV_WAL_BLOCKSIZE, TSERV_CLIENTPORT,
TSERV_PORTSEARCH, TSERV_MAX_MESSAGE_SIZE, TSERV_CACHE_MANAGER_IMPL, TSERV_DATACACHE_SIZE,
TSERV_INDEXCACHE_SIZE, TSERV_SUMMARYCACHE_SIZE, TSERV_DEFAULT_BLOCKSIZE, TSERV_MINTHREADS,
TSERV_MINTHREADS_TIMEOUT, TSERV_NATIVEMAP_ENABLED, TSERV_MAXMEM, TSERV_SCAN_MAX_OPENFILES,

// GC options
GC_CANDIDATE_BATCH_SIZE, GC_CYCLE_START, GC_PORT,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@

import java.io.IOException;
import java.nio.ByteBuffer;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Collections;
Expand Down Expand Up @@ -131,6 +132,8 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import com.google.common.base.Supplier;
import com.google.common.base.Suppliers;
import com.google.common.cache.Cache;

import io.opentelemetry.api.trace.Span;
Expand All @@ -139,7 +142,7 @@
public class TabletClientHandler implements TabletClientService.Iface {

private static final Logger log = LoggerFactory.getLogger(TabletClientHandler.class);
private final long MAX_TIME_TO_WAIT_FOR_SCAN_RESULT_MILLIS;
private final Supplier<Long> scanResultWaitTime;
private static final long RECENTLY_SPLIT_MILLIES = MINUTES.toMillis(1);
private final TabletServer server;
protected final TransactionWatcher watcher;
Expand All @@ -155,8 +158,12 @@ public TabletClientHandler(TabletServer server, TransactionWatcher watcher,
this.writeTracker = writeTracker;
this.security = context.getSecurityOperation();
this.server = server;
MAX_TIME_TO_WAIT_FOR_SCAN_RESULT_MILLIS = server.getContext().getConfiguration()
.getTimeInMillis(Property.TSERV_SCAN_RESULTS_MAX_TIMEOUT);
scanResultWaitTime =
Suppliers
.memoizeWithExpiration(
() -> server.getContext().getConfiguration()
.getTimeInMillis(Property.TSERV_SCAN_RESULTS_MAX_TIMEOUT),
Duration.ofSeconds(10));
log.debug("{} created", TabletClientHandler.class.getName());
}

Expand Down Expand Up @@ -1534,8 +1541,7 @@ public void removeLogs(TInfo tinfo, TCredentials credentials, List<String> filen

private TSummaries getSummaries(Future<SummaryCollection> future) throws TimeoutException {
try {
SummaryCollection sc =
future.get(MAX_TIME_TO_WAIT_FOR_SCAN_RESULT_MILLIS, TimeUnit.MILLISECONDS);
SummaryCollection sc = future.get(scanResultWaitTime.get(), TimeUnit.MILLISECONDS);
return sc.toThrift();
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@

import java.io.IOException;
import java.nio.ByteBuffer;
import java.time.Duration;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
Expand Down Expand Up @@ -88,6 +89,8 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import com.google.common.base.Supplier;
import com.google.common.base.Suppliers;
import com.google.common.collect.Collections2;

public class ThriftScanClientHandler implements TabletScanClientService.Iface {
Expand All @@ -98,15 +101,19 @@ public class ThriftScanClientHandler implements TabletScanClientService.Iface {
protected final ServerContext context;
protected final AuditedSecurityOperation security;
private final WriteTracker writeTracker;
private final long MAX_TIME_TO_WAIT_FOR_SCAN_RESULT_MILLIS;
private final Supplier<Long> scanResultWaitTime;

public ThriftScanClientHandler(TabletHostingServer server, WriteTracker writeTracker) {
this.server = server;
this.context = server.getContext();
this.writeTracker = writeTracker;
this.security = context.getSecurityOperation();
MAX_TIME_TO_WAIT_FOR_SCAN_RESULT_MILLIS = server.getContext().getConfiguration()
.getTimeInMillis(Property.TSERV_SCAN_RESULTS_MAX_TIMEOUT);
scanResultWaitTime =
Suppliers
.memoizeWithExpiration(
() -> server.getContext().getConfiguration()
.getTimeInMillis(Property.TSERV_SCAN_RESULTS_MAX_TIMEOUT),
Duration.ofSeconds(10));
}

public NamespaceId getNamespaceId(TCredentials credentials, TableId tableId)
Expand Down Expand Up @@ -264,7 +271,7 @@ protected ScanResult continueScan(TInfo tinfo, long scanID, SingleScanSession sc

ScanBatch bresult;
try {
bresult = scanSession.getScanTask().get(busyTimeout, MAX_TIME_TO_WAIT_FOR_SCAN_RESULT_MILLIS,
bresult = scanSession.getScanTask().get(busyTimeout, scanResultWaitTime.get(),
TimeUnit.MILLISECONDS);
scanSession.clearScanTask();
} catch (ExecutionException e) {
Expand All @@ -277,7 +284,7 @@ protected ScanResult continueScan(TInfo tinfo, long scanID, SingleScanSession sc
} else if (e.getCause() instanceof SampleNotPresentException) {
throw new TSampleNotPresentException(scanSession.extent.toThrift());
} else if (e.getCause() instanceof IOException) {
sleepUninterruptibly(MAX_TIME_TO_WAIT_FOR_SCAN_RESULT_MILLIS, TimeUnit.MILLISECONDS);
sleepUninterruptibly(scanResultWaitTime.get(), TimeUnit.MILLISECONDS);
List<KVEntry> empty = Collections.emptyList();
bresult = new ScanBatch(empty, true);
scanSession.clearScanTask();
Expand Down Expand Up @@ -482,8 +489,8 @@ private MultiScanResult continueMultiScan(long scanID, MultiScanSession session,

try {

MultiScanResult scanResult = session.getScanTask().get(busyTimeout,
MAX_TIME_TO_WAIT_FOR_SCAN_RESULT_MILLIS, TimeUnit.MILLISECONDS);
MultiScanResult scanResult =
session.getScanTask().get(busyTimeout, scanResultWaitTime.get(), TimeUnit.MILLISECONDS);
session.clearScanTask();
return scanResult;
} catch (ExecutionException e) {
Expand Down
Loading