mirror of https://github.com/apache/cassandra
Pig: support for cql3 tables
Patch by Alex Liu, reviewed by brandonwilliams for CASSANDRA-5234
This commit is contained in:
parent
6def8223f5
commit
764bcd3f3d
|
|
@ -218,22 +218,40 @@ public class ColumnFamilyRecordReader extends RecordReader<ByteBuffer, SortedMap
|
|||
|
||||
private RowIterator()
|
||||
{
|
||||
CfDef cfDef = new CfDef();
|
||||
try
|
||||
{
|
||||
partitioner = FBUtilities.newPartitioner(client.describe_partitioner());
|
||||
partitioner = FBUtilities.newPartitioner(client.describe_partitioner());
|
||||
// get CF meta data
|
||||
String query = "SELECT comparator," +
|
||||
" subcomparator," +
|
||||
" type " +
|
||||
"FROM system.schema_columnfamilies " +
|
||||
"WHERE keyspace_name = '%s' " +
|
||||
" AND columnfamily_name = '%s' ";
|
||||
|
||||
// Get the Keyspace metadata, then get the specific CF metadata
|
||||
// in order to populate the sub/comparator.
|
||||
KsDef ks_def = client.describe_keyspace(keyspace);
|
||||
List<String> cfnames = new ArrayList<String>(ks_def.cf_defs.size());
|
||||
for (CfDef cfd : ks_def.cf_defs)
|
||||
cfnames.add(cfd.name);
|
||||
int idx = cfnames.indexOf(cfName);
|
||||
CfDef cf_def = ks_def.cf_defs.get(idx);
|
||||
CqlResult result = client.execute_cql3_query(
|
||||
ByteBufferUtil.bytes(String.format(query, keyspace, cfName)),
|
||||
Compression.NONE,
|
||||
ConsistencyLevel.ONE);
|
||||
|
||||
isSuper = cf_def.column_type.equals("Super");
|
||||
comparator = TypeParser.parse(cf_def.comparator_type);
|
||||
subComparator = cf_def.subcomparator_type == null ? null : TypeParser.parse(cf_def.subcomparator_type);
|
||||
Iterator<CqlRow> iteraRow = result.rows.iterator();
|
||||
|
||||
if (iteraRow.hasNext())
|
||||
{
|
||||
CqlRow cqlRow = iteraRow.next();
|
||||
cfDef.comparator_type = ByteBufferUtil.string(cqlRow.columns.get(0).value);
|
||||
ByteBuffer subComparator = cqlRow.columns.get(1).value;
|
||||
if (subComparator != null)
|
||||
cfDef.subcomparator_type = ByteBufferUtil.string(subComparator);
|
||||
|
||||
ByteBuffer type = cqlRow.columns.get(2).value;
|
||||
if (type != null)
|
||||
cfDef.column_type = ByteBufferUtil.string(type);
|
||||
}
|
||||
|
||||
comparator = TypeParser.parse(cfDef.comparator_type);
|
||||
subComparator = cfDef.subcomparator_type == null ? null : TypeParser.parse(cfDef.subcomparator_type);
|
||||
}
|
||||
catch (ConfigurationException e)
|
||||
{
|
||||
|
|
@ -247,6 +265,7 @@ public class ColumnFamilyRecordReader extends RecordReader<ByteBuffer, SortedMap
|
|||
{
|
||||
throw new RuntimeException("unable to load keyspace " + keyspace, e);
|
||||
}
|
||||
isSuper = "Super".equalsIgnoreCase(cfDef.column_type);
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
|
|||
|
|
@ -0,0 +1,744 @@
|
|||
/*
|
||||
* 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.hadoop.pig;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.math.BigInteger;
|
||||
import java.nio.ByteBuffer;
|
||||
import java.nio.charset.CharacterCodingException;
|
||||
import java.util.*;
|
||||
|
||||
|
||||
import org.apache.cassandra.exceptions.ConfigurationException;
|
||||
import org.apache.cassandra.exceptions.SyntaxException;
|
||||
import org.apache.cassandra.auth.IAuthenticator;
|
||||
import org.apache.cassandra.db.Column;
|
||||
import org.apache.cassandra.db.marshal.*;
|
||||
import org.apache.cassandra.db.marshal.AbstractCompositeType.CompositeComponent;
|
||||
import org.apache.cassandra.hadoop.*;
|
||||
import org.apache.cassandra.thrift.*;
|
||||
import org.apache.cassandra.utils.ByteBufferUtil;
|
||||
import org.apache.cassandra.utils.FBUtilities;
|
||||
import org.apache.cassandra.utils.Hex;
|
||||
import org.apache.cassandra.utils.UUIDGen;
|
||||
|
||||
import org.apache.hadoop.conf.Configuration;
|
||||
import org.apache.hadoop.fs.Path;
|
||||
import org.apache.hadoop.mapreduce.*;
|
||||
import org.apache.pig.*;
|
||||
import org.apache.pig.backend.executionengine.ExecException;
|
||||
import org.apache.pig.data.*;
|
||||
import org.apache.pig.impl.util.UDFContext;
|
||||
import org.apache.thrift.TDeserializer;
|
||||
import org.apache.thrift.TException;
|
||||
import org.apache.thrift.TSerializer;
|
||||
import org.apache.thrift.protocol.TBinaryProtocol;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
/**
|
||||
* A LoadStoreFunc for retrieving data from and storing data to Cassandra
|
||||
*/
|
||||
public abstract class AbstractCassandraStorage extends LoadFunc implements StoreFuncInterface, LoadMetadata
|
||||
{
|
||||
protected enum MarshallerType { COMPARATOR, DEFAULT_VALIDATOR, KEY_VALIDATOR, SUBCOMPARATOR };
|
||||
|
||||
// system environment variables that can be set to configure connection info:
|
||||
// alternatively, Hadoop JobConf variables can be set using keys from ConfigHelper
|
||||
public final static String PIG_INPUT_RPC_PORT = "PIG_INPUT_RPC_PORT";
|
||||
public final static String PIG_INPUT_INITIAL_ADDRESS = "PIG_INPUT_INITIAL_ADDRESS";
|
||||
public final static String PIG_INPUT_PARTITIONER = "PIG_INPUT_PARTITIONER";
|
||||
public final static String PIG_OUTPUT_RPC_PORT = "PIG_OUTPUT_RPC_PORT";
|
||||
public final static String PIG_OUTPUT_INITIAL_ADDRESS = "PIG_OUTPUT_INITIAL_ADDRESS";
|
||||
public final static String PIG_OUTPUT_PARTITIONER = "PIG_OUTPUT_PARTITIONER";
|
||||
public final static String PIG_RPC_PORT = "PIG_RPC_PORT";
|
||||
public final static String PIG_INITIAL_ADDRESS = "PIG_INITIAL_ADDRESS";
|
||||
public final static String PIG_PARTITIONER = "PIG_PARTITIONER";
|
||||
public final static String PIG_INPUT_FORMAT = "PIG_INPUT_FORMAT";
|
||||
public final static String PIG_OUTPUT_FORMAT = "PIG_OUTPUT_FORMAT";
|
||||
public final static String PIG_INPUT_SPLIT_SIZE = "PIG_INPUT_SPLIT_SIZE";
|
||||
|
||||
protected String DEFAULT_INPUT_FORMAT;
|
||||
protected String DEFAULT_OUTPUT_FORMAT;
|
||||
|
||||
protected static final Logger logger = LoggerFactory.getLogger(AbstractCassandraStorage.class);
|
||||
|
||||
protected String username;
|
||||
protected String password;
|
||||
protected String keyspace;
|
||||
protected String column_family;
|
||||
protected String loadSignature;
|
||||
protected String storeSignature;
|
||||
|
||||
protected Configuration conf;
|
||||
protected String inputFormatClass;
|
||||
protected String outputFormatClass;
|
||||
protected int splitSize = 64 * 1024;
|
||||
protected String partitionerClass;
|
||||
|
||||
public AbstractCassandraStorage()
|
||||
{
|
||||
super();
|
||||
}
|
||||
|
||||
/** Deconstructs a composite type to a Tuple. */
|
||||
protected Tuple composeComposite(AbstractCompositeType comparator, ByteBuffer name) throws IOException
|
||||
{
|
||||
List<CompositeComponent> result = comparator.deconstruct(name);
|
||||
Tuple t = TupleFactory.getInstance().newTuple(result.size());
|
||||
for (int i=0; i<result.size(); i++)
|
||||
setTupleValue(t, i, result.get(i).comparator.compose(result.get(i).value));
|
||||
|
||||
return t;
|
||||
}
|
||||
|
||||
/** convert a column to a tuple */
|
||||
protected Tuple columnToTuple(Column col, CfDef cfDef, AbstractType comparator) throws IOException
|
||||
{
|
||||
Tuple pair = TupleFactory.getInstance().newTuple(2);
|
||||
|
||||
// name
|
||||
if(comparator instanceof AbstractCompositeType)
|
||||
setTupleValue(pair, 0, composeComposite((AbstractCompositeType)comparator,col.name()));
|
||||
else
|
||||
setTupleValue(pair, 0, comparator.compose(col.name()));
|
||||
|
||||
// value
|
||||
Map<ByteBuffer,AbstractType> validators = getValidatorMap(cfDef);
|
||||
if (validators.get(col.name()) == null)
|
||||
{
|
||||
Map<MarshallerType, AbstractType> marshallers = getDefaultMarshallers(cfDef);
|
||||
setTupleValue(pair, 1, marshallers.get(MarshallerType.DEFAULT_VALIDATOR).compose(col.value()));
|
||||
}
|
||||
else
|
||||
setTupleValue(pair, 1, validators.get(col.name()).compose(col.value()));
|
||||
return pair;
|
||||
}
|
||||
|
||||
/** set the value to the position of the tuple */
|
||||
protected void setTupleValue(Tuple pair, int position, Object value) throws ExecException
|
||||
{
|
||||
if (value instanceof BigInteger)
|
||||
pair.set(position, ((BigInteger) value).intValue());
|
||||
else if (value instanceof ByteBuffer)
|
||||
pair.set(position, new DataByteArray(ByteBufferUtil.getArray((ByteBuffer) value)));
|
||||
else if (value instanceof UUID)
|
||||
pair.set(position, new DataByteArray(UUIDGen.decompose((java.util.UUID) value)));
|
||||
else if (value instanceof Date)
|
||||
pair.set(position, DateType.instance.decompose((Date) value).getLong());
|
||||
else
|
||||
pair.set(position, value);
|
||||
}
|
||||
|
||||
/** get the columnfamily definition for the signature */
|
||||
protected CfDef getCfDef(String signature)
|
||||
{
|
||||
UDFContext context = UDFContext.getUDFContext();
|
||||
Properties property = context.getUDFProperties(AbstractCassandraStorage.class);
|
||||
return cfdefFromString(property.getProperty(signature));
|
||||
}
|
||||
|
||||
/** construct a map to store the mashaller type to cassandra data type mapping */
|
||||
protected Map<MarshallerType, AbstractType> getDefaultMarshallers(CfDef cfDef) throws IOException
|
||||
{
|
||||
Map<MarshallerType, AbstractType> marshallers = new EnumMap<MarshallerType, AbstractType>(MarshallerType.class);
|
||||
AbstractType comparator;
|
||||
AbstractType subcomparator;
|
||||
AbstractType default_validator;
|
||||
AbstractType key_validator;
|
||||
|
||||
comparator = parseType(cfDef.getComparator_type());
|
||||
subcomparator = parseType(cfDef.getSubcomparator_type());
|
||||
default_validator = parseType(cfDef.getDefault_validation_class());
|
||||
key_validator = parseType(cfDef.getKey_validation_class());
|
||||
|
||||
marshallers.put(MarshallerType.COMPARATOR, comparator);
|
||||
marshallers.put(MarshallerType.DEFAULT_VALIDATOR, default_validator);
|
||||
marshallers.put(MarshallerType.KEY_VALIDATOR, key_validator);
|
||||
marshallers.put(MarshallerType.SUBCOMPARATOR, subcomparator);
|
||||
return marshallers;
|
||||
}
|
||||
|
||||
/** get the validators */
|
||||
protected Map<ByteBuffer, AbstractType> getValidatorMap(CfDef cfDef) throws IOException
|
||||
{
|
||||
Map<ByteBuffer, AbstractType> validators = new HashMap<ByteBuffer, AbstractType>();
|
||||
for (ColumnDef cd : cfDef.getColumn_metadata())
|
||||
{
|
||||
if (cd.getValidation_class() != null && !cd.getValidation_class().isEmpty())
|
||||
{
|
||||
AbstractType validator = null;
|
||||
try
|
||||
{
|
||||
validator = TypeParser.parse(cd.getValidation_class());
|
||||
validators.put(cd.name, validator);
|
||||
}
|
||||
catch (ConfigurationException e)
|
||||
{
|
||||
throw new IOException(e);
|
||||
}
|
||||
catch (SyntaxException e)
|
||||
{
|
||||
throw new IOException(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
return validators;
|
||||
}
|
||||
|
||||
/** parse the string to a cassandra data type */
|
||||
protected AbstractType parseType(String type) throws IOException
|
||||
{
|
||||
try
|
||||
{
|
||||
// always treat counters like longs, specifically CCT.compose is not what we need
|
||||
if (type != null && type.equals("org.apache.cassandra.db.marshal.CounterColumnType"))
|
||||
return LongType.instance;
|
||||
return TypeParser.parse(type);
|
||||
}
|
||||
catch (ConfigurationException e)
|
||||
{
|
||||
throw new IOException(e);
|
||||
}
|
||||
catch (SyntaxException e)
|
||||
{
|
||||
throw new IOException(e);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
public InputFormat getInputFormat()
|
||||
{
|
||||
try
|
||||
{
|
||||
return FBUtilities.construct(inputFormatClass, "inputformat");
|
||||
}
|
||||
catch (ConfigurationException e)
|
||||
{
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
|
||||
/** decompose the query to store the parameters in a map*/
|
||||
public static Map<String, String> getQueryMap(String query)
|
||||
{
|
||||
String[] params = query.split("&");
|
||||
Map<String, String> map = new HashMap<String, String>();
|
||||
for (String param : params)
|
||||
{
|
||||
String[] keyValue = param.split("=");
|
||||
map.put(keyValue[0], keyValue[1]);
|
||||
}
|
||||
return map;
|
||||
}
|
||||
|
||||
/** set hadoop cassandra connection settings */
|
||||
protected void setConnectionInformation() throws IOException
|
||||
{
|
||||
if (System.getenv(PIG_RPC_PORT) != null)
|
||||
{
|
||||
ConfigHelper.setInputRpcPort(conf, System.getenv(PIG_RPC_PORT));
|
||||
ConfigHelper.setOutputRpcPort(conf, System.getenv(PIG_RPC_PORT));
|
||||
}
|
||||
|
||||
if (System.getenv(PIG_INPUT_RPC_PORT) != null)
|
||||
ConfigHelper.setInputRpcPort(conf, System.getenv(PIG_INPUT_RPC_PORT));
|
||||
if (System.getenv(PIG_OUTPUT_RPC_PORT) != null)
|
||||
ConfigHelper.setOutputRpcPort(conf, System.getenv(PIG_OUTPUT_RPC_PORT));
|
||||
|
||||
if (System.getenv(PIG_INITIAL_ADDRESS) != null)
|
||||
{
|
||||
ConfigHelper.setInputInitialAddress(conf, System.getenv(PIG_INITIAL_ADDRESS));
|
||||
ConfigHelper.setOutputInitialAddress(conf, System.getenv(PIG_INITIAL_ADDRESS));
|
||||
}
|
||||
if (System.getenv(PIG_INPUT_INITIAL_ADDRESS) != null)
|
||||
ConfigHelper.setInputInitialAddress(conf, System.getenv(PIG_INPUT_INITIAL_ADDRESS));
|
||||
if (System.getenv(PIG_OUTPUT_INITIAL_ADDRESS) != null)
|
||||
ConfigHelper.setOutputInitialAddress(conf, System.getenv(PIG_OUTPUT_INITIAL_ADDRESS));
|
||||
|
||||
if (System.getenv(PIG_PARTITIONER) != null)
|
||||
{
|
||||
ConfigHelper.setInputPartitioner(conf, System.getenv(PIG_PARTITIONER));
|
||||
ConfigHelper.setOutputPartitioner(conf, System.getenv(PIG_PARTITIONER));
|
||||
}
|
||||
if(System.getenv(PIG_INPUT_PARTITIONER) != null)
|
||||
ConfigHelper.setInputPartitioner(conf, System.getenv(PIG_INPUT_PARTITIONER));
|
||||
if(System.getenv(PIG_OUTPUT_PARTITIONER) != null)
|
||||
ConfigHelper.setOutputPartitioner(conf, System.getenv(PIG_OUTPUT_PARTITIONER));
|
||||
if (System.getenv(PIG_INPUT_FORMAT) != null)
|
||||
inputFormatClass = getFullyQualifiedClassName(System.getenv(PIG_INPUT_FORMAT));
|
||||
else
|
||||
inputFormatClass = DEFAULT_INPUT_FORMAT;
|
||||
if (System.getenv(PIG_OUTPUT_FORMAT) != null)
|
||||
outputFormatClass = getFullyQualifiedClassName(System.getenv(PIG_OUTPUT_FORMAT));
|
||||
else
|
||||
outputFormatClass = DEFAULT_OUTPUT_FORMAT;
|
||||
}
|
||||
|
||||
/** get the full class name */
|
||||
protected String getFullyQualifiedClassName(String classname)
|
||||
{
|
||||
return classname.contains(".") ? classname : "org.apache.cassandra.hadoop." + classname;
|
||||
}
|
||||
|
||||
/** get pig type for the cassandra data type*/
|
||||
protected byte getPigType(AbstractType type)
|
||||
{
|
||||
if (type instanceof LongType || type instanceof DateType) // DateType is bad and it should feel bad
|
||||
return DataType.LONG;
|
||||
else if (type instanceof IntegerType || type instanceof Int32Type) // IntegerType will overflow at 2**31, but is kept for compatibility until pig has a BigInteger
|
||||
return DataType.INTEGER;
|
||||
else if (type instanceof AsciiType)
|
||||
return DataType.CHARARRAY;
|
||||
else if (type instanceof UTF8Type)
|
||||
return DataType.CHARARRAY;
|
||||
else if (type instanceof FloatType)
|
||||
return DataType.FLOAT;
|
||||
else if (type instanceof DoubleType)
|
||||
return DataType.DOUBLE;
|
||||
else if (type instanceof AbstractCompositeType )
|
||||
return DataType.TUPLE;
|
||||
|
||||
return DataType.BYTEARRAY;
|
||||
}
|
||||
|
||||
public ResourceStatistics getStatistics(String location, Job job)
|
||||
{
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String relativeToAbsolutePath(String location, Path curDir) throws IOException
|
||||
{
|
||||
return location;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void setUDFContextSignature(String signature)
|
||||
{
|
||||
this.loadSignature = signature;
|
||||
}
|
||||
|
||||
/** StoreFunc methods */
|
||||
public void setStoreFuncUDFContextSignature(String signature)
|
||||
{
|
||||
this.storeSignature = signature;
|
||||
}
|
||||
|
||||
public String relToAbsPathForStoreLocation(String location, Path curDir) throws IOException
|
||||
{
|
||||
return relativeToAbsolutePath(location, curDir);
|
||||
}
|
||||
|
||||
/** output format */
|
||||
public OutputFormat getOutputFormat()
|
||||
{
|
||||
try
|
||||
{
|
||||
return FBUtilities.construct(outputFormatClass, "outputformat");
|
||||
}
|
||||
catch (ConfigurationException e)
|
||||
{
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
|
||||
public void checkSchema(ResourceSchema schema) throws IOException
|
||||
{
|
||||
// we don't care about types, they all get casted to ByteBuffers
|
||||
}
|
||||
|
||||
/** convert object to ByteBuffer */
|
||||
protected ByteBuffer objToBB(Object o)
|
||||
{
|
||||
if (o == null)
|
||||
return (ByteBuffer)o;
|
||||
if (o instanceof java.lang.String)
|
||||
return ByteBuffer.wrap(new DataByteArray((String)o).get());
|
||||
if (o instanceof Integer)
|
||||
return Int32Type.instance.decompose((Integer)o);
|
||||
if (o instanceof Long)
|
||||
return LongType.instance.decompose((Long)o);
|
||||
if (o instanceof Float)
|
||||
return FloatType.instance.decompose((Float)o);
|
||||
if (o instanceof Double)
|
||||
return DoubleType.instance.decompose((Double)o);
|
||||
if (o instanceof UUID)
|
||||
return ByteBuffer.wrap(UUIDGen.decompose((UUID) o));
|
||||
if(o instanceof Tuple) {
|
||||
List<Object> objects = ((Tuple)o).getAll();
|
||||
List<ByteBuffer> serialized = new ArrayList<ByteBuffer>(objects.size());
|
||||
int totalLength = 0;
|
||||
for(Object sub : objects)
|
||||
{
|
||||
ByteBuffer buffer = objToBB(sub);
|
||||
serialized.add(buffer);
|
||||
totalLength += 2 + buffer.remaining() + 1;
|
||||
}
|
||||
ByteBuffer out = ByteBuffer.allocate(totalLength);
|
||||
for (ByteBuffer bb : serialized)
|
||||
{
|
||||
int length = bb.remaining();
|
||||
out.put((byte) ((length >> 8) & 0xFF));
|
||||
out.put((byte) (length & 0xFF));
|
||||
out.put(bb);
|
||||
out.put((byte) 0);
|
||||
}
|
||||
out.flip();
|
||||
return out;
|
||||
}
|
||||
|
||||
return ByteBuffer.wrap(((DataByteArray) o).get());
|
||||
}
|
||||
|
||||
public void cleanupOnFailure(String failure, Job job)
|
||||
{
|
||||
}
|
||||
|
||||
/** Methods to get the column family schema from Cassandra */
|
||||
protected void initSchema(String signature)
|
||||
{
|
||||
Properties properties = UDFContext.getUDFContext().getUDFProperties(AbstractCassandraStorage.class);
|
||||
|
||||
// Only get the schema if we haven't already gotten it
|
||||
if (!properties.containsKey(signature))
|
||||
{
|
||||
try
|
||||
{
|
||||
Cassandra.Client client = ConfigHelper.getClientFromInputAddressList(conf);
|
||||
client.set_keyspace(keyspace);
|
||||
|
||||
if (username != null && password != null)
|
||||
{
|
||||
Map<String, String> credentials = new HashMap<String, String>(2);
|
||||
credentials.put(IAuthenticator.USERNAME_KEY, username);
|
||||
credentials.put(IAuthenticator.PASSWORD_KEY, password);
|
||||
|
||||
try
|
||||
{
|
||||
client.login(new AuthenticationRequest(credentials));
|
||||
}
|
||||
catch (AuthenticationException e)
|
||||
{
|
||||
logger.error("Authentication exception: invalid username and/or password");
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
catch (AuthorizationException e)
|
||||
{
|
||||
throw new AssertionError(e); // never actually throws AuthorizationException.
|
||||
}
|
||||
}
|
||||
|
||||
// compose the CfDef for the columfamily
|
||||
CfDef cfDef = getCfDef(client);
|
||||
|
||||
if (cfDef != null)
|
||||
properties.setProperty(signature, cfdefToString(cfDef));
|
||||
else
|
||||
throw new RuntimeException(String.format("Column family '%s' not found in keyspace '%s'",
|
||||
column_family,
|
||||
keyspace));
|
||||
}
|
||||
catch (TException e)
|
||||
{
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
catch (CharacterCodingException e)
|
||||
{
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
catch (IOException e)
|
||||
{
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** convert CfDef to string */
|
||||
protected static String cfdefToString(CfDef cfDef)
|
||||
{
|
||||
assert cfDef != null;
|
||||
// this is so awful it's kind of cool!
|
||||
TSerializer serializer = new TSerializer(new TBinaryProtocol.Factory());
|
||||
try
|
||||
{
|
||||
return Hex.bytesToHex(serializer.serialize(cfDef));
|
||||
}
|
||||
catch (TException e)
|
||||
{
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
}
|
||||
|
||||
/** convert string back to CfDef */
|
||||
protected static CfDef cfdefFromString(String st)
|
||||
{
|
||||
assert st != null;
|
||||
TDeserializer deserializer = new TDeserializer(new TBinaryProtocol.Factory());
|
||||
CfDef cfDef = new CfDef();
|
||||
try
|
||||
{
|
||||
deserializer.deserialize(cfDef, Hex.hexToBytes(st));
|
||||
}
|
||||
catch (TException e)
|
||||
{
|
||||
throw new RuntimeException(e);
|
||||
}
|
||||
return cfDef;
|
||||
}
|
||||
|
||||
/** return the CfDef for the column family */
|
||||
protected CfDef getCfDef(Cassandra.Client client)
|
||||
throws InvalidRequestException,
|
||||
UnavailableException,
|
||||
TimedOutException,
|
||||
SchemaDisagreementException,
|
||||
TException,
|
||||
CharacterCodingException
|
||||
{
|
||||
// get CF meta data
|
||||
String query = "SELECT type, " +
|
||||
" comparator," +
|
||||
" subcomparator," +
|
||||
" default_validator, " +
|
||||
" key_validator," +
|
||||
" key_aliases " +
|
||||
"FROM system.schema_columnfamilies " +
|
||||
"WHERE keyspace_name = '%s' " +
|
||||
" AND columnfamily_name = '%s' ";
|
||||
|
||||
CqlResult result = client.execute_cql3_query(
|
||||
ByteBufferUtil.bytes(String.format(query, keyspace, column_family)),
|
||||
Compression.NONE,
|
||||
ConsistencyLevel.ONE);
|
||||
|
||||
if (result == null || result.rows == null || result.rows.isEmpty())
|
||||
return null;
|
||||
|
||||
Iterator<CqlRow> iteraRow = result.rows.iterator();
|
||||
CfDef cfDef = new CfDef();
|
||||
cfDef.keyspace = keyspace;
|
||||
cfDef.name = column_family;
|
||||
boolean cql3Table = false;
|
||||
if (iteraRow.hasNext())
|
||||
{
|
||||
CqlRow cqlRow = iteraRow.next();
|
||||
|
||||
cfDef.column_type = ByteBufferUtil.string(cqlRow.columns.get(0).value);
|
||||
cfDef.comparator_type = ByteBufferUtil.string(cqlRow.columns.get(1).value);
|
||||
ByteBuffer subComparator = cqlRow.columns.get(2).value;
|
||||
if (subComparator != null)
|
||||
cfDef.subcomparator_type = ByteBufferUtil.string(subComparator);
|
||||
cfDef.default_validation_class = ByteBufferUtil.string(cqlRow.columns.get(3).value);
|
||||
cfDef.key_validation_class = ByteBufferUtil.string(cqlRow.columns.get(4).value);
|
||||
List<String> keys = null;
|
||||
if (cqlRow.columns.get(5).value != null)
|
||||
{
|
||||
String keyAliases = ByteBufferUtil.string(cqlRow.columns.get(5).value);
|
||||
keys = FBUtilities.fromJsonList(keyAliases);
|
||||
}
|
||||
// get column meta data
|
||||
if (keys != null && keys.size() > 0)
|
||||
cql3Table = true;
|
||||
}
|
||||
cfDef.column_metadata = getColumnMetadata(client, cql3Table);
|
||||
return cfDef;
|
||||
}
|
||||
|
||||
/** get a list of columns */
|
||||
protected abstract List<ColumnDef> getColumnMetadata(Cassandra.Client client, boolean cql3Table)
|
||||
throws InvalidRequestException,
|
||||
UnavailableException,
|
||||
TimedOutException,
|
||||
SchemaDisagreementException,
|
||||
TException,
|
||||
CharacterCodingException;
|
||||
|
||||
/** get column meta data */
|
||||
protected List<ColumnDef> getColumnMeta(Cassandra.Client client)
|
||||
throws InvalidRequestException,
|
||||
UnavailableException,
|
||||
TimedOutException,
|
||||
SchemaDisagreementException,
|
||||
TException,
|
||||
CharacterCodingException
|
||||
{
|
||||
String query = "SELECT column_name, " +
|
||||
" validator, " +
|
||||
" index_type " +
|
||||
"FROM system.schema_columns " +
|
||||
"WHERE keyspace_name = '%s' " +
|
||||
" AND columnfamily_name = '%s'";
|
||||
|
||||
CqlResult result = client.execute_cql3_query(
|
||||
ByteBufferUtil.bytes(String.format(query, keyspace, column_family)),
|
||||
Compression.NONE,
|
||||
ConsistencyLevel.ONE);
|
||||
|
||||
List<CqlRow> rows = result.rows;
|
||||
List<ColumnDef> columnDefs = new ArrayList<ColumnDef>();
|
||||
if (rows == null || rows.isEmpty())
|
||||
return columnDefs;
|
||||
|
||||
Iterator<CqlRow> iterator = rows.iterator();
|
||||
while (iterator.hasNext())
|
||||
{
|
||||
CqlRow row = iterator.next();
|
||||
ColumnDef cDef = new ColumnDef();
|
||||
cDef.setName(ByteBufferUtil.clone(row.getColumns().get(0).value));
|
||||
cDef.validation_class = ByteBufferUtil.string(row.getColumns().get(1).value);
|
||||
ByteBuffer indexType = row.getColumns().get(2).value;
|
||||
if (indexType != null)
|
||||
cDef.index_type = getIndexType(ByteBufferUtil.string(indexType));
|
||||
columnDefs.add(cDef);
|
||||
}
|
||||
return columnDefs;
|
||||
}
|
||||
|
||||
/** get keys meta data */
|
||||
protected List<ColumnDef> getKeysMeta(Cassandra.Client client)
|
||||
throws InvalidRequestException,
|
||||
UnavailableException,
|
||||
TimedOutException,
|
||||
SchemaDisagreementException,
|
||||
TException,
|
||||
IOException
|
||||
{
|
||||
String query = "SELECT key_aliases, " +
|
||||
" column_aliases, " +
|
||||
" key_validator, " +
|
||||
" comparator, " +
|
||||
" keyspace_name, " +
|
||||
" value_alias, " +
|
||||
" default_validator " +
|
||||
"FROM system.schema_columnfamilies " +
|
||||
"WHERE keyspace_name = '%s'" +
|
||||
" AND columnfamily_name = '%s' ";
|
||||
|
||||
CqlResult result = client.execute_cql3_query(
|
||||
ByteBufferUtil.bytes(String.format(query, keyspace, column_family)),
|
||||
Compression.NONE,
|
||||
ConsistencyLevel.ONE);
|
||||
|
||||
if (result == null || result.rows == null || result.rows.isEmpty())
|
||||
return null;
|
||||
|
||||
List<CqlRow> rows = result.rows;
|
||||
Iterator<CqlRow> iteraRow = rows.iterator();
|
||||
List<ColumnDef> keys = new ArrayList<ColumnDef>();
|
||||
if (iteraRow.hasNext())
|
||||
{
|
||||
CqlRow cqlRow = iteraRow.next();
|
||||
String name = ByteBufferUtil.string(cqlRow.columns.get(4).value);
|
||||
logger.debug("Found ksDef name: {}", name);
|
||||
String keyString = ByteBufferUtil.string(ByteBuffer.wrap(cqlRow.columns.get(0).getValue()));
|
||||
|
||||
logger.debug("partition keys: " + keyString);
|
||||
List<String> keyNames = FBUtilities.fromJsonList(keyString);
|
||||
|
||||
Iterator<String> iterator = keyNames.iterator();
|
||||
while (iterator.hasNext())
|
||||
{
|
||||
ColumnDef cDef = new ColumnDef();
|
||||
cDef.name = ByteBufferUtil.bytes(iterator.next());
|
||||
keys.add(cDef);
|
||||
}
|
||||
|
||||
keyString = ByteBufferUtil.string(ByteBuffer.wrap(cqlRow.columns.get(1).getValue()));
|
||||
|
||||
logger.debug("cluster keys: " + keyString);
|
||||
keyNames = FBUtilities.fromJsonList(keyString);
|
||||
|
||||
iterator = keyNames.iterator();
|
||||
while (iterator.hasNext())
|
||||
{
|
||||
ColumnDef cDef = new ColumnDef();
|
||||
cDef.name = ByteBufferUtil.bytes(iterator.next());
|
||||
keys.add(cDef);
|
||||
}
|
||||
|
||||
String validator = ByteBufferUtil.string(ByteBuffer.wrap(cqlRow.columns.get(2).getValue()));
|
||||
logger.debug("row key validator: " + validator);
|
||||
AbstractType<?> keyValidator = parseType(validator);
|
||||
|
||||
Iterator<ColumnDef> keyItera = keys.iterator();
|
||||
if (keyValidator instanceof CompositeType)
|
||||
{
|
||||
Iterator<AbstractType<?>> typeItera = ((CompositeType) keyValidator).types.iterator();
|
||||
while (typeItera.hasNext())
|
||||
keyItera.next().validation_class = typeItera.next().toString();
|
||||
}
|
||||
else
|
||||
keyItera.next().validation_class = keyValidator.toString();
|
||||
|
||||
validator = ByteBufferUtil.string(ByteBuffer.wrap(cqlRow.columns.get(3).getValue()));
|
||||
logger.debug("cluster key validator: " + validator);
|
||||
|
||||
if (keyItera.hasNext() && validator != null && !validator.isEmpty())
|
||||
{
|
||||
AbstractType<?> clusterKeyValidator = parseType(validator);
|
||||
|
||||
if (clusterKeyValidator instanceof CompositeType)
|
||||
{
|
||||
Iterator<AbstractType<?>> typeItera = ((CompositeType) clusterKeyValidator).types.iterator();
|
||||
while (keyItera.hasNext())
|
||||
keyItera.next().validation_class = typeItera.next().toString();
|
||||
}
|
||||
else
|
||||
keyItera.next().validation_class = clusterKeyValidator.toString();
|
||||
}
|
||||
|
||||
// compact value_alias column
|
||||
if (cqlRow.columns.get(5).value != null)
|
||||
{
|
||||
try
|
||||
{
|
||||
String compactValidator = ByteBufferUtil.string(ByteBuffer.wrap(cqlRow.columns.get(6).getValue()));
|
||||
logger.debug("default validator: " + compactValidator);
|
||||
AbstractType<?> defaultValidator = parseType(compactValidator);
|
||||
|
||||
ColumnDef cDef = new ColumnDef();
|
||||
cDef.name = cqlRow.columns.get(5).value;
|
||||
cDef.validation_class = defaultValidator.toString();
|
||||
keys.add(cDef);
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
// no compact column at value_alias
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
return keys;
|
||||
}
|
||||
|
||||
/** get index type from string */
|
||||
protected IndexType getIndexType(String type)
|
||||
{
|
||||
type = type.toLowerCase();
|
||||
if ("keys".equals(type))
|
||||
return IndexType.KEYS;
|
||||
else if("custom".equals(type))
|
||||
return IndexType.CUSTOM;
|
||||
else if("composites".equals(type))
|
||||
return IndexType.COMPOSITES;
|
||||
else
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
|
|
@ -0,0 +1,446 @@
|
|||
/*
|
||||
* 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.hadoop.pig;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.ByteBuffer;
|
||||
import java.nio.charset.CharacterCodingException;
|
||||
import java.util.*;
|
||||
|
||||
|
||||
import org.apache.cassandra.db.Column;
|
||||
import org.apache.cassandra.db.marshal.*;
|
||||
import org.apache.cassandra.hadoop.*;
|
||||
import org.apache.cassandra.hadoop.cql3.CqlConfigHelper;
|
||||
import org.apache.cassandra.thrift.*;
|
||||
import org.apache.cassandra.utils.ByteBufferUtil;
|
||||
import org.apache.hadoop.mapreduce.*;
|
||||
import org.apache.pig.Expression;
|
||||
import org.apache.pig.ResourceSchema;
|
||||
import org.apache.pig.ResourceSchema.ResourceFieldSchema;
|
||||
import org.apache.pig.backend.hadoop.executionengine.mapReduceLayer.PigSplit;
|
||||
import org.apache.pig.data.*;
|
||||
import org.apache.thrift.TException;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
/**
|
||||
* A LoadStoreFunc for retrieving data from and storing data to Cassandra
|
||||
*
|
||||
* A row from a standard CF will be returned as nested tuples:
|
||||
* (((key1, value1), (key2, value2)), ((name1, val1), (name2, val2))).
|
||||
*/
|
||||
public class CqlStorage extends AbstractCassandraStorage
|
||||
{
|
||||
private static final Logger logger = LoggerFactory.getLogger(CqlStorage.class);
|
||||
|
||||
private RecordReader<Map<String, ByteBuffer>, Map<String, ByteBuffer>> reader;
|
||||
private RecordWriter<Map<String, ByteBuffer>, List<ByteBuffer>> writer;
|
||||
|
||||
private int pageSize = 1000;
|
||||
private String columns;
|
||||
private String outputQuery;
|
||||
private String whereClause;
|
||||
|
||||
public CqlStorage()
|
||||
{
|
||||
this(1000);
|
||||
}
|
||||
|
||||
/** @param limit number of CQL rows to fetch in a thrift request */
|
||||
public CqlStorage(int pageSize)
|
||||
{
|
||||
super();
|
||||
this.pageSize = pageSize;
|
||||
DEFAULT_INPUT_FORMAT = "org.apache.cassandra.hadoop.cql3.CqlPagingInputFormat";
|
||||
DEFAULT_OUTPUT_FORMAT = "org.apache.cassandra.hadoop.cql3.CqlOutputFormat";
|
||||
}
|
||||
|
||||
public void prepareToRead(RecordReader reader, PigSplit split)
|
||||
{
|
||||
this.reader = reader;
|
||||
}
|
||||
|
||||
/** get next row */
|
||||
public Tuple getNext() throws IOException
|
||||
{
|
||||
try
|
||||
{
|
||||
// load the next pair
|
||||
if (!reader.nextKeyValue())
|
||||
return null;
|
||||
|
||||
CfDef cfDef = getCfDef(loadSignature);
|
||||
Map<String, ByteBuffer> keys = reader.getCurrentKey();
|
||||
Map<String, ByteBuffer> columns = reader.getCurrentValue();
|
||||
assert keys != null && columns != null;
|
||||
|
||||
// add key columns to the map
|
||||
for (Map.Entry<String,ByteBuffer> key : keys.entrySet())
|
||||
columns.put(key.getKey(), key.getValue());
|
||||
|
||||
Tuple tuple = TupleFactory.getInstance().newTuple(cfDef.column_metadata.size());
|
||||
Iterator<ColumnDef> itera = cfDef.column_metadata.iterator();
|
||||
int i = 0;
|
||||
while (itera.hasNext())
|
||||
{
|
||||
ColumnDef cdef = itera.next();
|
||||
ByteBuffer columnValue = columns.get(ByteBufferUtil.string(cdef.name.duplicate()));
|
||||
if (columnValue != null)
|
||||
{
|
||||
Column column = new Column(cdef.name, columnValue);
|
||||
tuple.set(i, columnToTuple(column, cfDef, UTF8Type.instance));
|
||||
}
|
||||
else
|
||||
tuple.set(i, TupleFactory.getInstance().newTuple());
|
||||
i++;
|
||||
}
|
||||
return tuple;
|
||||
}
|
||||
catch (InterruptedException e)
|
||||
{
|
||||
throw new IOException(e.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
/** set read configuration settings */
|
||||
public void setLocation(String location, Job job) throws IOException
|
||||
{
|
||||
conf = job.getConfiguration();
|
||||
setLocationFromUri(location);
|
||||
|
||||
if (username != null && password != null)
|
||||
ConfigHelper.setInputKeyspaceUserNameAndPassword(conf, username, password);
|
||||
if (splitSize > 0)
|
||||
ConfigHelper.setInputSplitSize(conf, splitSize);
|
||||
if (partitionerClass!= null)
|
||||
ConfigHelper.setInputPartitioner(conf, partitionerClass);
|
||||
|
||||
ConfigHelper.setInputColumnFamily(conf, keyspace, column_family);
|
||||
setConnectionInformation();
|
||||
|
||||
CqlConfigHelper.setInputCQLPageRowSize(conf, String.valueOf(pageSize));
|
||||
if (columns != null && !columns.trim().isEmpty())
|
||||
CqlConfigHelper.setInputColumns(conf, columns);
|
||||
if (whereClause != null && !whereClause.trim().isEmpty())
|
||||
CqlConfigHelper.setInputWhereClauses(conf, whereClause);
|
||||
|
||||
if (System.getenv(PIG_INPUT_SPLIT_SIZE) != null)
|
||||
{
|
||||
try
|
||||
{
|
||||
ConfigHelper.setInputSplitSize(conf, Integer.valueOf(System.getenv(PIG_INPUT_SPLIT_SIZE)));
|
||||
}
|
||||
catch (NumberFormatException e)
|
||||
{
|
||||
throw new RuntimeException("PIG_INPUT_SPLIT_SIZE is not a number", e);
|
||||
}
|
||||
}
|
||||
|
||||
if (ConfigHelper.getInputRpcPort(conf) == 0)
|
||||
throw new IOException("PIG_INPUT_RPC_PORT or PIG_RPC_PORT environment variable not set");
|
||||
if (ConfigHelper.getInputInitialAddress(conf) == null)
|
||||
throw new IOException("PIG_INPUT_INITIAL_ADDRESS or PIG_INITIAL_ADDRESS environment variable not set");
|
||||
if (ConfigHelper.getInputPartitioner(conf) == null)
|
||||
throw new IOException("PIG_INPUT_PARTITIONER or PIG_PARTITIONER environment variable not set");
|
||||
if (loadSignature == null)
|
||||
loadSignature = location;
|
||||
|
||||
initSchema(loadSignature);
|
||||
}
|
||||
|
||||
/** set store configuration settings */
|
||||
public void setStoreLocation(String location, Job job) throws IOException
|
||||
{
|
||||
conf = job.getConfiguration();
|
||||
setLocationFromUri(location);
|
||||
|
||||
if (username != null && password != null)
|
||||
ConfigHelper.setOutputKeyspaceUserNameAndPassword(conf, username, password);
|
||||
if (splitSize > 0)
|
||||
ConfigHelper.setInputSplitSize(conf, splitSize);
|
||||
if (partitionerClass!= null)
|
||||
ConfigHelper.setOutputPartitioner(conf, partitionerClass);
|
||||
|
||||
ConfigHelper.setOutputColumnFamily(conf, keyspace, column_family);
|
||||
CqlConfigHelper.setOutputCql(conf, outputQuery);
|
||||
|
||||
setConnectionInformation();
|
||||
|
||||
if (ConfigHelper.getOutputRpcPort(conf) == 0)
|
||||
throw new IOException("PIG_OUTPUT_RPC_PORT or PIG_RPC_PORT environment variable not set");
|
||||
if (ConfigHelper.getOutputInitialAddress(conf) == null)
|
||||
throw new IOException("PIG_OUTPUT_INITIAL_ADDRESS or PIG_INITIAL_ADDRESS environment variable not set");
|
||||
if (ConfigHelper.getOutputPartitioner(conf) == null)
|
||||
throw new IOException("PIG_OUTPUT_PARTITIONER or PIG_PARTITIONER environment variable not set");
|
||||
|
||||
initSchema(storeSignature);
|
||||
}
|
||||
|
||||
/** schema: ((name, value), (name, value), (name, value)) where keys are in the front. */
|
||||
public ResourceSchema getSchema(String location, Job job) throws IOException
|
||||
{
|
||||
setLocation(location, job);
|
||||
CfDef cfDef = getCfDef(loadSignature);
|
||||
|
||||
// top-level schema, no type
|
||||
ResourceSchema schema = new ResourceSchema();
|
||||
|
||||
// get default marshallers and validators
|
||||
Map<MarshallerType, AbstractType> marshallers = getDefaultMarshallers(cfDef);
|
||||
Map<ByteBuffer, AbstractType> validators = getValidatorMap(cfDef);
|
||||
|
||||
// will contain all fields for this schema
|
||||
List<ResourceFieldSchema> allSchemaFields = new ArrayList<ResourceFieldSchema>();
|
||||
|
||||
// defined validators/indexes
|
||||
for (ColumnDef cdef : cfDef.column_metadata)
|
||||
{
|
||||
// make a new tuple for each col/val pair
|
||||
ResourceSchema innerTupleSchema = new ResourceSchema();
|
||||
ResourceFieldSchema innerTupleField = new ResourceFieldSchema();
|
||||
innerTupleField.setType(DataType.TUPLE);
|
||||
innerTupleField.setSchema(innerTupleSchema);
|
||||
innerTupleField.setName(new String(cdef.getName()));
|
||||
ResourceFieldSchema idxColSchema = new ResourceFieldSchema();
|
||||
idxColSchema.setName("name");
|
||||
idxColSchema.setType(getPigType(UTF8Type.instance));
|
||||
|
||||
ResourceFieldSchema valSchema = new ResourceFieldSchema();
|
||||
AbstractType validator = validators.get(cdef.name);
|
||||
if (validator == null)
|
||||
validator = marshallers.get(MarshallerType.DEFAULT_VALIDATOR);
|
||||
valSchema.setName("value");
|
||||
valSchema.setType(getPigType(validator));
|
||||
|
||||
innerTupleSchema.setFields(new ResourceFieldSchema[] { idxColSchema, valSchema });
|
||||
allSchemaFields.add(innerTupleField);
|
||||
}
|
||||
|
||||
// top level schema contains everything
|
||||
schema.setFields(allSchemaFields.toArray(new ResourceFieldSchema[allSchemaFields.size()]));
|
||||
return schema;
|
||||
}
|
||||
|
||||
|
||||
/** We use CQL3 where clause to define the partition, so do nothing here*/
|
||||
public String[] getPartitionKeys(String location, Job job)
|
||||
{
|
||||
return null;
|
||||
}
|
||||
|
||||
/** We use CQL3 where clause to define the partition, so do nothing here*/
|
||||
public void setPartitionFilter(Expression partitionFilter)
|
||||
{
|
||||
}
|
||||
|
||||
public void prepareToWrite(RecordWriter writer)
|
||||
{
|
||||
this.writer = writer;
|
||||
}
|
||||
|
||||
/** output: (((name, value), (name, value)), (value ... value), (value...value)) */
|
||||
public void putNext(Tuple t) throws IOException
|
||||
{
|
||||
if (t.size() < 1)
|
||||
{
|
||||
// simply nothing here, we can't even delete without a key
|
||||
logger.warn("Empty output skipped, filter empty tuples to suppress this warning");
|
||||
return;
|
||||
}
|
||||
|
||||
if (t.getType(0) == DataType.TUPLE)
|
||||
{
|
||||
Map<String, ByteBuffer> key = tupleToKeyMap((Tuple)t.get(0));
|
||||
if (t.getType(1) == DataType.TUPLE)
|
||||
cqlQueryFromTuple(key, t, 1);
|
||||
else
|
||||
throw new IOException("Second argument in output must be a tuple");
|
||||
}
|
||||
else
|
||||
throw new IOException("First argument in output must be a tuple");
|
||||
}
|
||||
|
||||
/** convert key tuple to key map */
|
||||
private Map<String, ByteBuffer> tupleToKeyMap(Tuple t) throws IOException
|
||||
{
|
||||
Map<String, ByteBuffer> keys = new HashMap<String, ByteBuffer>();
|
||||
for (int i = 0; i < t.size(); i++)
|
||||
{
|
||||
if (t.getType(i) == DataType.TUPLE)
|
||||
{
|
||||
Tuple inner = (Tuple) t.get(i);
|
||||
if (inner.size() == 2)
|
||||
{
|
||||
Object name = inner.get(0);
|
||||
if (name != null)
|
||||
{
|
||||
keys.put(name.toString(), objToBB(inner.get(1)));
|
||||
}
|
||||
else
|
||||
throw new IOException("Key name was empty");
|
||||
}
|
||||
else
|
||||
throw new IOException("Keys were not in name and value pairs");
|
||||
}
|
||||
else
|
||||
{
|
||||
throw new IOException("keys was not a tuple");
|
||||
}
|
||||
}
|
||||
return keys;
|
||||
}
|
||||
|
||||
/** send CQL query request using data from tuple */
|
||||
private void cqlQueryFromTuple(Map<String, ByteBuffer> key, Tuple t, int offset) throws IOException
|
||||
{
|
||||
for (int i = offset; i < t.size(); i++)
|
||||
{
|
||||
if (t.getType(i) == DataType.TUPLE)
|
||||
{
|
||||
Tuple inner = (Tuple) t.get(i);
|
||||
if (inner.size() > 0)
|
||||
{
|
||||
|
||||
List<ByteBuffer> bindedVariables = bindedVariablesFromTuple(inner);
|
||||
if (bindedVariables.size() > 0)
|
||||
sendCqlQuery(key, bindedVariables);
|
||||
else
|
||||
throw new IOException("Missing binded variables");
|
||||
}
|
||||
}
|
||||
else
|
||||
{
|
||||
throw new IOException("Output type was not a tuple");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** compose a list of binded variables */
|
||||
private List<ByteBuffer> bindedVariablesFromTuple(Tuple t) throws IOException
|
||||
{
|
||||
List<ByteBuffer> variables = new ArrayList<ByteBuffer>();
|
||||
for (int i = 0; i < t.size(); i++)
|
||||
variables.add(objToBB(t.get(i)));
|
||||
return variables;
|
||||
}
|
||||
|
||||
/** writer write the data by executing CQL query */
|
||||
private void sendCqlQuery(Map<String, ByteBuffer> key, List<ByteBuffer> bindedVariables) throws IOException
|
||||
{
|
||||
try
|
||||
{
|
||||
writer.write(key, bindedVariables);
|
||||
}
|
||||
catch (InterruptedException e)
|
||||
{
|
||||
throw new IOException(e);
|
||||
}
|
||||
}
|
||||
|
||||
/** include key columns */
|
||||
protected List<ColumnDef> getColumnMetadata(Cassandra.Client client, boolean cql3Table)
|
||||
throws InvalidRequestException,
|
||||
UnavailableException,
|
||||
TimedOutException,
|
||||
SchemaDisagreementException,
|
||||
TException,
|
||||
CharacterCodingException
|
||||
{
|
||||
List<ColumnDef> keyColumns = null;
|
||||
// get key columns
|
||||
try
|
||||
{
|
||||
keyColumns = getKeysMeta(client);
|
||||
}
|
||||
catch(IOException e)
|
||||
{
|
||||
logger.error("Error in retrieving key columns" , e);
|
||||
}
|
||||
|
||||
// get other columns
|
||||
List<ColumnDef> columns = getColumnMeta(client);
|
||||
|
||||
// combine all columns in a list
|
||||
if (keyColumns != null && columns != null)
|
||||
keyColumns.addAll(columns);
|
||||
|
||||
return keyColumns;
|
||||
}
|
||||
|
||||
/** cql://[username:password@]<keyspace>/<columnfamily>[?[page_size=<size>]
|
||||
* [&columns=<col1,col2>][&output_query=<prepared_statement_query>][&where_clause=<clause>]
|
||||
* [&split_size=<size>][&partitioner=<partitioner>]] */
|
||||
private void setLocationFromUri(String location) throws IOException
|
||||
{
|
||||
try
|
||||
{
|
||||
if (!location.startsWith("cql://"))
|
||||
throw new Exception("Bad scheme: " + location);
|
||||
|
||||
String[] urlParts = location.split("\\?");
|
||||
if (urlParts.length > 1)
|
||||
{
|
||||
Map<String, String> urlQuery = getQueryMap(urlParts[1]);
|
||||
|
||||
// each page row size
|
||||
if (urlQuery.containsKey("page_size"))
|
||||
pageSize = Integer.parseInt(urlQuery.get("page_size"));
|
||||
|
||||
// input query select columns
|
||||
if (urlQuery.containsKey("columns"))
|
||||
columns = urlQuery.get("columns");
|
||||
|
||||
// output prepared statement
|
||||
if (urlQuery.containsKey("output_query"))
|
||||
outputQuery = urlQuery.get("output_query").replaceAll("#", "?").replaceAll("@", "=");
|
||||
|
||||
// user defined where clause
|
||||
if (urlQuery.containsKey("where_clause"))
|
||||
whereClause = urlQuery.get("where_clause");
|
||||
|
||||
//split size
|
||||
if (urlQuery.containsKey("split_size"))
|
||||
splitSize = Integer.parseInt(urlQuery.get("split_size"));
|
||||
if (urlQuery.containsKey("partitioner"))
|
||||
partitionerClass = urlQuery.get("partitioner");
|
||||
}
|
||||
String[] parts = urlParts[0].split("/+");
|
||||
String[] credentialsAndKeyspace = parts[1].split("@");
|
||||
if (credentialsAndKeyspace.length > 1)
|
||||
{
|
||||
String[] credentials = credentialsAndKeyspace[0].split(":");
|
||||
username = credentials[0];
|
||||
password = credentials[1];
|
||||
keyspace = credentialsAndKeyspace[1];
|
||||
}
|
||||
else
|
||||
{
|
||||
keyspace = parts[1];
|
||||
}
|
||||
column_family = parts[2];
|
||||
}
|
||||
catch (Exception e)
|
||||
{
|
||||
throw new IOException("Expected 'cql://[username:password@]<keyspace>/<columnfamily>" +
|
||||
"[?[page_size=<size>][&columns=<col1,col2>][&output_query=<prepared_statement>]" +
|
||||
"[&where_clause=<clause>][&split_size=<size>][&partitioner=<partitioner>]]': " + e.getMessage());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Loading…
Reference in New Issue