From 9c4c63ccae2dee3cdead7d2e83928e26d8a84a48 Mon Sep 17 00:00:00 2001 From: Brandon Williams Date: Wed, 28 Sep 2011 22:01:08 +0000 Subject: [PATCH] Fix handling of integer types in pig. Patch by brandonwilliams, reviewed by Jeremy Hanna for CASSANDRA-2810 git-svn-id: https://svn.apache.org/repos/asf/cassandra/branches/cassandra-0.8@1177084 13f79535-47bb-0310-9956-ffa450edef68 --- .../hadoop/pig/CassandraStorage.java | 36 +++++++++---------- 1 file changed, 18 insertions(+), 18 deletions(-) diff --git a/contrib/pig/src/java/org/apache/cassandra/hadoop/pig/CassandraStorage.java b/contrib/pig/src/java/org/apache/cassandra/hadoop/pig/CassandraStorage.java index 76e8698e1f..044ea3f33b 100644 --- a/contrib/pig/src/java/org/apache/cassandra/hadoop/pig/CassandraStorage.java +++ b/contrib/pig/src/java/org/apache/cassandra/hadoop/pig/CassandraStorage.java @@ -17,11 +17,13 @@ package org.apache.cassandra.hadoop.pig; import java.io.IOException; +import java.math.BigInteger; import java.nio.ByteBuffer; import java.util.*; import org.apache.cassandra.config.ConfigurationException; import org.apache.cassandra.db.marshal.BytesType; +import org.apache.cassandra.db.marshal.IntegerType; import org.apache.cassandra.db.marshal.TypeParser; import org.apache.cassandra.thrift.*; import org.apache.cassandra.utils.FBUtilities; @@ -143,18 +145,14 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface, Lo List marshallers = getDefaultMarshallers(cfDef); Map validators = getValidatorMap(cfDef); + setTupleValue(pair, 0, marshallers.get(0).compose(name)); if (col instanceof Column) { // standard - pair.set(0, marshallers.get(0).compose(name)); if (validators.get(name) == null) - // Have to special case BytesType because compose returns a ByteBuffer - if (marshallers.get(1) instanceof BytesType) - pair.set(1, new DataByteArray(ByteBufferUtil.getArray(col.value()))); - else - pair.set(1, marshallers.get(1).compose(col.value())); + setTupleValue(pair, 1, marshallers.get(1).compose(col.value())); else - pair.set(1, validators.get(name).compose(col.value())); + setTupleValue(pair, 1, validators.get(name).compose(col.value())); return pair; } @@ -167,6 +165,16 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface, Lo return pair; } + private 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 + pair.set(position, value); + } + private CfDef getCfDef(String signature) { UDFContext context = UDFContext.getUDFContext(); @@ -453,8 +461,6 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface, Lo DefaultDataBag pairs = (DefaultDataBag) t.get(1); ArrayList mutationList = new ArrayList(); CfDef cfDef = getCfDef(storeSignature); - List marshallers = getDefaultMarshallers(cfDef); - Map validators = getValidatorMap(cfDef); try { for (Tuple pair : pairs) @@ -498,15 +504,8 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface, Lo else { org.apache.cassandra.thrift.Column column = new org.apache.cassandra.thrift.Column(); - column.name = marshallers.get(0).decompose((pair.get(0))); - if (validators.get(column.name) == null) - // Have to special case BytesType to convert DataByteArray into ByteBuffer - if (marshallers.get(1) instanceof BytesType) - column.value = objToBB(pair.get(1)); - else - column.value = marshallers.get(1).decompose(pair.get(1)); - else - column.value = validators.get(column.name).decompose(pair.get(1)); + column.name = objToBB(pair.get(0)); + column.value = objToBB(pair.get(1)); column.setTimestamp(System.currentTimeMillis() * 1000); mutation.column_or_supercolumn = new ColumnOrSuperColumn(); mutation.column_or_supercolumn.column = column; @@ -626,3 +625,4 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface, Lo return cfDef; } } +