mirror of https://github.com/apache/cassandra
merge from 1.0.0
git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-1.0@1181094 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
commit
917af84403
|
|
@ -19,6 +19,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)
|
||||
|
|
|
|||
|
|
@ -25,7 +25,7 @@
|
|||
<property name="debuglevel" value="source,lines,vars"/>
|
||||
|
||||
<!-- default version and SCM information (we need the default SCM info as people may checkout with git-svn) -->
|
||||
<property name="base.version" value="1.0.0-rc2"/>
|
||||
<property name="base.version" value="1.0.0"/>
|
||||
<property name="scm.default.path" value="cassandra/branches/cassandra-1.0.0"/>
|
||||
<property name="scm.default.connection" value="scm:svn:http://svn.apache.org/repos/asf/${scm.default.path}"/>
|
||||
<property name="scm.default.developerConnection" value="scm:svn:https://svn.apache.org/repos/asf/${scm.default.path}"/>
|
||||
|
|
@ -350,7 +350,7 @@ url=${svn.entry.url}?pathrev=${svn.entry.commit.revision}
|
|||
<license name="The Apache Software License, Version 2.0" url="http://www.apache.org/licenses/LICENSE-2.0.txt"/>
|
||||
<scm connection="${scm.connection}" developerConnection="${scm.developerConnection}" url="${scm.url}"/>
|
||||
<dependencyManagement>
|
||||
<dependency groupId="org.xerial.snappy" artifactId="snappy-java" version="1.0.3.3"/>
|
||||
<dependency groupId="org.xerial.snappy" artifactId="snappy-java" version="1.0.3"/>
|
||||
<dependency groupId="com.ning" artifactId="compress-lzf" version="0.8.4"/>
|
||||
<dependency groupId="com.google.guava" artifactId="guava" version="r08"/>
|
||||
<dependency groupId="commons-cli" artifactId="commons-cli" version="1.1"/>
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
#
|
||||
|
|
|
|||
|
|
@ -1,3 +1,9 @@
|
|||
cassandra (1.0.0) unstable; urgency=low
|
||||
|
||||
* New release
|
||||
|
||||
-- Sylvain Lebresne <slebresne@apache.org> Sat, 08 Oct 2011 12:35:41 +0200
|
||||
|
||||
cassandra (1.0.0~rc2) unstable; urgency=low
|
||||
|
||||
* New release candidate
|
||||
|
|
|
|||
|
|
@ -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<Runnable> 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 <tt>Executor</tt> 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.
|
||||
*
|
||||
* <p>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}<tt>("modifyThread")</tt>,
|
||||
* or the security manager's <tt>checkAccess</tt> 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.
|
||||
*
|
||||
* <p>This method does not wait for actively executing tasks to
|
||||
* terminate. Use {@link #awaitTermination awaitTermination} to
|
||||
* do that.
|
||||
*
|
||||
* <p>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}<tt>("modifyThread")</tt>,
|
||||
* or the security manager's <tt>checkAccess</tt> method
|
||||
* denies access.
|
||||
*/
|
||||
public List<Runnable> shutdownNow()
|
||||
{
|
||||
return executorService_.shutdownNow();
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns <tt>true</tt> if this executor has been shut down.
|
||||
*
|
||||
* @return <tt>true</tt> if this executor has been shut down
|
||||
*/
|
||||
public boolean isShutdown()
|
||||
{
|
||||
return executorService_.isShutdown();
|
||||
}
|
||||
|
||||
/**
|
||||
* Returns <tt>true</tt> if all tasks have completed following shut down.
|
||||
* Note that <tt>isTerminated</tt> is never <tt>true</tt> unless
|
||||
* either <tt>shutdown</tt> or <tt>shutdownNow</tt> was called first.
|
||||
*
|
||||
* @return <tt>true</tt> 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 <tt>true</tt> if this executor terminated and
|
||||
* <tt>false</tt> 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 <tt>get</tt> method will return the task's result upon
|
||||
* successful completion.
|
||||
*
|
||||
* <p>
|
||||
* If you would like to immediately block waiting
|
||||
* for a task, you can use constructions of the form
|
||||
* <tt>result = exec.submit(aCallable).get();</tt>
|
||||
*
|
||||
* <p> 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 <T> Future<T> submit(Callable<T> task)
|
||||
{
|
||||
return executorService_.submit(task);
|
||||
}
|
||||
|
||||
/**
|
||||
* Submits a Runnable task for execution and returns a Future
|
||||
* representing that task. The Future's <tt>get</tt> 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 <T> Future<T> 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 <tt>get</tt> method will
|
||||
* return <tt>null</tt> upon <em>successful</em> 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 <tt>true</tt> for each
|
||||
* element of the returned list.
|
||||
* Note that a <em>completed</em> 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 <tt>null</tt>
|
||||
* @throws RejectedExecutionException if any task cannot be
|
||||
* scheduled for execution
|
||||
*/
|
||||
|
||||
public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> 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 <tt>true</tt> for each
|
||||
* element of the returned list.
|
||||
* Upon return, tasks that have not completed are cancelled.
|
||||
* Note that a <em>completed</em> 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 <tt>null</tt>
|
||||
* @throws RejectedExecutionException if any task cannot be scheduled
|
||||
* for execution
|
||||
*/
|
||||
public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> 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 <tt>null</tt>
|
||||
* @throws IllegalArgumentException if tasks is empty
|
||||
* @throws ExecutionException if no task successfully completes
|
||||
* @throws RejectedExecutionException if tasks cannot be scheduled
|
||||
* for execution
|
||||
*/
|
||||
public <T> T invokeAny(Collection<? extends Callable<T>> 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 <tt>null</tt>
|
||||
* @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> T invokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException
|
||||
{
|
||||
return executorService_.invokeAny(tasks, timeout, unit);
|
||||
}
|
||||
}
|
||||
|
|
@ -368,7 +368,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")
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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));
|
||||
}
|
||||
|
|
|
|||
|
|
@ -100,7 +100,7 @@ public class SizeTieredCompactionStrategy extends AbstractCompactionStrategy
|
|||
{
|
||||
List<Pair<SSTableReader, Long>> tableLengthPairs = new ArrayList<Pair<SSTableReader, Long>>();
|
||||
for(SSTableReader table: collection)
|
||||
tableLengthPairs.add(new Pair<SSTableReader, Long>(table, table.length()));
|
||||
tableLengthPairs.add(new Pair<SSTableReader, Long>(table, table.onDiskLength()));
|
||||
return tableLengthPairs;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -257,7 +257,7 @@ public abstract class SSTable
|
|||
long sum = 0;
|
||||
for (SSTableReader sstable : sstables)
|
||||
{
|
||||
sum += sstable.length();
|
||||
sum += sstable.onDiskLength();
|
||||
}
|
||||
return sum;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
|
|||
Loading…
Reference in New Issue