mirror of https://github.com/apache/cassandra
Merge branch '3628' into cassandra-1.0
This commit is contained in:
commit
cbb809d35a
|
|
@ -67,6 +67,7 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface, Lo
|
|||
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_ALLOW_DELETES = "PIG_ALLOW_DELETES";
|
||||
|
||||
private final static ByteBuffer BOUND = ByteBufferUtil.EMPTY_BYTE_BUFFER;
|
||||
private static final Log logger = LogFactory.getLog(CassandraStorage.class);
|
||||
|
|
@ -74,6 +75,7 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface, Lo
|
|||
private ByteBuffer slice_start = BOUND;
|
||||
private ByteBuffer slice_end = BOUND;
|
||||
private boolean slice_reverse = false;
|
||||
private boolean allow_deletes = false;
|
||||
private String keyspace;
|
||||
private String column_family;
|
||||
private String loadSignature;
|
||||
|
|
@ -331,6 +333,8 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface, Lo
|
|||
ConfigHelper.setPartitioner(conf, System.getenv(PIG_PARTITIONER));
|
||||
else if (ConfigHelper.getPartitioner(conf) == null)
|
||||
throw new IOException("PIG_PARTITIONER environment variable not set");
|
||||
if (System.getenv(PIG_ALLOW_DELETES) != null)
|
||||
allow_deletes = Boolean.valueOf(System.getenv(PIG_ALLOW_DELETES));
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
@ -585,11 +589,15 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface, Lo
|
|||
Mutation mutation = new Mutation();
|
||||
if (t.get(1) == null)
|
||||
{
|
||||
// TODO: optional deletion
|
||||
mutation.deletion = new Deletion();
|
||||
mutation.deletion.predicate = new org.apache.cassandra.thrift.SlicePredicate();
|
||||
mutation.deletion.predicate.column_names = Arrays.asList(objToBB(t.get(0)));
|
||||
mutation.deletion.setTimestamp(FBUtilities.timestampMicros());
|
||||
if (allow_deletes)
|
||||
{
|
||||
mutation.deletion = new Deletion();
|
||||
mutation.deletion.predicate = new org.apache.cassandra.thrift.SlicePredicate();
|
||||
mutation.deletion.predicate.column_names = Arrays.asList(objToBB(t.get(0)));
|
||||
mutation.deletion.setTimestamp(FBUtilities.timestampMicros());
|
||||
}
|
||||
else
|
||||
throw new IOException("null found but deletes are disabled, set " + PIG_ALLOW_DELETES + "=true to enable");
|
||||
}
|
||||
else
|
||||
{
|
||||
|
|
@ -622,11 +630,16 @@ public class CassandraStorage extends LoadFunc implements StoreFuncInterface, Lo
|
|||
column.setTimestamp(FBUtilities.timestampMicros());
|
||||
columns.add(column);
|
||||
}
|
||||
if (columns.isEmpty()) // TODO: optional deletion
|
||||
if (columns.isEmpty())
|
||||
{
|
||||
mutation.deletion = new Deletion();
|
||||
mutation.deletion.super_column = objToBB(pair.get(0));
|
||||
mutation.deletion.setTimestamp(FBUtilities.timestampMicros());
|
||||
if (allow_deletes)
|
||||
{
|
||||
mutation.deletion = new Deletion();
|
||||
mutation.deletion.super_column = objToBB(pair.get(0));
|
||||
mutation.deletion.setTimestamp(FBUtilities.timestampMicros());
|
||||
}
|
||||
else
|
||||
throw new IOException("SuperColumn deletion attempted with empty bag, but deletes are disabled, set " + PIG_ALLOW_DELETES + "=true to enable");
|
||||
}
|
||||
else
|
||||
{
|
||||
|
|
|
|||
Loading…
Reference in New Issue