From 21eb905986115b705a4743e61823e927b39d1ae1 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Fri, 7 Oct 2011 21:56:16 +0000 Subject: [PATCH 1/5] r/m unused code git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1180260 13f79535-47bb-0310-9956-ffa450edef68 --- .../concurrent/AIOExecutorService.java | 309 ------------------ 1 file changed, 309 deletions(-) delete mode 100644 src/java/org/apache/cassandra/concurrent/AIOExecutorService.java diff --git a/src/java/org/apache/cassandra/concurrent/AIOExecutorService.java b/src/java/org/apache/cassandra/concurrent/AIOExecutorService.java deleted file mode 100644 index f2655649ef..0000000000 --- a/src/java/org/apache/cassandra/concurrent/AIOExecutorService.java +++ /dev/null @@ -1,309 +0,0 @@ -/** - * 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.cassandra.concurrent; - -import java.util.Collection; -import java.util.List; -import java.util.concurrent.*; - -public class AIOExecutorService implements ExecutorService -{ - private ExecutorService executorService_; - - public AIOExecutorService(int corePoolSize, - int maximumPoolSize, - long keepAliveTime, - TimeUnit unit, - BlockingQueue workQueue, - ThreadFactory threadFactory) - { - executorService_ = new ThreadPoolExecutor(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, threadFactory); - } - - /** - * Executes the given command at some time in the future. The command - * may execute in a new thread, in a pooled thread, or in the calling - * thread, at the discretion of the Executor implementation. - * - * @param command the runnable task - * @throws RejectedExecutionException if this task cannot be - * accepted for execution. - * @throws NullPointerException if command is null - */ - public void execute(Runnable command) - { - executorService_.execute(command); - } - - /** - * Initiates an orderly shutdown in which previously submitted - * tasks are executed, but no new tasks will be accepted. - * Invocation has no additional effect if already shut down. - * - *

This method does not wait for previously submitted tasks to - * complete execution. Use {@link #awaitTermination awaitTermination} - * to do that. - * - * @throws SecurityException if a security manager exists and - * shutting down this ExecutorService may manipulate - * threads that the caller is not permitted to modify - * because it does not hold {@link - * java.lang.RuntimePermission}("modifyThread"), - * or the security manager's checkAccess method - * denies access. - */ - public void shutdown() - { - /* This is a noop. */ - } - - /** - * Attempts to stop all actively executing tasks, halts the - * processing of waiting tasks, and returns a list of the tasks - * that were awaiting execution. - * - *

This method does not wait for actively executing tasks to - * terminate. Use {@link #awaitTermination awaitTermination} to - * do that. - * - *

There are no guarantees beyond best-effort attempts to stop - * processing actively executing tasks. For example, typical - * implementations will cancel via {@link Thread#interrupt}, so any - * task that fails to respond to interrupts may never terminate. - * - * @return list of tasks that never commenced execution - * @throws SecurityException if a security manager exists and - * shutting down this ExecutorService may manipulate - * threads that the caller is not permitted to modify - * because it does not hold {@link - * java.lang.RuntimePermission}("modifyThread"), - * or the security manager's checkAccess method - * denies access. - */ - public List shutdownNow() - { - return executorService_.shutdownNow(); - } - - /** - * Returns true if this executor has been shut down. - * - * @return true if this executor has been shut down - */ - public boolean isShutdown() - { - return executorService_.isShutdown(); - } - - /** - * Returns true if all tasks have completed following shut down. - * Note that isTerminated is never true unless - * either shutdown or shutdownNow was called first. - * - * @return true if all tasks have completed following shut down - */ - public boolean isTerminated() - { - return executorService_.isTerminated(); - } - - /** - * Blocks until all tasks have completed execution after a shutdown - * request, or the timeout occurs, or the current thread is - * interrupted, whichever happens first. - * - * @param timeout the maximum time to wait - * @param unit the time unit of the timeout argument - * @return true if this executor terminated and - * false if the timeout elapsed before termination - * @throws InterruptedException if interrupted while waiting - */ - public boolean awaitTermination(long timeout, TimeUnit unit) throws InterruptedException - { - return executorService_.awaitTermination(timeout, unit); - } - - /** - * Submits a value-returning task for execution and returns a - * Future representing the pending results of the task. The - * Future's get method will return the task's result upon - * successful completion. - * - *

- * If you would like to immediately block waiting - * for a task, you can use constructions of the form - * result = exec.submit(aCallable).get(); - * - *

Note: The {@link Executors} class includes a set of methods - * that can convert some other common closure-like objects, - * for example, {@link java.security.PrivilegedAction} to - * {@link Callable} form so they can be submitted. - * - * @param task the task to submit - * @return a Future representing pending completion of the task - * @throws RejectedExecutionException if the task cannot be - * scheduled for execution - * @throws NullPointerException if the task is null - */ - public Future submit(Callable task) - { - return executorService_.submit(task); - } - - /** - * Submits a Runnable task for execution and returns a Future - * representing that task. The Future's get method will - * return the given result upon successful completion. - * - * @param task the task to submit - * @param result the result to return - * @return a Future representing pending completion of the task - * @throws RejectedExecutionException if the task cannot be - * scheduled for execution - * @throws NullPointerException if the task is null - */ - public Future submit(Runnable task, T result) - { - return executorService_.submit(task, result); - } - - /** - * Submits a Runnable task for execution and returns a Future - * representing that task. The Future's get method will - * return null upon successful completion. - * - * @param task the task to submit - * @return a Future representing pending completion of the task - * @throws RejectedExecutionException if the task cannot be - * scheduled for execution - * @throws NullPointerException if the task is null - */ - public Future submit(Runnable task) - { - return executorService_.submit(task); - } - - /** - * Executes the given tasks, returning a list of Futures holding - * their status and results when all complete. - * {@link Future#isDone} is true for each - * element of the returned list. - * Note that a completed task could have - * terminated either normally or by throwing an exception. - * The results of this method are undefined if the given - * collection is modified while this operation is in progress. - * - * @param tasks the collection of tasks - * @return A list of Futures representing the tasks, in the same - * sequential order as produced by the iterator for the - * given task list, each of which has completed. - * @throws InterruptedException if interrupted while waiting, in - * which case unfinished tasks are cancelled. - * @throws NullPointerException if tasks or any of its elements are null - * @throws RejectedExecutionException if any task cannot be - * scheduled for execution - */ - - public List> invokeAll(Collection> tasks) throws InterruptedException - { - return executorService_.invokeAll(tasks); - } - - /** - * Executes the given tasks, returning a list of Futures holding - * their status and results - * when all complete or the timeout expires, whichever happens first. - * {@link Future#isDone} is true for each - * element of the returned list. - * Upon return, tasks that have not completed are cancelled. - * Note that a completed task could have - * terminated either normally or by throwing an exception. - * The results of this method are undefined if the given - * collection is modified while this operation is in progress. - * - * @param tasks the collection of tasks - * @param timeout the maximum time to wait - * @param unit the time unit of the timeout argument - * @return a list of Futures representing the tasks, in the same - * sequential order as produced by the iterator for the - * given task list. If the operation did not time out, - * each task will have completed. If it did time out, some - * of these tasks will not have completed. - * @throws InterruptedException if interrupted while waiting, in - * which case unfinished tasks are cancelled - * @throws NullPointerException if tasks, any of its elements, or - * unit are null - * @throws RejectedExecutionException if any task cannot be scheduled - * for execution - */ - public List> invokeAll(Collection> tasks, long timeout, TimeUnit unit) throws InterruptedException - { - return executorService_.invokeAll(tasks, timeout, unit); - } - - /** - * Executes the given tasks, returning the result - * of one that has completed successfully (i.e., without throwing - * an exception), if any do. Upon normal or exceptional return, - * tasks that have not completed are cancelled. - * The results of this method are undefined if the given - * collection is modified while this operation is in progress. - * - * @param tasks the collection of tasks - * @return the result returned by one of the tasks - * @throws InterruptedException if interrupted while waiting - * @throws NullPointerException if tasks or any of its elements - * are null - * @throws IllegalArgumentException if tasks is empty - * @throws ExecutionException if no task successfully completes - * @throws RejectedExecutionException if tasks cannot be scheduled - * for execution - */ - public T invokeAny(Collection> tasks) throws InterruptedException, ExecutionException - { - return executorService_.invokeAny(tasks); - } - - /** - * Executes the given tasks, returning the result - * of one that has completed successfully (i.e., without throwing - * an exception), if any do before the given timeout elapses. - * Upon normal or exceptional return, tasks that have not - * completed are cancelled. - * The results of this method are undefined if the given - * collection is modified while this operation is in progress. - * - * @param tasks the collection of tasks - * @param timeout the maximum time to wait - * @param unit the time unit of the timeout argument - * @return the result returned by one of the tasks. - * @throws InterruptedException if interrupted while waiting - * @throws NullPointerException if tasks, any of its elements, or - * unit are null - * @throws TimeoutException if the given timeout elapses before - * any task successfully completes - * @throws ExecutionException if no task successfully completes - * @throws RejectedExecutionException if tasks cannot be scheduled - * for execution - */ - public T invokeAny(Collection> tasks, long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException - { - return executorService_.invokeAny(tasks, timeout, unit); - } -} From a4362ca906efa77799b8a4845810e9c5dfa8a682 Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Fri, 7 Oct 2011 23:21:26 +0000 Subject: [PATCH 2/5] Bump hsha threads to cores * 4 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1180277 13f79535-47bb-0310-9956-ffa450edef68 --- conf/cassandra.yaml | 3 ++- src/java/org/apache/cassandra/config/DatabaseDescriptor.java | 2 +- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/conf/cassandra.yaml b/conf/cassandra.yaml index bdd1ea88c1..9875148ef0 100644 --- a/conf/cassandra.yaml +++ b/conf/cassandra.yaml @@ -219,7 +219,8 @@ rpc_server_type: sync # disconnects before accepting more. The defaults for sync are min of 16 and max # unlimited. # -# For the Hsha server, the min and max both default to the number of CPU cores. +# For the Hsha server, the min and max both default to quadruple the number of +# CPU cores. # # This configuration is ignored by the async server. # diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index 5fd92913b4..8469f63f99 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -367,7 +367,7 @@ public class DatabaseDescriptor throw new ConfigurationException("Unknown rpc_server_type: " + conf.rpc_server_type); if (conf.rpc_min_threads == null) conf.rpc_min_threads = conf.rpc_server_type.toLowerCase().equals("hsha") - ? Runtime.getRuntime().availableProcessors() + ? Runtime.getRuntime().availableProcessors() * 4 : 16; if (conf.rpc_max_threads == null) conf.rpc_max_threads = conf.rpc_server_type.toLowerCase().equals("hsha") From 4f39d3e52f82d060bf96c2be0df6ff6782bc48e5 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Sat, 8 Oct 2011 10:37:19 +0000 Subject: [PATCH 3/5] Update version for 1.0.0 final release git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1180352 13f79535-47bb-0310-9956-ffa450edef68 --- build.xml | 2 +- debian/changelog | 6 ++++++ 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/build.xml b/build.xml index e6aff27412..8d7b43740c 100644 --- a/build.xml +++ b/build.xml @@ -25,7 +25,7 @@ - + diff --git a/debian/changelog b/debian/changelog index 50fd43b305..67eb6465ca 100644 --- a/debian/changelog +++ b/debian/changelog @@ -1,3 +1,9 @@ +cassandra (1.0.0) unstable; urgency=low + + * New release + + -- Sylvain Lebresne Sat, 08 Oct 2011 12:35:41 +0200 + cassandra (1.0.0~rc2) unstable; urgency=low * New release candidate From d790ed1bc871a9bd3a8c8640629532c0bed77656 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Mon, 10 Oct 2011 13:56:27 +0000 Subject: [PATCH 5/5] Fix places where uncompressed sstable size is used in place of the compressed one. patch by slebresne; reviewed by jbellis for CASSANDRA-3338 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0.0@1180970 13f79535-47bb-0310-9956-ffa450edef68 --- CHANGES.txt | 2 ++ build.xml | 2 +- .../apache/cassandra/db/ColumnFamilyStore.java | 6 +++--- .../db/compaction/CompactionManager.java | 6 +++--- .../SizeTieredCompactionStrategy.java | 2 +- .../apache/cassandra/io/sstable/SSTable.java | 2 +- .../cassandra/io/sstable/SSTableReader.java | 18 +++++++++++++++--- .../io/util/CompressedSegmentedFile.java | 2 +- .../cassandra/io/util/SegmentedFile.java | 10 ++++++++++ 9 files changed, 37 insertions(+), 13 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 7d10908b75..3d16db3fdb 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -12,6 +12,8 @@ * run compaction and hinted handoff threads at MIN_PRIORITY (CASSANDRA-3308) * default hsha thrift server to cpu core count in rpc pool (CASSANDRA-3329) * add bin\daemon to binary tarball for Windows service (CASSANDRA-3331) + * Fix places where uncompressed size of sstables was use in place of the + compressed one (CASSANDRA-3338) Fixes merged from 0.8 below: * Fix tool .bat files when CASSANDRA_HOME contains spaces (CASSANDRA-3258) * Force flush of status table when removing/updating token (CASSANDRA-3243) diff --git a/build.xml b/build.xml index 8d7b43740c..f687d124cf 100644 --- a/build.xml +++ b/build.xml @@ -350,7 +350,7 @@ url=${svn.entry.url}?pathrev=${svn.entry.commit.revision} - + diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index 3ef4ec6f19..3ff5451ea9 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -915,7 +915,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean long expectedFileSize = 0; for (SSTableReader sstable : sstables) { - long size = sstable.length(); + long size = sstable.onDiskLength(); expectedFileSize = expectedFileSize + size; } return expectedFileSize; @@ -930,9 +930,9 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean SSTableReader maxFile = null; for (SSTableReader sstable : sstables) { - if (sstable.length() > maxSize) + if (sstable.onDiskLength() > maxSize) { - maxSize = sstable.length(); + maxSize = sstable.onDiskLength(); maxFile = sstable; } } diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java index 796f79c4ea..020c6ff23b 100644 --- a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java @@ -472,7 +472,7 @@ public class CompactionManager implements CompactionManagerMBean boolean isCommutative = cfs.metadata.getDefaultValidator().isCommutative(); // Calculate the expected compacted filesize - String compactionFileLocation = cfs.table.getDataFileLocation(sstable.length()); + String compactionFileLocation = cfs.table.getDataFileLocation(sstable.onDiskLength()); if (compactionFileLocation == null) throw new IOException("disk full"); int expectedBloomFilterSize = Math.max(DatabaseDescriptor.getIndexInterval(), @@ -765,8 +765,8 @@ public class CompactionManager implements CompactionManagerMBean String format = "Cleaned up to %s. %,d to %,d (~%d%% of original) bytes for %,d keys. Time: %,dms."; long dTime = System.currentTimeMillis() - startTime; - long startsize = sstable.length(); - long endsize = newSstable.length(); + long startsize = sstable.onDiskLength(); + long endsize = newSstable.onDiskLength(); double ratio = (double)endsize / (double)startsize; logger.info(String.format(format, writer.getFilename(), startsize, endsize, (int)(ratio*100), totalkeysWritten, dTime)); } diff --git a/src/java/org/apache/cassandra/db/compaction/SizeTieredCompactionStrategy.java b/src/java/org/apache/cassandra/db/compaction/SizeTieredCompactionStrategy.java index b9909ec75b..e8c5c6fffc 100644 --- a/src/java/org/apache/cassandra/db/compaction/SizeTieredCompactionStrategy.java +++ b/src/java/org/apache/cassandra/db/compaction/SizeTieredCompactionStrategy.java @@ -100,7 +100,7 @@ public class SizeTieredCompactionStrategy extends AbstractCompactionStrategy { List> tableLengthPairs = new ArrayList>(); for(SSTableReader table: collection) - tableLengthPairs.add(new Pair(table, table.length())); + tableLengthPairs.add(new Pair(table, table.onDiskLength())); return tableLengthPairs; } diff --git a/src/java/org/apache/cassandra/io/sstable/SSTable.java b/src/java/org/apache/cassandra/io/sstable/SSTable.java index 5b7576b879..56b00a6542 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTable.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTable.java @@ -257,7 +257,7 @@ public abstract class SSTable long sum = 0; for (SSTableReader sstable : sstables) { - sum += sstable.length(); + sum += sstable.onDiskLength(); } return sum; } diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableReader.java b/src/java/org/apache/cassandra/io/sstable/SSTableReader.java index 5e7b5ad34c..934b2b8745 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableReader.java @@ -555,7 +555,7 @@ public class SSTableReader extends SSTable long right = getPosition(new DecoratedKey(range.right, null), Operator.GT); if (right == -1 || Range.isWrapAround(range.left, range.right)) // right is past the end of the file, or it wraps - right = length(); + right = uncompressedLength(); if (left == right) // empty range continue; @@ -669,13 +669,25 @@ public class SSTableReader extends SSTable } /** - * @return The length in bytes of the data file for this SSTable. + * @return The length in bytes of the data for this SSTable. For + * compressed files, this is not the same thing as the on disk size (see + * onDiskLength()) */ - public long length() + public long uncompressedLength() { return dfile.length; } + /** + * @return The length in bytes of the on disk size for this SSTable. For + * compressed files, this is not the same thing as the data length (see + * length()) + */ + public long onDiskLength() + { + return dfile.onDiskLength; + } + public boolean acquireReference() { while (true) diff --git a/src/java/org/apache/cassandra/io/util/CompressedSegmentedFile.java b/src/java/org/apache/cassandra/io/util/CompressedSegmentedFile.java index 76539d3e5f..4e0e7d8e3b 100644 --- a/src/java/org/apache/cassandra/io/util/CompressedSegmentedFile.java +++ b/src/java/org/apache/cassandra/io/util/CompressedSegmentedFile.java @@ -30,7 +30,7 @@ public class CompressedSegmentedFile extends SegmentedFile public CompressedSegmentedFile(String path, CompressionMetadata metadata) { - super(path, metadata.dataLength); + super(path, metadata.dataLength, metadata.compressedFileLength); this.metadata = metadata; } diff --git a/src/java/org/apache/cassandra/io/util/SegmentedFile.java b/src/java/org/apache/cassandra/io/util/SegmentedFile.java index 322e7426cb..602909fb89 100644 --- a/src/java/org/apache/cassandra/io/util/SegmentedFile.java +++ b/src/java/org/apache/cassandra/io/util/SegmentedFile.java @@ -42,13 +42,23 @@ public abstract class SegmentedFile public final String path; public final long length; + // This differs from length for compressed files (but we still need length for + // SegmentIterator because offsets in the file are relative to the uncompressed size) + public final long onDiskLength; + /** * Use getBuilder to get a Builder to construct a SegmentedFile. */ SegmentedFile(String path, long length) + { + this(path, length, length); + } + + protected SegmentedFile(String path, long length, long onDiskLength) { this.path = path; this.length = length; + this.onDiskLength = onDiskLength; } /**