diff --git a/CHANGES.txt b/CHANGES.txt index dc912f1a8c..8c41f50ae3 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -5,6 +5,8 @@ 3.1 +Merged from 2.2: + * (Hadoop) ensure that Cluster instances are always closed (CASSANDRA-10058) Merged from 2.1: * Reject counter writes in CQLSSTableWriter (CASSANDRA-10258) * Remove superfluous COUNTER_MUTATION stage mapping (CASSANDRA-10605) diff --git a/src/java/org/apache/cassandra/hadoop/cql3/CqlRecordWriter.java b/src/java/org/apache/cassandra/hadoop/cql3/CqlRecordWriter.java index 6b4caa5eea..23beba3953 100644 --- a/src/java/org/apache/cassandra/hadoop/cql3/CqlRecordWriter.java +++ b/src/java/org/apache/cassandra/hadoop/cql3/CqlRecordWriter.java @@ -112,27 +112,25 @@ class CqlRecordWriter extends RecordWriter, List(); - try + try (Cluster cluster = CqlConfigHelper.getOutputCluster(ConfigHelper.getOutputInitialAddress(conf), conf)) { String keyspace = ConfigHelper.getOutputKeyspace(conf); - try (Session client = CqlConfigHelper.getOutputCluster(ConfigHelper.getOutputInitialAddress(conf), conf).connect(keyspace)) + Session client = cluster.connect(keyspace); + ringCache = new NativeRingCache(conf); + if (client != null) { - ringCache = new NativeRingCache(conf); - if (client != null) - { - TableMetadata tableMetadata = client.getCluster().getMetadata().getKeyspace(client.getLoggedKeyspace()).getTable(ConfigHelper.getOutputColumnFamily(conf)); - clusterColumns = tableMetadata.getClusteringColumns(); - partitionKeyColumns = tableMetadata.getPartitionKey(); + TableMetadata tableMetadata = client.getCluster().getMetadata().getKeyspace(client.getLoggedKeyspace()).getTable(ConfigHelper.getOutputColumnFamily(conf)); + clusterColumns = tableMetadata.getClusteringColumns(); + partitionKeyColumns = tableMetadata.getPartitionKey(); - String cqlQuery = CqlConfigHelper.getOutputCql(conf).trim(); - if (cqlQuery.toLowerCase().startsWith("insert")) - throw new UnsupportedOperationException("INSERT with CqlRecordWriter is not supported, please use UPDATE/DELETE statement"); - cql = appendKeyWhereClauses(cqlQuery); - } - else - { - throw new IllegalArgumentException("Invalid configuration specified " + conf); - } + String cqlQuery = CqlConfigHelper.getOutputCql(conf).trim(); + if (cqlQuery.toLowerCase().startsWith("insert")) + throw new UnsupportedOperationException("INSERT with CqlRecordWriter is not supported, please use UPDATE/DELETE statement"); + cql = appendKeyWhereClauses(cqlQuery); + } + else + { + throw new IllegalArgumentException("Invalid configuration specified " + conf); } } catch (Exception e) @@ -234,7 +232,7 @@ class CqlRecordWriter extends RecordWriter, List endpoints; - protected Session client; + protected Cluster cluster = null; // A bounded queue of incoming mutations for this range protected final BlockingQueue> queue = new ArrayBlockingQueue>(queueSize); @@ -280,6 +278,7 @@ class CqlRecordWriter extends RecordWriter, List, List= batchThreshold) - break; - bindVariables = queue.poll(); + if (i >= batchThreshold) + break; + bindVariables = queue.poll(); + } + break; } - break; - } - catch (Exception e) - { - closeInternal(); - if (!iter.hasNext()) + catch (Exception e) { - lastException = new IOException(e); - break outer; + closeInternal(); + if (!iter.hasNext()) + { + lastException = new IOException(e); + break outer; + } } } @@ -333,7 +335,8 @@ class CqlRecordWriter extends RecordWriter, List, List, List(); metadata = session.getCluster().getMetadata(); Set ranges = metadata.getTokenRanges(); for (TokenRange range : ranges) - { rangeMap.put(range, metadata.getReplicas(keyspace, range)); - } } } diff --git a/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java b/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java index 602a1086c5..557bebaa7d 100644 --- a/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/CQLSSTableWriterTest.java @@ -158,7 +158,7 @@ public class CQLSSTableWriterTest String insert = String.format("UPDATE cql_keyspace.counter1 SET my_counter = my_counter - ? WHERE my_id = ?"); CQLSSTableWriter.builder().inDirectory(dataDir) .forTable(schema) - .withPartitioner(StorageService.instance.getPartitioner()) + .withPartitioner(Murmur3Partitioner.instance) .using(insert).build(); }