This is an automated email from the ASF dual-hosted git repository. dlmarion pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/accumulo.git
commit 82a141b3a37dc944a8f5dfce8ed6db9d0274f3d3 Merge: 20aeef52e8 0c814bace6 Author: Dave Marion <[email protected]> AuthorDate: Mon Jul 20 21:30:53 2026 +0000 Merge branch '2.1' .../apache/accumulo/core/file/FileOperations.java | 32 ++++++++++++++ .../file/blockfile/impl/CachableBlockFile.java | 51 ++++++---------------- .../core/file/rfile/bcfile/PrintBCInfo.java | 14 ++++-- .../server/compaction/CompactionPluginUtils.java | 8 ++-- .../apache/accumulo/server/fs/VolumeManager.java | 3 ++ .../accumulo/server/fs/VolumeManagerImpl.java | 6 +++ .../apache/accumulo/tserver/logger/LogReader.java | 8 ++-- 7 files changed, 74 insertions(+), 48 deletions(-) diff --cc core/src/main/java/org/apache/accumulo/core/file/FileOperations.java index 1ea4fb1384,e3d4d492cd..7afc7b5577 --- a/core/src/main/java/org/apache/accumulo/core/file/FileOperations.java +++ b/core/src/main/java/org/apache/accumulo/core/file/FileOperations.java @@@ -32,10 -36,10 +35,11 @@@ import org.apache.accumulo.core.data.Ra import org.apache.accumulo.core.data.TableId; import org.apache.accumulo.core.file.blockfile.impl.CacheProvider; import org.apache.accumulo.core.file.rfile.RFile; +import org.apache.accumulo.core.metadata.TabletFile; +import org.apache.accumulo.core.metadata.UnreferencedTabletFile; import org.apache.accumulo.core.spi.crypto.CryptoService; -import org.apache.accumulo.core.util.ratelimit.RateLimiter; import org.apache.hadoop.conf.Configuration; + import org.apache.hadoop.fs.FSDataInputStream; import org.apache.hadoop.fs.FSDataOutputStream; import org.apache.hadoop.fs.FileStatus; import org.apache.hadoop.fs.FileSystem; diff --cc core/src/main/java/org/apache/accumulo/core/file/blockfile/impl/CachableBlockFile.java index 6a6412e303,ba1c305bc5..e445da175a --- a/core/src/main/java/org/apache/accumulo/core/file/blockfile/impl/CachableBlockFile.java +++ b/core/src/main/java/org/apache/accumulo/core/file/blockfile/impl/CachableBlockFile.java @@@ -26,9 -26,7 +26,6 @@@ import java.io.UncheckedIOException import java.util.Collections; import java.util.Map; import java.util.Objects; - import java.util.concurrent.CancellationException; - import java.util.concurrent.CompletableFuture; --import java.util.concurrent.ExecutionException; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Supplier; @@@ -47,8 -48,8 +45,7 @@@ import org.apache.hadoop.conf.Configura import org.apache.hadoop.fs.FSDataInputStream; import org.apache.hadoop.fs.FileStatus; import org.apache.hadoop.fs.FileSystem; - import org.apache.hadoop.fs.FutureDataInputStreamBuilder; import org.apache.hadoop.fs.Path; -import org.apache.hadoop.fs.Seekable; import org.slf4j.Logger; import org.slf4j.LoggerFactory; diff --cc core/src/main/java/org/apache/accumulo/core/file/rfile/bcfile/PrintBCInfo.java index 3ed0d5b6ba,cf26aadac8..1923d22121 --- a/core/src/main/java/org/apache/accumulo/core/file/rfile/bcfile/PrintBCInfo.java +++ b/core/src/main/java/org/apache/accumulo/core/file/rfile/bcfile/PrintBCInfo.java @@@ -23,7 -23,9 +23,8 @@@ import java.io.PrintStream import java.util.Map.Entry; import java.util.Set; -import org.apache.accumulo.core.cli.ConfigOpts; -import org.apache.accumulo.core.conf.SiteConfiguration; +import org.apache.accumulo.core.cli.ClientOpts; + import org.apache.accumulo.core.file.FileOperations; import org.apache.accumulo.core.file.rfile.bcfile.BCFile.MetaIndexEntry; import org.apache.accumulo.core.spi.crypto.CryptoService; import org.apache.accumulo.core.spi.crypto.NoCryptoServiceFactory; @@@ -44,23 -48,11 +46,25 @@@ public class PrintBCInfo Path path; CryptoService cryptoService = NoCryptoServiceFactory.NONE; + public PrintBCInfo(String[] args) throws Exception { + BCInfoOpts opts = new BCInfoOpts(); + opts.parseArgs(getClass().getSimpleName(), args); + conf = new Configuration(); + FileSystem hadoopFs = FileSystem.get(conf); + FileSystem localFs = FileSystem.getLocal(conf); + path = new Path(opts.file); + if (opts.file.contains(":")) { + fs = path.getFileSystem(conf); + } else { + fs = hadoopFs.exists(path) ? hadoopFs : localFs; // fall back to local + } + } + public void printMetaBlockInfo() throws IOException { - try (FSDataInputStream fsin = fs.open(path); BCFile.Reader bcfr = - new BCFile.Reader(fsin, fs.getFileStatus(path).getLen(), conf, cryptoService)) { + FileStatus status = fs.getFileStatus(path); + + try (FSDataInputStream fsin = FileOperations.openFile(fs, path, status); + BCFile.Reader bcfr = new BCFile.Reader(fsin, status.getLen(), conf, cryptoService)) { Set<Entry<String,MetaIndexEntry>> es = bcfr.metaIndex.index.entrySet(); @@@ -78,14 -70,27 +82,16 @@@ } } - static class Opts extends ConfigOpts { - @Parameter(description = " <file>") - String file; - } + public String getCompressionType() throws IOException { - try (FSDataInputStream fsin = fs.open(path); BCFile.Reader bcfr = - new BCFile.Reader(fsin, fs.getFileStatus(path).getLen(), conf, cryptoService)) { ++ FileStatus status = fs.getFileStatus(path); + - public PrintBCInfo(String[] args) throws Exception { - Opts opts = new Opts(); - opts.parseArgs("PrintInfo", args); - if (opts.file.isEmpty()) { - System.err.println("No files were given"); - System.exit(-1); - } - siteConfig = opts.getSiteConfiguration(); - conf = new Configuration(); - FileSystem hadoopFs = FileSystem.get(conf); - FileSystem localFs = FileSystem.getLocal(conf); - path = new Path(opts.file); - if (opts.file.contains(":")) { - fs = path.getFileSystem(conf); - } else { - fs = hadoopFs.exists(path) ? hadoopFs : localFs; // fall back to local ++ try (FSDataInputStream fsin = FileOperations.openFile(fs, path, status); ++ BCFile.Reader bcfr = new BCFile.Reader(fsin, status.getLen(), conf, cryptoService)) { + + Set<Entry<String,MetaIndexEntry>> es = bcfr.metaIndex.index.entrySet(); + + return es.stream().filter(entry -> entry.getKey().equals("RFile.index")).findFirst() + .map(entry -> entry.getValue().getCompressionAlgorithm().getName()).orElse(null); } } diff --cc server/base/src/main/java/org/apache/accumulo/server/compaction/CompactionPluginUtils.java index 514a110f2a,0000000000..6f0d415e97 mode 100644,000000..100644 --- a/server/base/src/main/java/org/apache/accumulo/server/compaction/CompactionPluginUtils.java +++ b/server/base/src/main/java/org/apache/accumulo/server/compaction/CompactionPluginUtils.java @@@ -1,383 -1,0 +1,385 @@@ +/* + * 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 + * + * https://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.accumulo.server.compaction; + +import java.io.IOException; +import java.io.UncheckedIOException; +import java.net.URI; +import java.util.Collection; +import java.util.Collections; +import java.util.HashMap; +import java.util.HashSet; +import java.util.Map; +import java.util.Optional; +import java.util.Set; +import java.util.function.Predicate; +import java.util.function.Supplier; +import java.util.stream.Collectors; + +import org.apache.accumulo.core.classloader.ClassLoaderUtil; +import org.apache.accumulo.core.client.PluginEnvironment; +import org.apache.accumulo.core.client.admin.CompactionConfig; +import org.apache.accumulo.core.client.admin.PluginConfig; +import org.apache.accumulo.core.client.admin.compaction.CompactableFile; +import org.apache.accumulo.core.client.admin.compaction.CompactionConfigurer; +import org.apache.accumulo.core.client.admin.compaction.CompactionSelector; +import org.apache.accumulo.core.client.rfile.RFileSource; +import org.apache.accumulo.core.client.sample.SamplerConfiguration; +import org.apache.accumulo.core.client.summary.SummarizerConfiguration; +import org.apache.accumulo.core.client.summary.Summary; +import org.apache.accumulo.core.clientImpl.UserCompactionUtils; +import org.apache.accumulo.core.conf.AccumuloConfiguration; +import org.apache.accumulo.core.conf.ConfigurationTypeHelper; +import org.apache.accumulo.core.conf.Property; +import org.apache.accumulo.core.data.AbstractId; +import org.apache.accumulo.core.data.Key; +import org.apache.accumulo.core.data.RowRange; +import org.apache.accumulo.core.data.TableId; +import org.apache.accumulo.core.data.TabletId; +import org.apache.accumulo.core.data.Value; +import org.apache.accumulo.core.dataImpl.KeyExtent; +import org.apache.accumulo.core.dataImpl.TabletIdImpl; +import org.apache.accumulo.core.file.FileOperations; +import org.apache.accumulo.core.iterators.SortedKeyValueIterator; +import org.apache.accumulo.core.metadata.CompactableFileImpl; +import org.apache.accumulo.core.metadata.ReferencedTabletFile; +import org.apache.accumulo.core.metadata.StoredTabletFile; +import org.apache.accumulo.core.metadata.schema.DataFileValue; +import org.apache.accumulo.core.sample.impl.SamplerConfigurationImpl; +import org.apache.accumulo.core.spi.common.ServiceEnvironment; +import org.apache.accumulo.core.spi.compaction.CompactionDispatcher; +import org.apache.accumulo.core.spi.compaction.CompactionPlanner; +import org.apache.accumulo.core.spi.compaction.CompactionServiceId; +import org.apache.accumulo.core.summary.SummarizerFactory; +import org.apache.accumulo.core.summary.SummaryCollection; +import org.apache.accumulo.core.summary.SummaryReader; +import org.apache.accumulo.core.util.compaction.CompactionPlannerInitParams; +import org.apache.accumulo.core.util.compaction.CompactionServicesConfig; +import org.apache.accumulo.server.ServerContext; +import org.apache.accumulo.server.ServiceEnvironmentImpl; +import org.apache.accumulo.server.tablets.TabletNameGenerator; +import org.apache.hadoop.conf.Configuration; - import org.apache.hadoop.fs.FSDataInputStream; ++import org.apache.hadoop.fs.FileStatus; +import org.apache.hadoop.fs.FileSystem; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.google.common.collect.Collections2; + +public class CompactionPluginUtils { + + private static final Logger log = LoggerFactory.getLogger(CompactionPluginUtils.class); + + private static <T> T newInstance(AccumuloConfiguration tableConfig, String className, + Class<T> baseClass) { + String context = ClassLoaderUtil.tableContext(tableConfig); + try { + return ConfigurationTypeHelper.getClassInstance(context, className, baseClass); + } catch (ReflectiveOperationException e) { + throw new IllegalArgumentException(e); + } + } + + public static Set<StoredTabletFile> selectFiles(ServerContext context, KeyExtent extent, + CompactionConfig compactionConfig, Map<StoredTabletFile,DataFileValue> allFiles) { + if (!UserCompactionUtils.isDefault(compactionConfig.getSelector())) { + return selectFiles(context, extent, allFiles, compactionConfig.getSelector()); + } else { + return allFiles.keySet(); + } + } + + private static Set<StoredTabletFile> selectFiles(ServerContext context, KeyExtent extent, + Map<StoredTabletFile,DataFileValue> datafiles, PluginConfig selectorConfig) { + + log.debug("Selecting files for {} using {}", extent, selectorConfig); + + CompactionSelector selector = newInstance(context.getTableConfiguration(extent.tableId()), + selectorConfig.getClassName(), CompactionSelector.class); + + final ServiceEnvironment senv = new ServiceEnvironmentImpl(context); + + selector.init(new CompactionSelector.InitParameters() { + @Override + public Map<String,String> getOptions() { + return selectorConfig.getOptions(); + } + + @Override + public PluginEnvironment getEnvironment() { + return senv; + } + + @Override + public TableId getTableId() { + return extent.tableId(); + } + }); + + CompactionSelector.Selection selection = + selector.select(new CompactionSelector.SelectionParameters() { + @Override + public PluginEnvironment getEnvironment() { + return senv; + } + + @Override + public Collection<CompactableFile> getAvailableFiles() { + return Collections2.transform(datafiles.entrySet(), + e -> new CompactableFileImpl(e.getKey(), e.getValue())); + } + + @Override + public Collection<Summary> getSummaries(Collection<CompactableFile> files, + Predicate<SummarizerConfiguration> summarySelector) { + + // ELASTICITY_TODO this may open files for user tables in the manager, need to avoid + // this. See #3526 + + try { + var tableConf = context.getTableConfiguration(extent.tableId()); + + SummaryCollection sc = new SummaryCollection(); + SummarizerFactory factory = new SummarizerFactory(tableConf); + for (CompactableFile cf : files) { + var file = CompactableFileImpl.toStoredTabletFile(cf); + FileSystem fs = context.getVolumeManager().getFileSystemByPath(file.getPath()); ++ FileStatus status = fs.getFileStatus(file.getPath()); + Configuration conf = context.getHadoopConf(); - RFileSource source = new RFileSource(new FSDataInputStream(fs.open(file.getPath())), - fs.getFileStatus(file.getPath()).getLen(), file.getRange()); ++ RFileSource source = ++ new RFileSource(FileOperations.openFile(fs, file.getPath(), status), ++ status.getLen(), file.getRange()); + + SummaryCollection fsc = SummaryReader + .load(conf, source, file.getFileName(), summarySelector, factory, + tableConf.getCryptoService()) + .getSummaries(Collections.singletonList( + RowRange.range(extent.prevEndRow(), false, extent.endRow(), true))); + + sc.merge(fsc, factory); + } + return sc.getSummaries(); + } catch (IOException ioe) { + throw new UncheckedIOException(ioe); + } + } + + @Override + public TableId getTableId() { + return extent.tableId(); + } + + @Override + public TabletId getTabletId() { + return new TabletIdImpl(extent); + } + + @Override + public Optional<SortedKeyValueIterator<Key,Value>> getSample(CompactableFile cf, + SamplerConfiguration sc) { + + // ELASTICITY_TODO this may open files for user tables in the manager, need to avoid + // this. See #3526 + + try { + var file = CompactableFileImpl.toStoredTabletFile(cf); + FileSystem fs = context.getVolumeManager().getFileSystemByPath(file.getPath()); + Configuration conf = context.getHadoopConf(); + var tableConf = context.getTableConfiguration(extent.tableId()); + var iter = FileOperations.getInstance().newReaderBuilder() + .forFile(file, fs, conf, tableConf.getCryptoService()) + .withTableConfiguration(tableConf).seekToBeginning().build(); + var sampleIter = iter.getSample(new SamplerConfigurationImpl(sc)); + if (sampleIter == null) { + iter.close(); + return Optional.empty(); + } + + return Optional.of(sampleIter); + } catch (IOException ioe) { + throw new UncheckedIOException(ioe); + } + } + }); + + return selection.getFilesToCompact().stream().map(CompactableFileImpl::toStoredTabletFile) + .collect(Collectors.toSet()); + } + + public static Map<String,String> computeOverrides(Optional<CompactionConfig> compactionConfig, + ServerContext context, KeyExtent extent, Set<CompactableFile> inputFiles, + Supplier<Set<CompactableFile>> selectedFiles, ReferencedTabletFile outputFile) { + + if (compactionConfig.isPresent() + && !UserCompactionUtils.isDefault(compactionConfig.orElseThrow().getConfigurer())) { + return CompactionPluginUtils.computeOverrides(context, extent, inputFiles, selectedFiles, + compactionConfig.orElseThrow().getConfigurer(), outputFile); + } + + var tableConf = context.getTableConfiguration(extent.tableId()); + + var configurorClass = tableConf.get(Property.TABLE_COMPACTION_CONFIGURER); + if (configurorClass == null || configurorClass.isBlank()) { + return Map.of(); + } + + var opts = + tableConf.getAllPropertiesWithPrefixStripped(Property.TABLE_COMPACTION_CONFIGURER_OPTS); + + return CompactionPluginUtils.computeOverrides(context, extent, inputFiles, selectedFiles, + new PluginConfig(configurorClass, opts), outputFile); + } + + public static Map<String,String> computeOverrides(ServerContext context, KeyExtent extent, + Set<CompactableFile> inputFiles, Supplier<Set<CompactableFile>> selectedFiles, + PluginConfig cfg, ReferencedTabletFile outputFile) { + + CompactionConfigurer configurer = newInstance(context.getTableConfiguration(extent.tableId()), + cfg.getClassName(), CompactionConfigurer.class); + + final ServiceEnvironment senv = new ServiceEnvironmentImpl(context); + + configurer.init(new CompactionConfigurer.InitParameters() { + @Override + public Map<String,String> getOptions() { + return cfg.getOptions(); + } + + @Override + public PluginEnvironment getEnvironment() { + return senv; + } + + @Override + public TableId getTableId() { + return extent.tableId(); + } + }); + + var overrides = configurer.override(new CompactionConfigurer.InputParameters() { + @Override + public Collection<CompactableFile> getInputFiles() { + return inputFiles; + } + + @Override + public Set<CompactableFile> getSelectedFiles() { + return selectedFiles.get(); + } + + @Override + public PluginEnvironment getEnvironment() { + return senv; + } + + @Override + public TableId getTableId() { + return extent.tableId(); + } + + @Override + public TabletId getTabletId() { + return new TabletIdImpl(extent); + } + + @Override + public URI getOutputFile() { + return TabletNameGenerator.computeCompactionFileDest(outputFile).getPath().toUri(); + } + }); + + if (overrides.getOverrides().isEmpty()) { + return null; + } + + return overrides.getOverrides(); + } + + static CompactionDispatcher createDispatcher(ServiceEnvironment env, TableId tableId) { + + var conf = env.getConfiguration(tableId); + + var className = conf.get(Property.TABLE_COMPACTION_DISPATCHER.getKey()); + + Map<String,String> opts = new HashMap<>(); + + conf.getWithPrefix(Property.TABLE_COMPACTION_DISPATCHER_OPTS.getKey()).forEach((k, v) -> { + opts.put(k.substring(Property.TABLE_COMPACTION_DISPATCHER_OPTS.getKey().length()), v); + }); + + var finalOpts = Collections.unmodifiableMap(opts); + + CompactionDispatcher.InitParameters initParameters = new CompactionDispatcher.InitParameters() { + @Override + public Map<String,String> getOptions() { + return finalOpts; + } + + @Override + public TableId getTableId() { + return tableId; + } + + @Override + public ServiceEnvironment getServiceEnv() { + return env; + } + }; + + CompactionDispatcher dispatcher = null; + try { + dispatcher = env.instantiate(tableId, className, CompactionDispatcher.class); + } catch (ReflectiveOperationException e) { + throw new RuntimeException(e); + } + + dispatcher.init(initParameters); + + return dispatcher; + } + + /** + * Inspect configuration and determines what resource groups are configured for compaction. + */ + public static Set<String> getConfiguredCompactionResourceGroups(ServerContext ctx) + throws ReflectiveOperationException { + + Set<String> groups = new HashSet<>(); + AccumuloConfiguration config = ctx.getConfiguration(); + CompactionServicesConfig servicesConfig = new CompactionServicesConfig(config); + + for (var entry : servicesConfig.getPlanners().entrySet()) { + String serviceId = entry.getKey(); + String plannerClassName = entry.getValue(); + + Class<? extends CompactionPlanner> plannerClass = + Class.forName(plannerClassName).asSubclass(CompactionPlanner.class); + CompactionPlanner planner = plannerClass.getDeclaredConstructor().newInstance(); + + var initParams = new CompactionPlannerInitParams(CompactionServiceId.of(serviceId), + servicesConfig.getPlannerPrefix(serviceId), servicesConfig.getOptions().get(serviceId), + new ServiceEnvironmentImpl(ctx)); + + planner.init(initParams); + + initParams.getRequestedGroups().stream().map(AbstractId::canonical).forEach(groups::add); + } + return groups; + } +} diff --cc server/base/src/main/java/org/apache/accumulo/server/fs/VolumeManagerImpl.java index 636a52b6c5,211e852915..7609bb08c4 --- a/server/base/src/main/java/org/apache/accumulo/server/fs/VolumeManagerImpl.java +++ b/server/base/src/main/java/org/apache/accumulo/server/fs/VolumeManagerImpl.java @@@ -48,11 -44,10 +48,12 @@@ import java.util.stream.Stream import org.apache.accumulo.core.conf.AccumuloConfiguration; import org.apache.accumulo.core.conf.DefaultConfiguration; import org.apache.accumulo.core.conf.Property; +import org.apache.accumulo.core.fate.FateId; + import org.apache.accumulo.core.file.FileOperations; import org.apache.accumulo.core.spi.fs.VolumeChooser; import org.apache.accumulo.core.util.Pair; -import org.apache.accumulo.core.util.threads.ThreadPools; +import org.apache.accumulo.core.util.cache.Caches; +import org.apache.accumulo.core.util.cache.Caches.CacheName; import org.apache.accumulo.core.volume.Volume; import org.apache.accumulo.core.volume.VolumeConfiguration; import org.apache.accumulo.core.volume.VolumeImpl; diff --cc server/tserver/src/main/java/org/apache/accumulo/tserver/logger/LogReader.java index 407a8f5939,704a57804c..c41455850f --- a/server/tserver/src/main/java/org/apache/accumulo/tserver/logger/LogReader.java +++ b/server/tserver/src/main/java/org/apache/accumulo/tserver/logger/LogReader.java @@@ -52,9 -48,8 +52,10 @@@ import org.apache.accumulo.start.spi.Ke import org.apache.accumulo.tserver.log.DfsLogger; import org.apache.accumulo.tserver.log.DfsLogger.LogHeaderIncompleteException; import org.apache.accumulo.tserver.log.RecoveryLogsIterator; +import org.apache.accumulo.tserver.log.ResolvedSortedLog; +import org.apache.accumulo.tserver.logger.LogReader.ReaderOpts; import org.apache.hadoop.fs.FSDataInputStream; + import org.apache.hadoop.fs.FileStatus; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.Text; import org.slf4j.Logger;
