diff --git a/CHANGES.txt b/CHANGES.txt
index b0f71aa2e7..40a35967c9 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -8,6 +8,14 @@
* Fixed the ability to set compaction strategy in cli using create column family command (CASSANDRA-2778)
* Add startup flag to renew counter node id (CASSANDRA-2788)
+
+0.8.2
+ * fix cache mbean getSize (CASSANDRA-2781)
+ * Add Date, Float, Double, and Boolean types (CASSANDRA-2530)
+ * fix repair hanging if a neighbor has nothing to send (CASSANDRA-2797)
+ * add jamm agent to cassandra.bat (CASSANDRA-2787)
+
+
0.8.1
* CQL:
- support for insert, delete in BATCH (CASSANDRA-2537)
@@ -15,7 +23,6 @@
- timestamp support for INSERT, UPDATE, and BATCH (CASSANDRA-2555)
- TTL support (CASSANDRA-2476)
- counter support (CASSANDRA-2473)
- - improve JDBC spec compliance (CASSANDRA-2720)
- ALTER COLUMNFAMILY (CASSANDRA-1709)
- DROP INDEX (CASSANDRA-2617)
- add SCHEMA/TABLE as aliases for KS/CF (CASSANDRA-2743)
@@ -58,6 +65,7 @@
by nio sockets (CASSANDRA-2654)
* restrict repair streaming to specific columnfamilies (CASSANDRA-2280)
* fix nodetool ring use with Ec2Snitch (CASSANDRA-2733)
+ * fix inconsistency window during bootstrap (CASSANDRA-833)
* fix removing columns and subcolumns that are supressed by a row or
supercolumn tombstone during replica resolution (CASSANDRA-2590)
* support sstable2json against snapshot sstables (CASSANDRA-2386)
@@ -68,7 +76,14 @@
* avoid skipping rows in scrub for counter column family (CASSANDRA-2759)
* fix ConcurrentModificationException in repair when dealing with 0.7 node
(CASSANDRA-2767)
+ * use threadsafe collections for StreamInSession (CASSANDRA-2766)
* avoid infinite loop when creating merkle tree (CASSANDRA-2758)
+ * avoids unmarking compacting sstable prematurely in cleanup (CASSANDRA-2769)
+ * fix NPE when the commit log is bypassed (CASSANDRA-2718)
+ * don't throw an exception in SS.isRPCServerRunning (CASSANDRA-2721)
+ * make stress.jar executable (CASSANDRA-2744)
+ * add daemon mode to java stress (CASSANDRA-2267)
+ * expose the DC and rack of a node through JMX and nodetool ring (CASSANDRA-2531)
0.8.0-final
@@ -175,7 +190,8 @@
* add key type information and alias (CASSANDRA-2311, 2396)
* cli no longer divides read_repair_chance by 100 (CASSANDRA-2458)
* made CompactionInfo.getTaskType return an enum (CASSANDRA-2482)
- * add a server-wide cap on measured memtable memory usage (CASSANDRA-2006)
+ * add a server-wide cap on measured memtable memory usage and aggressively
+ flush to keep under that threshold (CASSANDRA-2006)
* add unified UUIDType (CASSANDRA-2233)
diff --git a/NEWS.txt b/NEWS.txt
index c94da5d05d..85b7dec4d9 100644
--- a/NEWS.txt
+++ b/NEWS.txt
@@ -7,6 +7,31 @@ Upgrading
sstableloader tool instead.
+0.8.1
+=====
+
+Upgrading
+---------
+ - 0.8.1 is backwards compatible with 0.8, upgrade can be achieved by a
+ simple rolling restart.
+ - If upgrading for earlier version (0.7), please refer to the 0.8 section
+ for instructions.
+
+Features
+--------
+ - Numerous additions/improvements to CQL (support for counters, TTL, batch
+ inserts/deletes, index dropping, ...).
+ - Add two new AbstractTypes (comparator) to support compound keys
+ (CompositeType and DynamicCompositeType), as well as a ReverseType to
+ reverse the order of any existing comparator.
+ - New option to bypass the commit log on some keyspaces (for advanced
+ users).
+
+Tools
+-----
+ - Add new data bulk loading utility (sstableloader).
+
+
0.8
===
diff --git a/bin/cassandra.bat b/bin/cassandra.bat
index 0b927543b6..cd8bd68e0a 100644
--- a/bin/cassandra.bat
+++ b/bin/cassandra.bat
@@ -24,6 +24,7 @@ if NOT DEFINED JAVA_HOME goto err
REM ***** JAVA options *****
set JAVA_OPTS=^
-ea^
+ -javaagent:%CASSANDRA_HOME%\lib\jamm-0.2.2.jar^
-Xms1G^
-Xmx1G^
-XX:+HeapDumpOnOutOfMemoryError^
diff --git a/build.xml b/build.xml
index 7caa99fb19..4988be28b4 100644
--- a/build.xml
+++ b/build.xml
@@ -25,8 +25,8 @@
-
-
+
+
diff --git a/debian/changelog b/debian/changelog
index fd4892a03c..9bb7f182dc 100644
--- a/debian/changelog
+++ b/debian/changelog
@@ -1,3 +1,9 @@
+cassandra (0.8.1) unstable; urgency=low
+
+ * New release
+
+ -- Sylvain Lebresne Thu, 16 Jun 2011 09:37:27 +0200
+
cassandra (0.8.0) unstable; urgency=low
* New release
diff --git a/debian/rules b/debian/rules
index a7190ffd78..d76ff33a38 100755
--- a/debian/rules
+++ b/debian/rules
@@ -40,8 +40,6 @@ install: build
usr/share/cassandra
dh_install build/apache-cassandra-thrift-$(VERSION).jar \
usr/share/cassandra
- dh_install build/apache-cassandra-cql-*.jar \
- usr/share/cassandra
dh_link usr/share/cassandra/apache-cassandra-$(VERSION).jar \
usr/share/cassandra/apache-cassandra.jar
diff --git a/examples/hadoop_word_count/src/WordCount.java b/examples/hadoop_word_count/src/WordCount.java
index ddda2c3608..b21c05cb8b 100644
--- a/examples/hadoop_word_count/src/WordCount.java
+++ b/examples/hadoop_word_count/src/WordCount.java
@@ -136,13 +136,13 @@ public class WordCount extends Configured implements Tool
private static Mutation getMutation(Text word, int sum)
{
Column c = new Column();
- c.name = ByteBuffer.wrap(Arrays.copyOf(word.getBytes(), word.getLength()));
- c.value = ByteBufferUtil.bytes(String.valueOf(sum));
- c.timestamp = System.currentTimeMillis() * 1000;
+ c.setName(Arrays.copyOf(word.getBytes(), word.getLength()));
+ c.setValue(ByteBufferUtil.bytes(String.valueOf(sum)));
+ c.setTimestamp(System.currentTimeMillis());
Mutation m = new Mutation();
- m.column_or_supercolumn = new ColumnOrSuperColumn();
- m.column_or_supercolumn.column = c;
+ m.setColumn_or_supercolumn(new ColumnOrSuperColumn());
+ m.column_or_supercolumn.setColumn(c);
return m;
}
}
diff --git a/src/java/org/apache/cassandra/cache/InstrumentingCache.java b/src/java/org/apache/cassandra/cache/InstrumentingCache.java
index 00ba844d49..7c77f9f9a6 100644
--- a/src/java/org/apache/cassandra/cache/InstrumentingCache.java
+++ b/src/java/org/apache/cassandra/cache/InstrumentingCache.java
@@ -108,6 +108,11 @@ public class InstrumentingCache implements InstrumentingCacheMBean
return map.size();
}
+ public int getSize()
+ {
+ return size();
+ }
+
public long getHits()
{
return hits.get();
diff --git a/src/java/org/apache/cassandra/cache/InstrumentingCacheMBean.java b/src/java/org/apache/cassandra/cache/InstrumentingCacheMBean.java
index b7376960e7..3bcea58849 100644
--- a/src/java/org/apache/cassandra/cache/InstrumentingCacheMBean.java
+++ b/src/java/org/apache/cassandra/cache/InstrumentingCacheMBean.java
@@ -25,7 +25,7 @@ public interface InstrumentingCacheMBean
{
public int getCapacity();
public void setCapacity(int capacity);
- public int size();
+ public int getSize();
/** total request count since cache creation */
public long getRequests();
diff --git a/src/java/org/apache/cassandra/cli/CliClient.java b/src/java/org/apache/cassandra/cli/CliClient.java
index a2e96e7fb6..7513179fe3 100644
--- a/src/java/org/apache/cassandra/cli/CliClient.java
+++ b/src/java/org/apache/cassandra/cli/CliClient.java
@@ -265,6 +265,7 @@ public class CliClient
break;
case CliParser.NODE_CONSISTENCY_LEVEL:
executeConsistencyLevelStatement(tree);
+ break;
case CliParser.NODE_THRIFT_INCR:
executeIncr(tree, 1L);
break;
@@ -2087,17 +2088,6 @@ public class CliClient
return getBytesAccordingToType(superColumn, getFormatType(comparatorClass));
}
- /**
- * Converts column name into byte[] according to comparator type
- * @param superColumn - sub-column name from parser
- * @param columnFamily - column family name from parser
- * @return bytes[] - into which column name was converted according to comparator type
- */
- private byte[] subColumnNameAsByteArray(String superColumn, String columnFamily)
- {
- return TBaseHelper.byteBufferToByteArray(subColumnNameAsBytes(superColumn, columnFamily));
- }
-
/**
* Converts sub-column name into byte[] according to comparator type
* @param superColumn - sub-column name from parser
diff --git a/src/java/org/apache/cassandra/cql/Cql.g b/src/java/org/apache/cassandra/cql/Cql.g
index 02448cebd2..087485e9f4 100644
--- a/src/java/org/apache/cassandra/cql/Cql.g
+++ b/src/java/org/apache/cassandra/cql/Cql.g
@@ -427,11 +427,11 @@ dropColumnFamilyStatement returns [String cfam]
;
comparatorType
- : 'bytea' | 'ascii' | 'text' | 'varchar' | 'int' | 'varint' | 'bigint' | 'uuid' | 'counter'
+ : 'bytea' | 'ascii' | 'text' | 'varchar' | 'int' | 'varint' | 'bigint' | 'uuid' | 'counter' | 'boolean' | 'date' | 'float' | 'double'
;
term returns [Term item]
- : ( t=K_KEY | t=STRING_LITERAL | t=INTEGER | t=UUID | t=IDENT ) { $item = new Term($t.text, $t.type); }
+ : ( t=K_KEY | t=STRING_LITERAL | t=INTEGER | t=UUID | t=IDENT | t=FLOAT) { $item = new Term($t.text, $t.type); }
;
termList returns [List items]
diff --git a/src/java/org/apache/cassandra/cql/CreateColumnFamilyStatement.java b/src/java/org/apache/cassandra/cql/CreateColumnFamilyStatement.java
index 8a729c4a29..ac9e094f13 100644
--- a/src/java/org/apache/cassandra/cql/CreateColumnFamilyStatement.java
+++ b/src/java/org/apache/cassandra/cql/CreateColumnFamilyStatement.java
@@ -72,6 +72,10 @@ public class CreateColumnFamilyStatement
comparators.put("bigint", "LongType");
comparators.put("uuid", "UUIDType");
comparators.put("counter", "CounterColumnType");
+ comparators.put("boolean", "BooleanType");
+ comparators.put("date", "DateType");
+ comparators.put("float", "FloatType");
+ comparators.put("double", "DoubleType");
keywords.add(KW_COMPARATOR);
keywords.add(KW_COMMENT);
diff --git a/src/java/org/apache/cassandra/cql/Term.java b/src/java/org/apache/cassandra/cql/Term.java
index 5df293601c..b1519c8609 100644
--- a/src/java/org/apache/cassandra/cql/Term.java
+++ b/src/java/org/apache/cassandra/cql/Term.java
@@ -24,6 +24,7 @@ import java.nio.ByteBuffer;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.db.marshal.AsciiType;
+import org.apache.cassandra.db.marshal.FloatType;
import org.apache.cassandra.db.marshal.IntegerType;
import org.apache.cassandra.db.marshal.LexicalUUIDType;
import org.apache.cassandra.db.marshal.MarshalException;
@@ -107,6 +108,8 @@ public class Term
// we specifically want the Lexical class here, not "UUIDType," because we're supposed to have
// a uuid-shaped string here, and UUIDType also accepts integer or date strings (and turns them into version 1 uuids).
return LexicalUUIDType.instance.fromString(text);
+ case FLOAT:
+ return FloatType.instance.fromString(text);
}
// FIXME: handle scenario that should never happen
@@ -132,7 +135,7 @@ public class Term
enum TermType
{
- STRING, INTEGER, UUID;
+ STRING, INTEGER, UUID, FLOAT;
static TermType forInt(int type)
{
@@ -141,7 +144,9 @@ enum TermType
else if (type == CqlParser.INTEGER)
return INTEGER;
else if (type == CqlParser.UUID)
- return UUID;
+ return UUID;
+ else if (type == CqlParser.FLOAT)
+ return FLOAT;
// FIXME: handled scenario that should never occur.
return null;
diff --git a/src/java/org/apache/cassandra/db/EchoedRow.java b/src/java/org/apache/cassandra/db/EchoedRow.java
index 4e464f944d..c6d7231252 100644
--- a/src/java/org/apache/cassandra/db/EchoedRow.java
+++ b/src/java/org/apache/cassandra/db/EchoedRow.java
@@ -1,4 +1,25 @@
package org.apache.cassandra.db;
+/*
+ *
+ * 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.
+ *
+ */
+
import java.io.DataOutput;
import java.io.IOException;
diff --git a/src/java/org/apache/cassandra/db/commitlog/CommitLogHeader.java b/src/java/org/apache/cassandra/db/commitlog/CommitLogHeader.java
deleted file mode 100644
index e69de29bb2..0000000000
diff --git a/src/java/org/apache/cassandra/db/commitlog/ReplayPosition.java b/src/java/org/apache/cassandra/db/commitlog/ReplayPosition.java
index 6fea88d930..cd3db8a7e7 100644
--- a/src/java/org/apache/cassandra/db/commitlog/ReplayPosition.java
+++ b/src/java/org/apache/cassandra/db/commitlog/ReplayPosition.java
@@ -73,9 +73,9 @@ public class ReplayPosition implements Comparable
public int compare(ReplayPosition o1, ReplayPosition o2)
{
if (o1.segment != o2.segment)
- return new Long(o1.segment).compareTo(o2.segment);
+ return Long.valueOf(o1.segment).compareTo(o2.segment);
- return new Integer(o1.position).compareTo(o2.position);
+ return Integer.valueOf(o1.position).compareTo(o2.position);
}
};
diff --git a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java
index 134ac2eb65..0eb13f101c 100644
--- a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java
+++ b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java
@@ -652,56 +652,50 @@ public class CompactionManager implements CompactionManagerMBean
logger.debug("Expected bloom filter size : " + expectedBloomFilterSize);
SSTableWriter writer = null;
+
+ logger.info("Cleaning up " + sstable);
+ // Calculate the expected compacted filesize
+ long expectedRangeFileSize = cfs.getExpectedCompactedFileSize(Arrays.asList(sstable)) / 2;
+ String compactionFileLocation = table.getDataFileLocation(expectedRangeFileSize);
+ if (compactionFileLocation == null)
+ throw new IOException("disk full");
+
+ SSTableScanner scanner = sstable.getDirectScanner(CompactionIterator.FILE_BUFFER_SIZE);
+ SortedSet indexedColumns = cfs.getIndexedColumns();
+ CleanupInfo ci = new CleanupInfo(sstable, scanner);
+ executor.beginCompaction(ci);
try
{
- logger.info("Cleaning up " + sstable);
- // Calculate the expected compacted filesize
- long expectedRangeFileSize = cfs.getExpectedCompactedFileSize(Arrays.asList(sstable)) / 2;
- String compactionFileLocation = table.getDataFileLocation(expectedRangeFileSize);
- if (compactionFileLocation == null)
- throw new IOException("disk full");
-
- SSTableScanner scanner = sstable.getDirectScanner(CompactionIterator.FILE_BUFFER_SIZE);
- SortedSet indexedColumns = cfs.getIndexedColumns();
- CleanupInfo ci = new CleanupInfo(sstable, scanner);
- executor.beginCompaction(ci);
- try
+ while (scanner.hasNext())
{
- while (scanner.hasNext())
+ SSTableIdentityIterator row = (SSTableIdentityIterator) scanner.next();
+ if (Range.isTokenInRanges(row.getKey().token, ranges))
{
- SSTableIdentityIterator row = (SSTableIdentityIterator) scanner.next();
- if (Range.isTokenInRanges(row.getKey().token, ranges))
+ writer = maybeCreateWriter(cfs, compactionFileLocation, expectedBloomFilterSize, writer, Collections.singletonList(sstable));
+ writer.append(controller.getCompactedRow(row));
+ totalkeysWritten++;
+ }
+ else
+ {
+ cfs.invalidateCachedRow(row.getKey());
+ if (!indexedColumns.isEmpty() || isCommutative)
{
- writer = maybeCreateWriter(cfs, compactionFileLocation, expectedBloomFilterSize, writer, Collections.singletonList(sstable));
- writer.append(controller.getCompactedRow(row));
- totalkeysWritten++;
- }
- else
- {
- cfs.invalidateCachedRow(row.getKey());
- if (!indexedColumns.isEmpty() || isCommutative)
+ while (row.hasNext())
{
- while (row.hasNext())
- {
- IColumn column = row.next();
- if (column instanceof CounterColumn)
- renewer.maybeRenew((CounterColumn) column);
- if (indexedColumns.contains(column.name()))
- Table.cleanupIndexEntry(cfs, row.getKey().key, column);
- }
+ IColumn column = row.next();
+ if (column instanceof CounterColumn)
+ renewer.maybeRenew((CounterColumn) column);
+ if (indexedColumns.contains(column.name()))
+ Table.cleanupIndexEntry(cfs, row.getKey().key, column);
}
}
}
}
- finally
- {
- scanner.close();
- executor.finishCompaction(ci);
- }
}
finally
{
- cfs.getDataTracker().unmarkCompacting(Arrays.asList(sstable));
+ scanner.close();
+ executor.finishCompaction(ci);
}
List results = new ArrayList();
diff --git a/src/java/org/apache/cassandra/db/marshal/BooleanType.java b/src/java/org/apache/cassandra/db/marshal/BooleanType.java
new file mode 100644
index 0000000000..a77d2d2e32
--- /dev/null
+++ b/src/java/org/apache/cassandra/db/marshal/BooleanType.java
@@ -0,0 +1,136 @@
+package org.apache.cassandra.db.marshal;
+/*
+ *
+ * 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.
+ *
+ */
+
+import java.nio.ByteBuffer;
+import java.sql.Types;
+
+import org.apache.cassandra.utils.ByteBufferUtil;
+
+public class BooleanType extends AbstractType
+{
+ public static final BooleanType instance = new BooleanType();
+
+ BooleanType() {} // singleton
+
+ public Boolean compose(ByteBuffer bytes)
+ {
+ byte value = bytes.get(bytes.position());
+ return Boolean.valueOf(value ==0 ? false:true);
+ }
+
+ public ByteBuffer decompose(Boolean value)
+ {
+ return (value==null) ? ByteBufferUtil.EMPTY_BYTE_BUFFER
+ : value ? ByteBuffer.wrap(new byte[]{1}) // true
+ : ByteBuffer.wrap(new byte[]{0}); // false
+ }
+
+ public int compare(ByteBuffer o1, ByteBuffer o2)
+ {
+ if ((o1 == null) || (o1.remaining() != 1))
+ return ((o2 == null) || (o2.remaining() != 1)) ? 0 : -1;
+ if ((o2 == null) || (o2.remaining() != 1))
+ return 1;
+
+ return o1.compareTo(o2);
+ }
+
+ public String getString(ByteBuffer bytes)
+ {
+ if (bytes.remaining() == 0)
+ {
+ return Boolean.FALSE.toString();
+ }
+ if (bytes.remaining() != 1)
+ {
+ throw new MarshalException("A boolean is stored in exactly 1 byte: "+bytes.remaining());
+ }
+ byte value = bytes.get(bytes.position());
+
+ return value ==0 ? Boolean.FALSE.toString(): Boolean.TRUE.toString();
+ }
+
+ public String toString(Boolean b)
+ {
+ return b.toString();
+ }
+
+ public ByteBuffer fromString(String source) throws MarshalException
+ {
+
+ if (source.isEmpty()|| source.equalsIgnoreCase(Boolean.FALSE.toString()))
+ return decompose(false);
+
+ if (source.equalsIgnoreCase(Boolean.TRUE.toString()))
+ return decompose(true);
+
+ throw new MarshalException(String.format("unable to make boolean from '%s'", source));
+
+ }
+
+ public void validate(ByteBuffer bytes) throws MarshalException
+ {
+ if (bytes.remaining() != 1 && bytes.remaining() != 0)
+ throw new MarshalException(String.format("Expected 1 or 0 byte value (%d)", bytes.remaining()));
+ }
+
+ public Class getType()
+ {
+ return Boolean.class;
+ }
+
+ public boolean isSigned()
+ {
+ return false;
+ }
+
+ public boolean isCaseSensitive()
+ {
+ return false;
+ }
+
+ public boolean isCurrency()
+ {
+ return false;
+ }
+
+ public int getPrecision(Boolean obj)
+ {
+ return -1;
+ }
+
+ public int getScale(Boolean obj)
+ {
+ return -1;
+ }
+
+ public int getJdbcType()
+ {
+ return Types.BOOLEAN;
+ }
+
+ public boolean needsQuotes()
+ {
+ return false;
+ }
+
+}
diff --git a/src/java/org/apache/cassandra/db/marshal/DateType.java b/src/java/org/apache/cassandra/db/marshal/DateType.java
new file mode 100644
index 0000000000..dc1c87281e
--- /dev/null
+++ b/src/java/org/apache/cassandra/db/marshal/DateType.java
@@ -0,0 +1,178 @@
+package org.apache.cassandra.db.marshal;
+/*
+ *
+ * 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.
+ *
+ */
+import static org.apache.cassandra.db.marshal.TimeUUIDType.iso8601Patterns;
+
+import java.nio.ByteBuffer;
+import java.sql.Types;
+import java.text.ParseException;
+import java.text.SimpleDateFormat;
+import java.util.Date;
+
+import org.apache.cassandra.utils.ByteBufferUtil;
+import org.apache.commons.lang.time.DateUtils;
+
+public class DateType extends AbstractType
+{
+ public static final DateType instance = new DateType();
+
+ static final String DEFAULT_FORMAT = iso8601Patterns[3];
+
+ static final SimpleDateFormat FORMATTER = new SimpleDateFormat(DEFAULT_FORMAT);
+
+ DateType() {} // singleton
+
+ public Date compose(ByteBuffer bytes)
+ {
+ return new Date(ByteBufferUtil.toLong(bytes));
+ }
+
+ public ByteBuffer decompose(Date value)
+ {
+ return (value==null) ? ByteBufferUtil.EMPTY_BYTE_BUFFER
+ : ByteBufferUtil.bytes(value.getTime());
+ }
+
+
+ public int compare(ByteBuffer o1, ByteBuffer o2)
+ {
+ if (o1.remaining() == 0)
+ {
+ return o2.remaining() == 0 ? 0 : -1;
+ }
+ if (o2.remaining() == 0)
+ {
+ return 1;
+ }
+
+ return ByteBufferUtil.compareUnsigned(o1, o2);
+ }
+
+ public String getString(ByteBuffer bytes)
+ {
+ if (bytes.remaining() == 0)
+ {
+ return "";
+ }
+ if (bytes.remaining() != 8)
+ {
+ throw new MarshalException("A date is exactly 8 bytes (stored as a long): "+bytes.remaining());
+ }
+
+ // uses ISO-8601 formatted string
+ return FORMATTER.format(new Date(bytes.getLong(bytes.position())));
+ }
+
+ public String toString(Date d)
+ {
+ // uses ISO-8601 formatted string
+ return FORMATTER.format(d);
+ }
+
+ public ByteBuffer fromString(String source) throws MarshalException
+ {
+ // Return an empty ByteBuffer for an empty string.
+ if (source.isEmpty())
+ return ByteBufferUtil.EMPTY_BYTE_BUFFER;
+
+ long millis;
+ ByteBuffer idBytes = null;
+
+ if (source.toLowerCase().equals("now"))
+ {
+ millis = System.currentTimeMillis();
+ idBytes = ByteBufferUtil.bytes(millis);
+ }
+ // Milliseconds since epoch?
+ else if (source.matches("^\\d+$"))
+ {
+ try
+ {
+ idBytes = ByteBufferUtil.bytes(Long.parseLong(source));
+ }
+ catch (NumberFormatException e)
+ {
+ throw new MarshalException(String.format("unable to make long (for date) from: '%s'", source), e);
+ }
+ }
+ // Last chance, attempt to parse as date-time string
+ else
+ {
+ try
+ {
+ millis = DateUtils.parseDate(source, iso8601Patterns).getTime();
+ idBytes = ByteBufferUtil.bytes(millis);
+ }
+ catch (ParseException e1)
+ {
+ throw new MarshalException(String.format("unable to coerce '%s' to a formatted date (long)", source), e1);
+ }
+ }
+
+ return idBytes;
+ }
+
+ public void validate(ByteBuffer bytes) throws MarshalException
+ {
+ if (bytes.remaining() != 8 && bytes.remaining() != 0)
+ throw new MarshalException(String.format("Expected 8 or 0 byte long for date (%d)", bytes.remaining()));
+ }
+
+ public Class getType()
+ {
+ return Date.class;
+ }
+
+ public boolean isSigned()
+ {
+ return false;
+ }
+
+ public boolean isCaseSensitive()
+ {
+ return false;
+ }
+
+ public boolean isCurrency()
+ {
+ return false;
+ }
+
+ public int getPrecision(Date obj)
+ {
+ return -1;
+ }
+
+ public int getScale(Date obj)
+ {
+ return -1;
+ }
+
+ public int getJdbcType()
+ {
+ return Types.DATE;
+ }
+
+ public boolean needsQuotes()
+ {
+ return false;
+ }
+}
diff --git a/src/java/org/apache/cassandra/db/marshal/DoubleType.java b/src/java/org/apache/cassandra/db/marshal/DoubleType.java
new file mode 100644
index 0000000000..758f1b9aeb
--- /dev/null
+++ b/src/java/org/apache/cassandra/db/marshal/DoubleType.java
@@ -0,0 +1,143 @@
+package org.apache.cassandra.db.marshal;
+/*
+ *
+ * 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.
+ *
+ */
+
+import java.nio.ByteBuffer;
+import java.sql.Types;
+
+import org.apache.cassandra.utils.ByteBufferUtil;
+
+public class DoubleType extends AbstractType
+{
+ public static final DoubleType instance = new DoubleType();
+
+ DoubleType() {} // singleton
+
+ public Double compose(ByteBuffer bytes)
+ {
+ return ByteBufferUtil.toDouble(bytes);
+ }
+
+ public ByteBuffer decompose(Double value)
+ {
+ return (value==null) ? ByteBufferUtil.EMPTY_BYTE_BUFFER : ByteBufferUtil.bytes(value);
+ }
+
+
+ public int compare(ByteBuffer o1, ByteBuffer o2)
+ {
+ if (o1.remaining() == 0)
+ {
+ return o2.remaining() == 0 ? 0 : -1;
+ }
+ if (o2.remaining() == 0)
+ {
+ return 1;
+ }
+
+ return compose(o1).compareTo(compose(o2));
+ }
+
+ public String getString(ByteBuffer bytes)
+ {
+ if (bytes.remaining() == 0)
+ {
+ return "";
+ }
+ if (bytes.remaining() != 8)
+ {
+ throw new MarshalException("A double is exactly 8 bytes : "+bytes.remaining());
+ }
+
+ return compose(bytes).toString();
+ }
+
+ public String toString(Double d)
+ {
+ return d.toString();
+ }
+
+ public ByteBuffer fromString(String source) throws MarshalException
+ {
+ // Return an empty ByteBuffer for an empty string.
+ if (source.isEmpty())
+ return ByteBufferUtil.EMPTY_BYTE_BUFFER;
+
+ Double d;
+ try
+ {
+ d = Double.parseDouble(source);
+ }
+ catch (NumberFormatException e1)
+ {
+ throw new MarshalException(String.format("unable to coerce '%s' to a double", source), e1);
+ }
+
+ return decompose(d);
+ }
+
+ public void validate(ByteBuffer bytes) throws MarshalException
+ {
+ if (bytes.remaining() != 8 && bytes.remaining() != 0)
+ throw new MarshalException(String.format("Expected 8 or 0 byte value for a double (%d)", bytes.remaining()));
+ }
+
+ public Class getType()
+ {
+ return Double.class;
+ }
+
+ public boolean isSigned()
+ {
+ return true;
+ }
+
+ public boolean isCaseSensitive()
+ {
+ return false;
+ }
+
+ public boolean isCurrency()
+ {
+ return false;
+ }
+
+ public int getPrecision(Double obj) // see: http://teaching.idallen.org/dat2343/09f/notes/10FloatingPoint.htm
+ {
+ return 15;
+ }
+
+ public int getScale(Double obj) // see: http://teaching.idallen.org/dat2343/09f/notes/10FloatingPoint.htm
+ {
+ return 300;
+ }
+
+ public int getJdbcType()
+ {
+ return Types.DOUBLE;
+ }
+
+ public boolean needsQuotes()
+ {
+ return false;
+ }
+
+}
diff --git a/src/java/org/apache/cassandra/db/marshal/FloatType.java b/src/java/org/apache/cassandra/db/marshal/FloatType.java
new file mode 100644
index 0000000000..f6b4e34fd5
--- /dev/null
+++ b/src/java/org/apache/cassandra/db/marshal/FloatType.java
@@ -0,0 +1,144 @@
+package org.apache.cassandra.db.marshal;
+/*
+ *
+ * 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.
+ *
+ */
+
+import java.nio.ByteBuffer;
+import java.sql.Types;
+
+import org.apache.cassandra.utils.ByteBufferUtil;
+
+
+public class FloatType extends AbstractType
+{
+ public static final FloatType instance = new FloatType();
+
+ FloatType() {} // singleton
+
+ public Float compose(ByteBuffer bytes)
+ {
+ return ByteBufferUtil.toFloat(bytes);
+ }
+
+ public ByteBuffer decompose(Float value)
+ {
+ return (value==null) ? ByteBufferUtil.EMPTY_BYTE_BUFFER : ByteBufferUtil.bytes(value);
+ }
+
+
+ public int compare(ByteBuffer o1, ByteBuffer o2)
+ {
+ if (o1.remaining() == 0)
+ {
+ return o2.remaining() == 0 ? 0 : -1;
+ }
+ if (o2.remaining() == 0)
+ {
+ return 1;
+ }
+
+ return compose(o1).compareTo(compose(o2));
+ }
+
+ public String getString(ByteBuffer bytes)
+ {
+ if (bytes.remaining() == 0)
+ {
+ return "";
+ }
+ if (bytes.remaining() != 4)
+ {
+ throw new MarshalException("A float is exactly 4 bytes : "+bytes.remaining());
+ }
+
+ return compose(bytes).toString();
+ }
+
+ public String toString(Float d)
+ {
+ return d.toString();
+ }
+
+ public ByteBuffer fromString(String source) throws MarshalException
+ {
+ // Return an empty ByteBuffer for an empty string.
+ if (source.isEmpty())
+ return ByteBufferUtil.EMPTY_BYTE_BUFFER;
+
+ Float f;
+ try
+ {
+ f = Float.parseFloat(source);
+ }
+ catch (NumberFormatException e1)
+ {
+ throw new MarshalException(String.format("unable to coerce '%s' to a float", source), e1);
+ }
+
+ return ByteBufferUtil.bytes(f);
+ }
+
+ public void validate(ByteBuffer bytes) throws MarshalException
+ {
+ if (bytes.remaining() != 4 && bytes.remaining() != 0)
+ throw new MarshalException(String.format("Expected 4 or 0 byte value for a float (%d)", bytes.remaining()));
+ }
+
+ public Class getType()
+ {
+ return Float.class;
+ }
+
+ public boolean isSigned()
+ {
+ return true;
+ }
+
+ public boolean isCaseSensitive()
+ {
+ return false;
+ }
+
+ public boolean isCurrency()
+ {
+ return false;
+ }
+
+ public int getPrecision(Float obj) // see: http://teaching.idallen.org/dat2343/09f/notes/10FloatingPoint.htm
+ {
+ return 7;
+ }
+
+ public int getScale(Float obj) // see: http://teaching.idallen.org/dat2343/09f/notes/10FloatingPoint.htm
+ {
+ return 40;
+ }
+
+ public int getJdbcType()
+ {
+ return Types.FLOAT;
+ }
+
+ public boolean needsQuotes()
+ {
+ return false;
+ }
+
+}
diff --git a/src/java/org/apache/cassandra/gms/FailureDetector.java b/src/java/org/apache/cassandra/gms/FailureDetector.java
index c7b09b067a..ab1cc07772 100644
--- a/src/java/org/apache/cassandra/gms/FailureDetector.java
+++ b/src/java/org/apache/cassandra/gms/FailureDetector.java
@@ -257,7 +257,7 @@ class ArrivalWindow
}
// see CASSANDRA-2597 for an explanation of the math at work here.
- double phi(long tnow)
+ synchronized double phi(long tnow)
{
int size = arrivalIntervals_.size();
double t = tnow - tLast_;
diff --git a/src/java/org/apache/cassandra/gms/Gossiper.java b/src/java/org/apache/cassandra/gms/Gossiper.java
index a3333fac0c..7f71a64138 100644
--- a/src/java/org/apache/cassandra/gms/Gossiper.java
+++ b/src/java/org/apache/cassandra/gms/Gossiper.java
@@ -190,7 +190,7 @@ public class Gossiper implements IFailureDetectionEventListener
public void setVersion(InetAddress address, int version)
{
- Integer old = versions.put(address, version);
+ versions.put(address, version);
}
public Integer getVersion(InetAddress address)
@@ -717,7 +717,6 @@ public class Gossiper implements IFailureDetectionEventListener
{
// don't assert here, since if the node restarts the version will go back to zero
int oldVersion = localState.getHeartBeatState().getHeartBeatVersion();
- Map localAppStateMap = localState.getApplicationStateMap();
localState.setHeartBeatState(remoteState.getHeartBeatState());
if (logger.isTraceEnabled())
diff --git a/src/java/org/apache/cassandra/io/util/FileUtils.java b/src/java/org/apache/cassandra/io/util/FileUtils.java
index 4ffd3c799a..e19084aafe 100644
--- a/src/java/org/apache/cassandra/io/util/FileUtils.java
+++ b/src/java/org/apache/cassandra/io/util/FileUtils.java
@@ -136,26 +136,13 @@ public class FileUtils
}
}
- public static void createFile(String directory) throws IOException
- {
- File file = new File(directory);
- if ( !file.exists() )
- file.createNewFile();
- }
-
- public static boolean isExists(String filename) throws IOException
- {
- File file = new File(filename);
- return file.exists();
- }
-
public static boolean delete(String file)
{
File f = new File(file);
return f.delete();
}
- public static boolean delete(List files) throws IOException
+ public static boolean delete(List files)
{
boolean bVal = true;
for ( int i = 0; i < files.size(); ++i )
@@ -172,7 +159,7 @@ public class FileUtils
return bVal;
}
- public static void delete(File[] files) throws IOException
+ public static void delete(File[] files)
{
for ( File file : files )
{
diff --git a/src/java/org/apache/cassandra/streaming/StreamInSession.java b/src/java/org/apache/cassandra/streaming/StreamInSession.java
index c3aa0a02f8..4c32a40bf1 100644
--- a/src/java/org/apache/cassandra/streaming/StreamInSession.java
+++ b/src/java/org/apache/cassandra/streaming/StreamInSession.java
@@ -24,18 +24,20 @@ import java.util.*;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.Future;
+import java.util.concurrent.LinkedBlockingQueue;
-import org.apache.cassandra.gms.Gossiper;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.apache.cassandra.db.ColumnFamilyStore;
-import org.apache.cassandra.db.compaction.CompactionManager;
import org.apache.cassandra.db.Table;
+import org.apache.cassandra.db.compaction.CompactionManager;
+import org.apache.cassandra.gms.Gossiper;
import org.apache.cassandra.io.sstable.SSTableReader;
import org.apache.cassandra.net.MessagingService;
import org.apache.cassandra.utils.Pair;
import org.cliffc.high_scale_lib.NonBlockingHashMap;
+import org.cliffc.high_scale_lib.NonBlockingHashSet;
/** each context gets its own StreamInSession. So there may be >1 Session per host */
public class StreamInSession
@@ -44,11 +46,11 @@ public class StreamInSession
private static ConcurrentMap, StreamInSession> sessions = new NonBlockingHashMap, StreamInSession>();
- private final List files = new ArrayList();
+ private final Set files = new NonBlockingHashSet();
private final Pair context;
private final Runnable callback;
private String table;
- private final List> buildFutures = new ArrayList>();
+ private final Collection> buildFutures = new LinkedBlockingQueue>();
private PendingFile current;
private StreamInSession(Pair context, Runnable callback)
diff --git a/src/java/org/apache/cassandra/streaming/StreamOut.java b/src/java/org/apache/cassandra/streaming/StreamOut.java
index 6783ac7098..c9ac671811 100644
--- a/src/java/org/apache/cassandra/streaming/StreamOut.java
+++ b/src/java/org/apache/cassandra/streaming/StreamOut.java
@@ -122,15 +122,10 @@ public class StreamOut
{
List pending = createPendingFiles(sstables, ranges, type);
- if (pending.size() > 0)
- {
- session.addFilesToStream(pending);
- session.begin();
- }
- else
- {
- session.close();
- }
+ // Even if the list of pending files is empty, we need to initiate the transfer otherwise
+ // the remote end will hang in cases where this was a requested transfer.
+ session.addFilesToStream(pending);
+ session.begin();
}
// called prior to sending anything.
diff --git a/src/java/org/apache/cassandra/tools/NodeCmd.java b/src/java/org/apache/cassandra/tools/NodeCmd.java
index d4031f7e50..ec7bb24429 100644
--- a/src/java/org/apache/cassandra/tools/NodeCmd.java
+++ b/src/java/org/apache/cassandra/tools/NodeCmd.java
@@ -450,7 +450,7 @@ public class NodeCmd
if (keyCacheMBean.getCapacity() > 0)
{
outs.println("\t\tKey cache capacity: " + keyCacheMBean.getCapacity());
- outs.println("\t\tKey cache size: " + keyCacheMBean.size());
+ outs.println("\t\tKey cache size: " + keyCacheMBean.getSize());
outs.println("\t\tKey cache hit rate: " + keyCacheMBean.getRecentHitRate());
}
else
@@ -462,7 +462,7 @@ public class NodeCmd
if (rowCacheMBean.getCapacity() > 0)
{
outs.println("\t\tRow cache capacity: " + rowCacheMBean.getCapacity());
- outs.println("\t\tRow cache size: " + rowCacheMBean.size());
+ outs.println("\t\tRow cache size: " + rowCacheMBean.getSize());
outs.println("\t\tRow cache hit rate: " + rowCacheMBean.getRecentHitRate());
}
else
diff --git a/src/java/org/apache/cassandra/tools/SSTableExport.java b/src/java/org/apache/cassandra/tools/SSTableExport.java
index e5aca7de9e..00af39104f 100644
--- a/src/java/org/apache/cassandra/tools/SSTableExport.java
+++ b/src/java/org/apache/cassandra/tools/SSTableExport.java
@@ -25,6 +25,7 @@ import java.nio.ByteBuffer;
import java.util.*;
import org.apache.cassandra.config.CFMetaData;
+import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.db.*;
import org.apache.cassandra.db.marshal.AbstractType;
import org.apache.cassandra.io.util.BufferedRandomAccessFile;
@@ -33,11 +34,11 @@ import org.apache.cassandra.service.StorageService;
import org.apache.commons.cli.*;
import org.apache.cassandra.config.ConfigurationException;
-import org.apache.cassandra.config.DatabaseDescriptor;
import org.apache.cassandra.dht.IPartitioner;
import org.apache.cassandra.io.sstable.*;
import org.apache.cassandra.utils.ByteBufferUtil;
-import org.apache.cassandra.utils.Pair;
+import org.codehaus.jackson.JsonGenerator;
+import org.codehaus.jackson.map.ObjectMapper;
import static org.apache.cassandra.utils.ByteBufferUtil.bytesToHex;
import static org.apache.cassandra.utils.ByteBufferUtil.hexToBytes;
@@ -47,12 +48,12 @@ import static org.apache.cassandra.utils.ByteBufferUtil.hexToBytes;
*/
public class SSTableExport
{
- // size of the columns page
- private static final int PAGE_SIZE = 1000;
+ private static ObjectMapper jsonMapper = new ObjectMapper();
private static final String KEY_OPTION = "k";
private static final String EXCLUDEKEY_OPTION = "x";
private static final String ENUMERATEKEYS_OPTION = "e";
+
private static Options options;
private static CommandLine cmd;
@@ -72,42 +73,36 @@ public class SSTableExport
Option optEnumerate = new Option(ENUMERATEKEYS_OPTION, false, "enumerate keys only");
options.addOption(optEnumerate);
- }
- /**
- * Wraps given string into quotes
- * @param val string to quote
- * @return quoted string
- */
- private static String quote(String val)
- {
- return String.format("\"%s\"", val);
+ // disabling auto close of the stream
+ jsonMapper.configure(JsonGenerator.Feature.AUTO_CLOSE_TARGET, false);
}
/**
* JSON Hash Key serializer
- * @param val value to set as a key
- * @return JSON Hash key
+ *
+ * @param out The output steam to write data
+ * @param value value to set as a key
*/
- private static String asKey(String val)
+ private static void writeKey(PrintStream out, String value)
{
- return String.format("%s: ", quote(val));
+ writeJSON(out, value);
+ out.print(": ");
}
/**
* Serialize columns using given column iterator
+ *
* @param columns column iterator
* @param out output stream
* @param comparator columns comparator
* @param cfMetaData Column Family metadata (to get validator)
- * @return pair of (number of columns serialized, last column serialized)
*/
private static void serializeColumns(Iterator columns, PrintStream out, AbstractType comparator, CFMetaData cfMetaData)
{
while (columns.hasNext())
{
- IColumn column = columns.next();
- serializeColumn(column, out, comparator, cfMetaData);
+ writeJSON(out, serializeColumn(columns.next(), comparator, cfMetaData));
if (columns.hasNext())
out.print(", ");
@@ -116,47 +111,42 @@ public class SSTableExport
/**
* Serialize a given column to the JSON format
+ *
* @param column column presentation
- * @param out output stream
* @param comparator columns comparator
* @param cfMetaData Column Family metadata (to get validator)
+ *
+ * @return column as serialized list
*/
- private static void serializeColumn(IColumn column, PrintStream out, AbstractType comparator, CFMetaData cfMetaData)
+ private static List