mirror of https://github.com/apache/cassandra
remove avro rpc source
Patch by eevans; reviewed by Jeremy Hanna for CASSANDRA-926 git-svn-id: https://svn.apache.org/repos/asf/cassandra/trunk@1059458 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
parent
9f71fb6e5e
commit
e4d9524efc
|
|
@ -1,376 +0,0 @@
|
|||
#!/usr/bin/python
|
||||
# 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.
|
||||
|
||||
# expects a Cassandra server to be running and listening on port 9160.
|
||||
# (read tests expect insert tests to have run first too.)
|
||||
|
||||
have_multiproc = False
|
||||
try:
|
||||
from multiprocessing import Array as array, Process as Thread
|
||||
from uuid import uuid1 as get_ident
|
||||
Thread.isAlive = Thread.is_alive
|
||||
have_multiproc = True
|
||||
except ImportError:
|
||||
from threading import Thread
|
||||
from thread import get_ident
|
||||
from array import array
|
||||
from hashlib import md5
|
||||
import time, random, sys, os
|
||||
from random import randint, gauss
|
||||
from optparse import OptionParser
|
||||
|
||||
import avro.ipc as ipc
|
||||
import avro.protocol as protocol
|
||||
from avro.ipc import AvroRemoteException
|
||||
|
||||
L = os.path.abspath(__file__).split(os.path.sep)[:-3]
|
||||
root = os.path.sep.join(L)
|
||||
|
||||
|
||||
parser = OptionParser()
|
||||
parser.add_option('-n', '--num-keys', type="int", dest="numkeys",
|
||||
help="Number of keys", default=1000**2)
|
||||
parser.add_option('-t', '--threads', type="int", dest="threads",
|
||||
help="Number of threads/procs to use", default=50)
|
||||
parser.add_option('-c', '--columns', type="int", dest="columns",
|
||||
help="Number of columns per key", default=5)
|
||||
parser.add_option('-d', '--nodes', type="string", dest="nodes",
|
||||
help="Host nodes (comma separated)", default="localhost")
|
||||
parser.add_option('-s', '--stdev', type="float", dest="stdev", default=0.1,
|
||||
help="standard deviation factor")
|
||||
parser.add_option('-r', '--random', action="store_true", dest="random",
|
||||
help="use random key generator (stdev will have no effect)")
|
||||
parser.add_option('-f', '--file', type="string", dest="file",
|
||||
help="write output to file")
|
||||
parser.add_option('-p', '--port', type="int", default=9160, dest="port",
|
||||
help="thrift port")
|
||||
parser.add_option('-m', '--unframed', action="store_true", dest="unframed",
|
||||
help="use unframed transport")
|
||||
parser.add_option('-o', '--operation', type="choice", dest="operation",
|
||||
default="insert", choices=('insert', 'read', 'rangeslice'),
|
||||
help="operation to perform")
|
||||
parser.add_option('-u', '--supercolumns', type="int", dest="supers", default=1,
|
||||
help="number of super columns per key")
|
||||
parser.add_option('-y', '--family-type', type="choice", dest="cftype",
|
||||
choices=('regular','super'), default='regular',
|
||||
help="column family type")
|
||||
parser.add_option('-k', '--keep-going', action="store_true", dest="ignore",
|
||||
help="ignore errors inserting or reading")
|
||||
parser.add_option('-i', '--progress-interval', type="int", default=10,
|
||||
dest="interval", help="progress report interval (seconds)")
|
||||
parser.add_option('-g', '--get-range-slice-count', type="int", default=1000,
|
||||
dest="rangecount",
|
||||
help="amount of keys to get_range_slices per call")
|
||||
parser.add_option('-l', '--replication-factor', type="int", default=1,
|
||||
dest="replication",
|
||||
help="replication factor to use when creating needed column families")
|
||||
parser.add_option('-e', '--consistency-level', type="str", default='ONE',
|
||||
dest="consistency", help="consistency level to use")
|
||||
|
||||
(options, args) = parser.parse_args()
|
||||
|
||||
total_keys = options.numkeys
|
||||
n_threads = options.threads
|
||||
keys_per_thread = total_keys / n_threads
|
||||
columns_per_key = options.columns
|
||||
supers_per_key = options.supers
|
||||
# this allows client to round robin requests directly for
|
||||
# simple request load-balancing
|
||||
nodes = options.nodes.split(',')
|
||||
|
||||
# a generator that generates all keys according to a bell curve centered
|
||||
# around the middle of the keys generated (0..total_keys). Remember that
|
||||
# about 68% of keys will be within stdev away from the mean and
|
||||
# about 95% within 2*stdev.
|
||||
stdev = total_keys * options.stdev
|
||||
mean = total_keys / 2
|
||||
|
||||
c_levels = ['ZERO', 'ANY', 'ONE', 'QUORUM', 'DCQUORUM', 'DCQUORUMSYNC', 'ALL']
|
||||
consistency = options.consistency
|
||||
if not consistency in c_levels:
|
||||
print "%s is not a valid consistency level" % options.consistency
|
||||
sys.exit(3)
|
||||
|
||||
def key_generator_gauss():
|
||||
fmt = '%0' + str(len(str(total_keys))) + 'd'
|
||||
while True:
|
||||
guess = gauss(mean, stdev)
|
||||
if 0 <= guess < total_keys:
|
||||
return fmt % int(guess)
|
||||
|
||||
# a generator that will generate all keys w/ equal probability. this is the
|
||||
# worst case for caching.
|
||||
def key_generator_random():
|
||||
fmt = '%0' + str(len(str(total_keys))) + 'd'
|
||||
return fmt % randint(0, total_keys - 1)
|
||||
|
||||
key_generator = key_generator_gauss
|
||||
if options.random:
|
||||
key_generator = key_generator_random
|
||||
|
||||
def get_client(host='127.0.0.1', port=9170):
|
||||
schema = os.path.join(root, 'interface/avro', 'cassandra.avpr')
|
||||
proto = protocol.parse(open(schema).read())
|
||||
client = ipc.HTTPTransceiver(host, port)
|
||||
return ipc.Requestor(proto, client)
|
||||
|
||||
def make_keyspaces():
|
||||
keyspace1 = dict()
|
||||
keyspace1['name'] = 'Keyspace1'
|
||||
keyspace1['replication_factor'] = options.replication
|
||||
keyspace1['strategy_class'] = 'org.apache.cassandra.locator.SimpleStrategy'
|
||||
|
||||
keyspace1['cf_defs'] = [{
|
||||
'keyspace': 'Keyspace1',
|
||||
'name': 'Standard1',
|
||||
}]
|
||||
|
||||
keyspace1['cf_defs'].append({
|
||||
'keyspace': 'Keyspace1',
|
||||
'name': 'Super1',
|
||||
'column_type': 'Super',
|
||||
'comparator_type': 'BytesType',
|
||||
'subcomparator_type': 'BytesType',
|
||||
})
|
||||
client = get_client(nodes[0], options.port)
|
||||
try:
|
||||
client.request('system_add_keyspace', {'ks_def': keyspace1})
|
||||
except AvroRemoteException, e:
|
||||
print e
|
||||
client.transceiver.conn.close()
|
||||
|
||||
class Operation(Thread):
|
||||
def __init__(self, i, opcounts, keycounts, latencies):
|
||||
Thread.__init__(self)
|
||||
# generator of the keys to be used
|
||||
self.range = xrange(keys_per_thread * i, keys_per_thread * (i + 1))
|
||||
# we can't use a local counter, since that won't be visible to the parent
|
||||
# under multiprocessing. instead, the parent passes a "opcounts" array
|
||||
# and an index that is our assigned counter.
|
||||
self.idx = i
|
||||
self.opcounts = opcounts
|
||||
# similarly, a shared array for latency and key totals
|
||||
self.latencies = latencies
|
||||
self.keycounts = keycounts
|
||||
# random host for pseudo-load-balancing
|
||||
[hostname] = random.sample(nodes, 1)
|
||||
# open client
|
||||
self.cclient = get_client(hostname, options.port)
|
||||
self.cclient.request('set_keyspace', {'keyspace': 'Keyspace1'})
|
||||
|
||||
class Inserter(Operation):
|
||||
def run(self):
|
||||
data = md5(str(get_ident())).hexdigest()
|
||||
columns = [{'name': 'C' + str(j), 'value': data, 'timestamp': int(time.time() * 1000000)} for j in xrange(columns_per_key)]
|
||||
fmt = '%0' + str(len(str(total_keys))) + 'd'
|
||||
if 'super' == options.cftype:
|
||||
supers = [{'name': 'S' + str(j), 'columns': columns} for j in xrange(supers_per_key)]
|
||||
for i in self.range:
|
||||
key = fmt % i
|
||||
if 'super' == options.cftype:
|
||||
cfmap= {'key': key, 'mutations': {'Super1' : [{'column_or_supercolumn': {'super_column': s}} for s in supers]}}
|
||||
else:
|
||||
cfmap = {'key': key, 'mutations': {'Standard1': [{'column_or_supercolumn': {'column': c}} for c in columns]}}
|
||||
start = time.time()
|
||||
try:
|
||||
self.cclient.request('batch_mutate', {'mutation_map': [cfmap], 'consistency_level': consistency})
|
||||
except KeyboardInterrupt:
|
||||
raise
|
||||
except Exception, e:
|
||||
if options.ignore:
|
||||
print e
|
||||
else:
|
||||
raise
|
||||
self.latencies[self.idx] += time.time() - start
|
||||
self.opcounts[self.idx] += 1
|
||||
self.keycounts[self.idx] += 1
|
||||
|
||||
|
||||
class Reader(Operation):
|
||||
def run(self):
|
||||
p = {'slice_range': {'start': '', 'finish': '', 'reversed': False, 'count': columns_per_key}}
|
||||
if 'super' == options.cftype:
|
||||
for i in xrange(keys_per_thread):
|
||||
key = key_generator()
|
||||
for j in xrange(supers_per_key):
|
||||
parent = {'column_family': 'Super1', 'super_column': 'S' + str(j)}
|
||||
start = time.time()
|
||||
try:
|
||||
r = self.cclient.request('get_slice', {'key': key, 'column_parent': parent, 'predicate': p, 'consistency_level': consistency})
|
||||
if not r: raise RuntimeError("Key %s not found" % key)
|
||||
except KeyboardInterrupt:
|
||||
raise
|
||||
except Exception, e:
|
||||
if options.ignore:
|
||||
print e
|
||||
else:
|
||||
raise
|
||||
self.latencies[self.idx] += time.time() - start
|
||||
self.opcounts[self.idx] += 1
|
||||
self.keycounts[self.idx] += 1
|
||||
else:
|
||||
parent = {'column_family': 'Standard1'}
|
||||
for i in xrange(keys_per_thread):
|
||||
key = key_generator()
|
||||
start = time.time()
|
||||
try:
|
||||
r = self.cclient.request('get_slice', {'key': key, 'column_parent': parent, 'predicate': p, 'consistency_level': consistency})
|
||||
if not r: raise RuntimeError("Key %s not found" % key)
|
||||
except KeyboardInterrupt:
|
||||
raise
|
||||
except Exception, e:
|
||||
if options.ignore:
|
||||
print e
|
||||
else:
|
||||
raise
|
||||
self.latencies[self.idx] += time.time() - start
|
||||
self.opcounts[self.idx] += 1
|
||||
self.keycounts[self.idx] += 1
|
||||
|
||||
class RangeSlicer(Operation):
|
||||
def run(self):
|
||||
begin = self.range[0]
|
||||
end = self.range[-1]
|
||||
current = begin
|
||||
last = current + options.rangecount
|
||||
fmt = '%0' + str(len(str(total_keys))) + 'd'
|
||||
p = {'slice_range': {'start': '', 'finish': '', 'reversed': False, 'count': columns_per_key}}
|
||||
if 'super' == options.cftype:
|
||||
while current < end:
|
||||
keyrange = {'start_key': fmt % current, 'end_key': fmt % last, 'count': options.rangecount}
|
||||
res = []
|
||||
for j in xrange(supers_per_key):
|
||||
parent = {'column_family': 'Super1', 'super_column': 'S' + str(j)}
|
||||
begin = time.time()
|
||||
try:
|
||||
res = self.cclient.request('get_range_slices', {'column_parent': parent, 'predicate': p, 'range': keyrange, 'consistency_level': consistency})
|
||||
if not res: raise RuntimeError("Key %s not found" % key)
|
||||
except KeyboardInterrupt:
|
||||
raise
|
||||
except Exception, e:
|
||||
if options.ignore:
|
||||
print e
|
||||
else:
|
||||
raise
|
||||
self.latencies[self.idx] += time.time() - begin
|
||||
self.opcounts[self.idx] += 1
|
||||
current += len(r) + 1
|
||||
last = current + len(r) + 1
|
||||
self.keycounts[self.idx] += len(r)
|
||||
else:
|
||||
parent = {'column_family': 'Standard1'}
|
||||
while current < end:
|
||||
start = fmt % current
|
||||
finish = fmt % last
|
||||
keyrange = {'start_key': start, 'end_key': finish, 'count': options.rangecount}
|
||||
begin = time.time()
|
||||
try:
|
||||
r = self.cclient.request('get_range_slices', {'column_parent': parent, 'predicate': p, 'range': keyrange, 'consistency_level': consistency})
|
||||
if not r: raise RuntimeError("Range not found:", start, finish)
|
||||
except KeyboardInterrupt:
|
||||
raise
|
||||
except Exception, e:
|
||||
if options.ignore:
|
||||
print e
|
||||
else:
|
||||
print start, finish
|
||||
raise
|
||||
current += len(r) + 1
|
||||
last = current + len(r) + 1
|
||||
self.latencies[self.idx] += time.time() - begin
|
||||
self.opcounts[self.idx] += 1
|
||||
self.keycounts[self.idx] += len(r)
|
||||
|
||||
|
||||
class OperationFactory:
|
||||
@staticmethod
|
||||
def create(type, i, opcounts, keycounts, latencies):
|
||||
if type == 'read':
|
||||
return Reader(i, opcounts, keycounts, latencies)
|
||||
elif type == 'insert':
|
||||
return Inserter(i, opcounts, keycounts, latencies)
|
||||
elif type == 'rangeslice':
|
||||
return RangeSlicer(i, opcounts, keycounts, latencies)
|
||||
else:
|
||||
raise RuntimeError, 'Unsupported op!'
|
||||
|
||||
|
||||
class Stress(object):
|
||||
opcounts = array('i', [0] * n_threads)
|
||||
latencies = array('d', [0] * n_threads)
|
||||
keycounts = array('i', [0] * n_threads)
|
||||
|
||||
def create_threads(self,type):
|
||||
threads = []
|
||||
for i in xrange(n_threads):
|
||||
th = OperationFactory.create(type, i, self.opcounts, self.keycounts, self.latencies)
|
||||
threads.append(th)
|
||||
th.start()
|
||||
return threads
|
||||
|
||||
def run_test(self,filename,threads):
|
||||
start_t = time.time()
|
||||
if filename:
|
||||
outf = open(filename,'w')
|
||||
else:
|
||||
outf = sys.stdout
|
||||
outf.write('total,interval_op_rate,interval_key_rate,avg_latency,elapsed_time\n')
|
||||
epoch = total = old_total = latency = keycount = old_keycount = old_latency = 0
|
||||
epoch_intervals = (options.interval * 10) # 1 epoch = 1 tenth of a second
|
||||
terminate = False
|
||||
while not terminate:
|
||||
time.sleep(0.1)
|
||||
if not [th for th in threads if th.isAlive()]:
|
||||
terminate = True
|
||||
epoch = epoch + 1
|
||||
if terminate or epoch > epoch_intervals:
|
||||
epoch = 0
|
||||
old_total, old_latency, old_keycount = total, latency, keycount
|
||||
total = sum(self.opcounts[th.idx] for th in threads)
|
||||
latency = sum(self.latencies[th.idx] for th in threads)
|
||||
keycount = sum(self.keycounts[th.idx] for th in threads)
|
||||
opdelta = total - old_total
|
||||
keydelta = keycount - old_keycount
|
||||
delta_latency = latency - old_latency
|
||||
if opdelta > 0:
|
||||
delta_formatted = (delta_latency / opdelta)
|
||||
else:
|
||||
delta_formatted = 'NaN'
|
||||
elapsed_t = int(time.time() - start_t)
|
||||
outf.write('%d,%d,%d,%s,%d\n'
|
||||
% (total, opdelta / options.interval, keydelta / options.interval, delta_formatted, elapsed_t))
|
||||
|
||||
def insert(self):
|
||||
threads = self.create_threads('insert')
|
||||
self.run_test(options.file,threads);
|
||||
|
||||
def read(self):
|
||||
threads = self.create_threads('read')
|
||||
self.run_test(options.file,threads);
|
||||
|
||||
def rangeslice(self):
|
||||
threads = self.create_threads('rangeslice')
|
||||
self.run_test(options.file,threads);
|
||||
|
||||
stresser = Stress()
|
||||
benchmark = getattr(stresser, options.operation, None)
|
||||
if not have_multiproc:
|
||||
print """WARNING: multiprocessing not present, threading will be used.
|
||||
Benchmark may not be accurate!"""
|
||||
if options.operation == 'insert':
|
||||
make_keyspaces()
|
||||
benchmark()
|
||||
|
|
@ -1,94 +0,0 @@
|
|||
package org.apache.cassandra.avro;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
import org.apache.avro.util.Utf8;
|
||||
|
||||
public class AvroErrorFactory
|
||||
{
|
||||
public static InvalidRequestException newInvalidRequestException(Utf8 why)
|
||||
{
|
||||
InvalidRequestException exception = new InvalidRequestException();
|
||||
exception.why = why;
|
||||
return exception;
|
||||
}
|
||||
|
||||
public static InvalidRequestException newInvalidRequestException(String why)
|
||||
{
|
||||
return newInvalidRequestException(new Utf8(why));
|
||||
}
|
||||
|
||||
public static InvalidRequestException newInvalidRequestException(Throwable e)
|
||||
{
|
||||
InvalidRequestException exception = newInvalidRequestException(e.getMessage());
|
||||
exception.initCause(e);
|
||||
return exception;
|
||||
}
|
||||
|
||||
public static NotFoundException newNotFoundException(Utf8 why)
|
||||
{
|
||||
NotFoundException exception = new NotFoundException();
|
||||
exception.why = why;
|
||||
return exception;
|
||||
}
|
||||
|
||||
public static NotFoundException newNotFoundException(String why)
|
||||
{
|
||||
return newNotFoundException(new Utf8(why));
|
||||
}
|
||||
|
||||
public static NotFoundException newNotFoundException()
|
||||
{
|
||||
return newNotFoundException(new Utf8());
|
||||
}
|
||||
|
||||
public static TimedOutException newTimedOutException(Utf8 why)
|
||||
{
|
||||
TimedOutException exception = new TimedOutException();
|
||||
exception.why = why;
|
||||
return exception;
|
||||
}
|
||||
|
||||
public static TimedOutException newTimedOutException(String why)
|
||||
{
|
||||
return newTimedOutException(new Utf8(why));
|
||||
}
|
||||
|
||||
public static TimedOutException newTimedOutException()
|
||||
{
|
||||
return newTimedOutException(new Utf8());
|
||||
}
|
||||
|
||||
public static UnavailableException newUnavailableException(Utf8 why)
|
||||
{
|
||||
UnavailableException exception = new UnavailableException();
|
||||
exception.why = why;
|
||||
return exception;
|
||||
}
|
||||
|
||||
public static UnavailableException newUnavailableException(String why)
|
||||
{
|
||||
return newUnavailableException(new Utf8(why));
|
||||
}
|
||||
|
||||
public static UnavailableException newUnavailableException(Throwable t)
|
||||
{
|
||||
UnavailableException exception = newUnavailableException(t.getMessage());
|
||||
exception.initCause(t);
|
||||
return exception;
|
||||
}
|
||||
|
||||
public static UnavailableException newUnavailableException()
|
||||
{
|
||||
return newUnavailableException(new Utf8());
|
||||
}
|
||||
|
||||
public static TokenRange newTokenRange(String startRange, String endRange, List<? extends CharSequence> endpoints)
|
||||
{
|
||||
TokenRange tRange = new TokenRange();
|
||||
tRange.start_token = startRange;
|
||||
tRange.end_token = endRange;
|
||||
tRange.endpoints = (List<CharSequence>) endpoints;
|
||||
return tRange;
|
||||
}
|
||||
}
|
||||
|
|
@ -1,113 +0,0 @@
|
|||
package org.apache.cassandra.avro;
|
||||
/*
|
||||
*
|
||||
* 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.util.List;
|
||||
|
||||
import org.apache.avro.generic.GenericArray;
|
||||
import org.apache.avro.util.Utf8;
|
||||
|
||||
public class AvroRecordFactory
|
||||
{
|
||||
public static Column newColumn(ByteBuffer name, ByteBuffer value, long timestamp)
|
||||
{
|
||||
Column column = new Column();
|
||||
column.name = name;
|
||||
column.value = value;
|
||||
column.timestamp = timestamp;
|
||||
return column;
|
||||
}
|
||||
|
||||
public static Column newColumn(byte[] name, byte[] value, long timestamp)
|
||||
{
|
||||
return newColumn(ByteBuffer.wrap(name), ByteBuffer.wrap(value), timestamp);
|
||||
}
|
||||
|
||||
public static SuperColumn newSuperColumn(ByteBuffer name, List<Column> columns)
|
||||
{
|
||||
SuperColumn column = new SuperColumn();
|
||||
column.name = name;
|
||||
column.columns = columns;
|
||||
return column;
|
||||
}
|
||||
|
||||
public static SuperColumn newSuperColumn(byte[] name, List<Column> columns)
|
||||
{
|
||||
return newSuperColumn(ByteBuffer.wrap(name), columns);
|
||||
}
|
||||
|
||||
public static ColumnOrSuperColumn newColumnOrSuperColumn(Column column)
|
||||
{
|
||||
ColumnOrSuperColumn col = new ColumnOrSuperColumn();
|
||||
col.column = column;
|
||||
return col;
|
||||
}
|
||||
|
||||
public static ColumnOrSuperColumn newColumnOrSuperColumn(SuperColumn superColumn)
|
||||
{
|
||||
ColumnOrSuperColumn column = new ColumnOrSuperColumn();
|
||||
column.super_column = superColumn;
|
||||
return column;
|
||||
}
|
||||
|
||||
public static ColumnPath newColumnPath(String cfName, ByteBuffer superColumn, ByteBuffer column)
|
||||
{
|
||||
ColumnPath cPath = new ColumnPath();
|
||||
cPath.column_family = new Utf8(cfName);
|
||||
cPath.super_column = superColumn;
|
||||
cPath.column = column;
|
||||
return cPath;
|
||||
}
|
||||
|
||||
public static ColumnPath newColumnPath(String cfName, byte[] superColumn, byte[] column)
|
||||
{
|
||||
ByteBuffer wrappedSuperColumn = (superColumn != null) ? ByteBuffer.wrap(superColumn) : null;
|
||||
ByteBuffer wrappedColumn = (column != null) ? ByteBuffer.wrap(column) : null;
|
||||
return newColumnPath(cfName, wrappedSuperColumn, wrappedColumn);
|
||||
}
|
||||
|
||||
public static ColumnParent newColumnParent(String cfName, byte[] superColumn)
|
||||
{
|
||||
ColumnParent cp = new ColumnParent();
|
||||
cp.column_family = new Utf8(cfName);
|
||||
if (superColumn != null)
|
||||
cp.super_column = ByteBuffer.wrap(superColumn);
|
||||
return cp;
|
||||
}
|
||||
|
||||
public static CoscsMapEntry newCoscsMapEntry(ByteBuffer key, GenericArray<ColumnOrSuperColumn> columns)
|
||||
{
|
||||
CoscsMapEntry entry = new CoscsMapEntry();
|
||||
entry.key = key;
|
||||
entry.columns = columns;
|
||||
return entry;
|
||||
}
|
||||
|
||||
public static KeySlice newKeySlice(ByteBuffer key, List<ColumnOrSuperColumn> columns) {
|
||||
KeySlice slice = new KeySlice();
|
||||
slice.key = key;
|
||||
slice.columns = columns;
|
||||
return slice;
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -1,336 +0,0 @@
|
|||
package org.apache.cassandra.avro;
|
||||
/*
|
||||
*
|
||||
* 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.util.Arrays;
|
||||
import java.util.Comparator;
|
||||
import java.util.Set;
|
||||
|
||||
import org.apache.avro.util.Utf8;
|
||||
import org.apache.cassandra.config.DatabaseDescriptor;
|
||||
import org.apache.cassandra.db.ColumnFamily;
|
||||
import org.apache.cassandra.db.ColumnFamilyType;
|
||||
import org.apache.cassandra.db.IColumn;
|
||||
import org.apache.cassandra.db.Table;
|
||||
import org.apache.cassandra.db.marshal.AbstractType;
|
||||
import org.apache.cassandra.db.marshal.MarshalException;
|
||||
import org.apache.cassandra.dht.IPartitioner;
|
||||
import org.apache.cassandra.dht.RandomPartitioner;
|
||||
import org.apache.cassandra.dht.Token;
|
||||
import org.apache.cassandra.service.StorageService;
|
||||
import org.apache.cassandra.utils.FBUtilities;
|
||||
|
||||
import static org.apache.cassandra.avro.AvroErrorFactory.newInvalidRequestException;
|
||||
import static org.apache.cassandra.avro.AvroRecordFactory.newColumnPath;
|
||||
|
||||
/**
|
||||
* The Avro analogue to org.apache.cassandra.service.ThriftValidation
|
||||
*/
|
||||
public class AvroValidation
|
||||
{
|
||||
public static void validateKey(ByteBuffer key) throws InvalidRequestException
|
||||
{
|
||||
if (key == null || key.remaining() == 0)
|
||||
throw newInvalidRequestException("Key may not be empty");
|
||||
|
||||
// check that key can be handled by FBUtilities.writeShortByteArray
|
||||
if (key.remaining() > FBUtilities.MAX_UNSIGNED_SHORT)
|
||||
throw newInvalidRequestException("Key length of " + key.remaining() +
|
||||
" is longer than maximum of " + FBUtilities.MAX_UNSIGNED_SHORT);
|
||||
}
|
||||
|
||||
|
||||
// FIXME: could use method in ThriftValidation
|
||||
static void validateKeyspace(String keyspace) throws KeyspaceNotDefinedException
|
||||
{
|
||||
if (!DatabaseDescriptor.getTables().contains(keyspace))
|
||||
throw new KeyspaceNotDefinedException(new Utf8("Keyspace " + keyspace + " does not exist in this schema."));
|
||||
}
|
||||
|
||||
// FIXME: could use method in ThriftValidation
|
||||
public static ColumnFamilyType validateColumnFamily(String keyspace, String columnFamily) throws InvalidRequestException
|
||||
{
|
||||
if (columnFamily.isEmpty())
|
||||
throw newInvalidRequestException("non-empty columnfamily is required");
|
||||
|
||||
ColumnFamilyType cfType = DatabaseDescriptor.getColumnFamilyType(keyspace, columnFamily);
|
||||
if (cfType == null)
|
||||
throw newInvalidRequestException("unconfigured columnfamily " + columnFamily);
|
||||
|
||||
return cfType;
|
||||
}
|
||||
|
||||
static void validateColumnPath(String keyspace, ColumnPath cp) throws InvalidRequestException
|
||||
{
|
||||
validateKeyspace(keyspace);
|
||||
String column_family = cp.column_family.toString();
|
||||
ColumnFamilyType cfType = validateColumnFamily(keyspace, column_family);
|
||||
|
||||
|
||||
if (cfType == ColumnFamilyType.Standard)
|
||||
{
|
||||
if (cp.super_column != null)
|
||||
throw newInvalidRequestException("supercolumn parameter is invalid for standard CF " + column_family);
|
||||
|
||||
if (cp.column == null)
|
||||
throw newInvalidRequestException("column parameter is not optional for standard CF " + column_family);
|
||||
}
|
||||
else
|
||||
{
|
||||
if (cp.super_column == null)
|
||||
throw newInvalidRequestException("supercolumn parameter is not optional for super CF " + column_family);
|
||||
}
|
||||
|
||||
if (cp.column != null)
|
||||
validateColumns(keyspace, column_family, cp.super_column, Arrays.asList(cp.column));
|
||||
if (cp.super_column != null)
|
||||
validateColumns(keyspace, column_family, null, Arrays.asList(cp.super_column));
|
||||
}
|
||||
|
||||
static void validateColumnParent(String keyspace, ColumnParent parent) throws InvalidRequestException
|
||||
{
|
||||
validateKeyspace(keyspace);
|
||||
String cfName = parent.column_family.toString();
|
||||
ColumnFamilyType cfType = validateColumnFamily(keyspace, cfName);
|
||||
|
||||
if (cfType == ColumnFamilyType.Standard)
|
||||
if (parent.super_column != null)
|
||||
throw newInvalidRequestException("super column specified for standard column family");
|
||||
if (parent.super_column != null)
|
||||
validateColumns(keyspace, cfName, null, Arrays.asList(parent.super_column));
|
||||
}
|
||||
|
||||
// FIXME: could use method in ThriftValidation
|
||||
static void validateColumns(String keyspace, String cfName, ByteBuffer superColumnName, Iterable<ByteBuffer> columnNames)
|
||||
throws InvalidRequestException
|
||||
{
|
||||
if (superColumnName != null)
|
||||
{
|
||||
if (superColumnName.remaining() > IColumn.MAX_NAME_LENGTH)
|
||||
throw newInvalidRequestException("supercolumn name length must not be greater than " + IColumn.MAX_NAME_LENGTH);
|
||||
if (superColumnName.remaining() == 0)
|
||||
throw newInvalidRequestException("supercolumn name must not be empty");
|
||||
if (DatabaseDescriptor.getColumnFamilyType(keyspace, cfName) == ColumnFamilyType.Standard)
|
||||
throw newInvalidRequestException("supercolumn specified to ColumnFamily " + cfName + " containing normal columns");
|
||||
}
|
||||
|
||||
AbstractType comparator = ColumnFamily.getComparatorFor(keyspace, cfName, superColumnName);
|
||||
for (ByteBuffer buff : columnNames)
|
||||
{
|
||||
|
||||
if (buff.remaining() > IColumn.MAX_NAME_LENGTH)
|
||||
throw newInvalidRequestException("column name length must not be greater than " + IColumn.MAX_NAME_LENGTH);
|
||||
if (buff.remaining() == 0)
|
||||
throw newInvalidRequestException("column name must not be empty");
|
||||
|
||||
try
|
||||
{
|
||||
comparator.validate(buff);
|
||||
}
|
||||
catch (MarshalException e)
|
||||
{
|
||||
throw newInvalidRequestException(e.getMessage());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
static void validateColumns(String keyspace, ColumnParent parent, Iterable<ByteBuffer> columnNames)
|
||||
throws InvalidRequestException
|
||||
{
|
||||
validateColumns(keyspace,
|
||||
parent.column_family.toString(),
|
||||
parent.super_column,
|
||||
columnNames);
|
||||
}
|
||||
|
||||
static void validateColumn(String keyspace, ColumnParent parent, Column column)
|
||||
throws InvalidRequestException
|
||||
{
|
||||
validateTtl(column);
|
||||
validateColumns(keyspace, parent, Arrays.asList(column.name));
|
||||
}
|
||||
|
||||
static void validateColumnOrSuperColumn(String keyspace, String cfName, ColumnOrSuperColumn cosc)
|
||||
throws InvalidRequestException
|
||||
{
|
||||
if (cosc.column != null)
|
||||
AvroValidation.validateColumnPath(keyspace, newColumnPath(cfName, null, cosc.column.name));
|
||||
|
||||
if (cosc.super_column != null)
|
||||
for (Column c : cosc.super_column.columns)
|
||||
AvroValidation.validateColumnPath(keyspace, newColumnPath(cfName, cosc.super_column.name, c.name));
|
||||
|
||||
if ((cosc.column == null) && (cosc.super_column == null))
|
||||
throw newInvalidRequestException("ColumnOrSuperColumn must have one or both of Column or SuperColumn");
|
||||
}
|
||||
|
||||
static void validateRange(String keyspace, String cfName, ByteBuffer superName, SliceRange range)
|
||||
throws InvalidRequestException
|
||||
{
|
||||
AbstractType comparator = ColumnFamily.getComparatorFor(keyspace, cfName, superName);
|
||||
|
||||
|
||||
try
|
||||
{
|
||||
comparator.validate(range.start);
|
||||
comparator.validate(range.finish);
|
||||
}
|
||||
catch (MarshalException me)
|
||||
{
|
||||
throw newInvalidRequestException(me.getMessage());
|
||||
}
|
||||
|
||||
if (range.count < 0)
|
||||
throw newInvalidRequestException("Ranges require a non-negative count.");
|
||||
|
||||
Comparator<ByteBuffer> orderedComparator = range.reversed ? comparator.getReverseComparator() : comparator;
|
||||
if (range.start.remaining() > 0 && range.finish.remaining() > 0 && orderedComparator.compare(range.start, range.finish) > 0)
|
||||
throw newInvalidRequestException("range finish must come after start in the order of traversal");
|
||||
}
|
||||
|
||||
static void validateRange(String keyspace, ColumnParent cp, SliceRange range) throws InvalidRequestException
|
||||
{
|
||||
validateRange(keyspace, cp.column_family.toString(), cp.super_column, range);
|
||||
}
|
||||
|
||||
static void validateSlicePredicate(String keyspace, String cfName, ByteBuffer superName, SlicePredicate predicate)
|
||||
throws InvalidRequestException
|
||||
{
|
||||
if (predicate.column_names == null && predicate.slice_range == null)
|
||||
throw newInvalidRequestException("A SlicePredicate must be given a list of Columns, a SliceRange, or both");
|
||||
|
||||
if (predicate.slice_range != null)
|
||||
validateRange(keyspace, cfName, superName, predicate.slice_range);
|
||||
|
||||
if (predicate.column_names != null)
|
||||
validateColumns(keyspace, cfName, superName, predicate.column_names);
|
||||
}
|
||||
|
||||
static void validateDeletion(String keyspace, String cfName, Deletion del) throws InvalidRequestException
|
||||
{
|
||||
validateColumnFamily(keyspace, cfName);
|
||||
if (del.super_column == null && del.predicate == null)
|
||||
throw newInvalidRequestException("A Deletion must have a SuperColumn, a SlicePredicate, or both.");
|
||||
|
||||
if (del.predicate != null)
|
||||
{
|
||||
validateSlicePredicate(keyspace, cfName, del.super_column, del.predicate);
|
||||
if (del.predicate.slice_range != null)
|
||||
throw newInvalidRequestException("Deletion does not yet support SliceRange predicates.");
|
||||
}
|
||||
}
|
||||
|
||||
static void validateMutation(String keyspace, String cfName, Mutation mutation) throws InvalidRequestException
|
||||
{
|
||||
ColumnOrSuperColumn cosc = mutation.column_or_supercolumn;
|
||||
Deletion del = mutation.deletion;
|
||||
|
||||
if (cosc != null && del != null)
|
||||
throw newInvalidRequestException("Mutation may have either a ColumnOrSuperColumn or a Deletion, but not both");
|
||||
|
||||
if (cosc != null)
|
||||
{
|
||||
validateColumnOrSuperColumn(keyspace, cfName, cosc);
|
||||
}
|
||||
else if (del != null)
|
||||
{
|
||||
validateDeletion(keyspace, cfName, del);
|
||||
}
|
||||
else
|
||||
{
|
||||
throw newInvalidRequestException("Mutation must have one ColumnOrSuperColumn, or one Deletion");
|
||||
}
|
||||
}
|
||||
|
||||
static void validateTtl(Column column) throws InvalidRequestException
|
||||
{
|
||||
if (column.ttl != null && column.ttl < 0)
|
||||
throw newInvalidRequestException("ttl must be a positive value");
|
||||
}
|
||||
|
||||
static void validatePredicate(String keyspace, ColumnParent cp, SlicePredicate predicate)
|
||||
throws InvalidRequestException
|
||||
{
|
||||
if (predicate.column_names == null && predicate.slice_range == null)
|
||||
throw newInvalidRequestException("predicate column_names and slice_range may not both be null");
|
||||
|
||||
if (predicate.column_names != null && predicate.slice_range != null)
|
||||
throw newInvalidRequestException("predicate column_names and slice_range may not both be set");
|
||||
|
||||
if (predicate.slice_range != null)
|
||||
validateRange(keyspace, cp, predicate.slice_range);
|
||||
else
|
||||
validateColumns(keyspace, cp, predicate.column_names);
|
||||
}
|
||||
|
||||
public static void validateKeyRange(KeyRange range)
|
||||
throws InvalidRequestException
|
||||
{
|
||||
if ((range.start_key == null) != (range.end_key == null))
|
||||
{
|
||||
throw newInvalidRequestException("start key and end key must either both be non-null, or both be null");
|
||||
}
|
||||
if ((range.start_token == null) != (range.end_token == null))
|
||||
{
|
||||
throw newInvalidRequestException("start token and end token must either both be non-null, or both be null");
|
||||
}
|
||||
if ((range.start_key == null) == (range.start_token == null))
|
||||
{
|
||||
throw newInvalidRequestException("exactly one of {start key, end key} or {start token, end token} must be specified");
|
||||
}
|
||||
|
||||
if (range.start_key != null)
|
||||
{
|
||||
IPartitioner p = StorageService.getPartitioner();
|
||||
Token startToken = p.getToken(range.start_key);
|
||||
Token endToken = p.getToken(range.end_key);
|
||||
if (startToken.compareTo(endToken) > 0 && !endToken.equals(p.getMinimumToken()))
|
||||
{
|
||||
if (p instanceof RandomPartitioner)
|
||||
throw newInvalidRequestException("start key's md5 sorts after end key's md5. this is not allowed; you probably should not specify end key at all, under RandomPartitioner");
|
||||
else
|
||||
throw newInvalidRequestException("start key must sort before (or equal to) finish key in your partitioner!");
|
||||
}
|
||||
}
|
||||
|
||||
if (range.count <= 0)
|
||||
{
|
||||
throw newInvalidRequestException("maxRows must be positive");
|
||||
}
|
||||
}
|
||||
|
||||
static void validateIndexClauses(String keyspace, String columnFamily, IndexClause index_clause)
|
||||
throws InvalidRequestException
|
||||
{
|
||||
if (index_clause.expressions.isEmpty())
|
||||
throw newInvalidRequestException("index clause list may not be empty");
|
||||
Set<ByteBuffer> indexedColumns = Table.open(keyspace).getColumnFamilyStore(columnFamily).getIndexedColumns();
|
||||
for (IndexExpression expression : index_clause.expressions)
|
||||
{
|
||||
if (expression.op.equals(IndexOperator.EQ) && indexedColumns.contains(expression.column_name))
|
||||
return;
|
||||
}
|
||||
throw newInvalidRequestException("No indexed columns present in index clause with operator EQ");
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -1,85 +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.avro;
|
||||
|
||||
import java.io.IOException;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
import org.apache.avro.ipc.ResponderServlet;
|
||||
import org.apache.avro.specific.SpecificResponder;
|
||||
import org.mortbay.jetty.servlet.Context;
|
||||
import org.mortbay.jetty.servlet.ServletHolder;
|
||||
|
||||
/**
|
||||
* The Avro analogue to org.apache.cassandra.service.CassandraDaemon.
|
||||
*
|
||||
*/
|
||||
public class CassandraDaemon extends org.apache.cassandra.service.AbstractCassandraDaemon {
|
||||
private static Logger logger = LoggerFactory.getLogger(CassandraDaemon.class);
|
||||
private org.mortbay.jetty.Server server;
|
||||
|
||||
/** hook for JSVC */
|
||||
public void start() throws IOException
|
||||
{
|
||||
if (logger.isDebugEnabled())
|
||||
logger.debug(String.format("Binding avro service to %s:%s", listenAddr, listenPort));
|
||||
CassandraServer cassandraServer = new CassandraServer();
|
||||
SpecificResponder responder = new SpecificResponder(Cassandra.class, cassandraServer);
|
||||
|
||||
logger.info("Listening for avro clients...");
|
||||
|
||||
// FIXME: This isn't actually binding to listenAddr (it should).
|
||||
server = new org.mortbay.jetty.Server(listenPort);
|
||||
server.setThreadPool(new CleaningThreadPool(cassandraServer.clientState,
|
||||
MIN_WORKER_THREADS,
|
||||
Integer.MAX_VALUE));
|
||||
try
|
||||
{
|
||||
// see CASSANDRA-1440
|
||||
ResponderServlet servlet = new ResponderServlet(responder);
|
||||
new Context(server, "/").addServlet(new ServletHolder(servlet), "/*");
|
||||
|
||||
server.start();
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
throw new IOException("Could not start Avro server.", e);
|
||||
}
|
||||
}
|
||||
|
||||
/** hook for JSVC */
|
||||
public void stop()
|
||||
{
|
||||
logger.info("Cassandra shutting down...");
|
||||
try
|
||||
{
|
||||
server.stop();
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
logger.error("Avro server did not exit cleanly.", e);
|
||||
}
|
||||
}
|
||||
|
||||
public static void main(String[] args) {
|
||||
new CassandraDaemon().activate();
|
||||
}
|
||||
}
|
||||
File diff suppressed because it is too large
Load Diff
|
|
@ -1,34 +0,0 @@
|
|||
package org.apache.cassandra.avro;
|
||||
/*
|
||||
*
|
||||
* 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 org.apache.avro.util.Utf8;
|
||||
|
||||
// XXX: This is an analogue to org.apache.cassandra.db.KeyspaceNotDefinedException
|
||||
@SuppressWarnings("serial")
|
||||
public class KeyspaceNotDefinedException extends InvalidRequestException {
|
||||
|
||||
public KeyspaceNotDefinedException(Utf8 why)
|
||||
{
|
||||
this.why = why;
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue