mirror of https://github.com/apache/cassandra
Merge branch 'cassandra-5.0' into trunk
* cassandra-5.0: Support null column value tombstones in FQL batch statements
This commit is contained in:
commit
df96494801
|
|
@ -279,6 +279,7 @@ Merged from 4.1:
|
||||||
* Reduce info logging from automatic paxos repair (CASSANDRA-19445)
|
* Reduce info logging from automatic paxos repair (CASSANDRA-19445)
|
||||||
* Support legacy plain_text_auth section in credentials file removed unintentionally (CASSANDRA-19498)
|
* Support legacy plain_text_auth section in credentials file removed unintentionally (CASSANDRA-19498)
|
||||||
Merged from 4.0:
|
Merged from 4.0:
|
||||||
|
* Support null column value tombstones in FQL batch statements (CASSANDRA-20397)
|
||||||
* Fix millisecond and microsecond precision for commit log replay (CASSANDRA-19448)
|
* Fix millisecond and microsecond precision for commit log replay (CASSANDRA-19448)
|
||||||
* Improve accuracy of memtable heap usage tracking (CASSANDRA-17298)
|
* Improve accuracy of memtable heap usage tracking (CASSANDRA-17298)
|
||||||
* Fix rendering UNSET collection types in query tracing (CASSANDRA-19880)
|
* Fix rendering UNSET collection types in query tracing (CASSANDRA-19880)
|
||||||
|
|
|
||||||
|
|
@ -454,9 +454,7 @@ public class FullQueryLogger implements QueryEvents.Listener
|
||||||
{
|
{
|
||||||
valueOut.int32(subValues.size());
|
valueOut.int32(subValues.size());
|
||||||
for (ByteBuffer value : subValues)
|
for (ByteBuffer value : subValues)
|
||||||
{
|
valueOut.bytes(value == null ? null : BytesStore.wrap(value));
|
||||||
valueOut.bytes(BytesStore.wrap(value));
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -16,7 +16,7 @@
|
||||||
* limitations under the License.
|
* limitations under the License.
|
||||||
*/
|
*/
|
||||||
|
|
||||||
package org.apache.cassandra.distributed.test;
|
package org.apache.cassandra.distributed.test.fql;
|
||||||
|
|
||||||
import org.junit.Ignore;
|
import org.junit.Ignore;
|
||||||
import org.junit.Rule;
|
import org.junit.Rule;
|
||||||
|
|
@ -27,6 +27,7 @@ import com.datastax.driver.core.Session;
|
||||||
import org.apache.cassandra.distributed.Cluster;
|
import org.apache.cassandra.distributed.Cluster;
|
||||||
import org.apache.cassandra.distributed.api.IInvokableInstance;
|
import org.apache.cassandra.distributed.api.IInvokableInstance;
|
||||||
import org.apache.cassandra.distributed.api.QueryResults;
|
import org.apache.cassandra.distributed.api.QueryResults;
|
||||||
|
import org.apache.cassandra.distributed.test.TestBaseImpl;
|
||||||
import org.apache.cassandra.tools.ToolRunner;
|
import org.apache.cassandra.tools.ToolRunner;
|
||||||
import org.apache.cassandra.tools.ToolRunner.ToolResult;
|
import org.apache.cassandra.tools.ToolRunner.ToolResult;
|
||||||
|
|
||||||
|
|
@ -0,0 +1,115 @@
|
||||||
|
/*
|
||||||
|
* 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.distributed.test.fql;
|
||||||
|
|
||||||
|
import com.google.common.collect.Sets;
|
||||||
|
import org.junit.AfterClass;
|
||||||
|
import org.junit.BeforeClass;
|
||||||
|
import org.junit.Rule;
|
||||||
|
import org.junit.Test;
|
||||||
|
import org.junit.rules.TemporaryFolder;
|
||||||
|
|
||||||
|
import com.datastax.driver.core.BatchStatement;
|
||||||
|
import com.datastax.driver.core.PreparedStatement;
|
||||||
|
import com.datastax.driver.core.Session;
|
||||||
|
import org.apache.cassandra.distributed.Cluster;
|
||||||
|
import org.apache.cassandra.distributed.test.TestBaseImpl;
|
||||||
|
import org.apache.cassandra.tools.ToolRunner;
|
||||||
|
|
||||||
|
import static org.junit.Assert.assertEquals;
|
||||||
|
import static org.junit.Assert.assertTrue;
|
||||||
|
import static org.apache.cassandra.distributed.api.Feature.GOSSIP;
|
||||||
|
import static org.apache.cassandra.distributed.api.Feature.NATIVE_PROTOCOL;
|
||||||
|
import static org.apache.cassandra.distributed.api.Feature.NETWORK;
|
||||||
|
import static org.apache.cassandra.distributed.shared.AssertUtils.assertRows;
|
||||||
|
import static org.apache.cassandra.distributed.shared.AssertUtils.row;
|
||||||
|
|
||||||
|
public class FqlTombstoneHandlingTest extends TestBaseImpl
|
||||||
|
{
|
||||||
|
private static Cluster CLUSTER;
|
||||||
|
|
||||||
|
@Rule
|
||||||
|
public final TemporaryFolder temporaryFolder = new TemporaryFolder();
|
||||||
|
|
||||||
|
@BeforeClass
|
||||||
|
public static void beforeClass() throws Throwable
|
||||||
|
{
|
||||||
|
CLUSTER = init(Cluster.build(1).withConfig(updater -> updater.with(NETWORK, GOSSIP, NATIVE_PROTOCOL)).start());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
public void testNullCellBindingInBatch()
|
||||||
|
{
|
||||||
|
String tableName = "null_as_tombstone_in_batch";
|
||||||
|
CLUSTER.schemaChange(withKeyspace("CREATE TABLE %s." + tableName + " (k int, c int, s set<int>, primary key (k, c))"));
|
||||||
|
CLUSTER.get(1).nodetool("enablefullquerylog", "--path", temporaryFolder.getRoot().getAbsolutePath());
|
||||||
|
String insertTemplate = withKeyspace("INSERT INTO %s." + tableName + " (k, c, s) VALUES ( ?, ?, ?) USING TIMESTAMP 2");
|
||||||
|
String select = withKeyspace("SELECT * FROM %s." + tableName + " WHERE k = 0 AND c = 0");
|
||||||
|
|
||||||
|
com.datastax.driver.core.Cluster.Builder builder1 =com.datastax.driver.core.Cluster.builder().addContactPoint("127.0.0.1");
|
||||||
|
|
||||||
|
// Use the driver to write this initial row, since otherwise we won't hit the dispatcher
|
||||||
|
try (com.datastax.driver.core.Cluster cluster1 = builder1.build(); Session session1 = cluster1.connect())
|
||||||
|
{
|
||||||
|
BatchStatement batch = new BatchStatement(BatchStatement.Type.UNLOGGED);
|
||||||
|
PreparedStatement preparedWrite = session1.prepare(insertTemplate);
|
||||||
|
batch.add(preparedWrite.bind(0, 0, null));
|
||||||
|
session1.execute(batch);
|
||||||
|
}
|
||||||
|
|
||||||
|
CLUSTER.get(1).nodetool("disablefullquerylog");
|
||||||
|
|
||||||
|
// The dump should contain a null entry for our tombstone
|
||||||
|
ToolRunner.ToolResult runner = ToolRunner.invokeClass("org.apache.cassandra.fqltool.FullQueryLogTool",
|
||||||
|
"dump",
|
||||||
|
"--",
|
||||||
|
temporaryFolder.getRoot().getAbsolutePath());
|
||||||
|
assertTrue(runner.getStdout().contains(insertTemplate));
|
||||||
|
assertEquals(0, runner.getExitCode());
|
||||||
|
|
||||||
|
Object[][] preReplayResult = CLUSTER.get(1).executeInternal(select);
|
||||||
|
assertRows(preReplayResult, row(0, 0, null));
|
||||||
|
|
||||||
|
// Make sure the row no longer exists after truncate...
|
||||||
|
CLUSTER.get(1).executeInternal(withKeyspace("TRUNCATE %s." + tableName));
|
||||||
|
assertRows(CLUSTER.get(1).executeInternal(select));
|
||||||
|
|
||||||
|
// ...insert a new row with an actual value for the set at an earlier timestamp...
|
||||||
|
CLUSTER.get(1).executeInternal(withKeyspace("INSERT INTO %s." + tableName + " (k, c, s) VALUES ( ?, ?, ?) USING TIMESTAMP 1"), 0, 0, Sets.newHashSet(1));
|
||||||
|
assertRows(CLUSTER.get(1).executeInternal(select), row(0, 0, Sets.newHashSet(1)));
|
||||||
|
|
||||||
|
runner = ToolRunner.invokeClass("org.apache.cassandra.fqltool.FullQueryLogTool",
|
||||||
|
"replay",
|
||||||
|
"--keyspace", KEYSPACE,
|
||||||
|
"--target", "127.0.0.1",
|
||||||
|
"--", temporaryFolder.getRoot().getAbsolutePath());
|
||||||
|
assertEquals(0, runner.getExitCode());
|
||||||
|
|
||||||
|
// ...then ensure the replayed row deletes the one we wrote before replay.
|
||||||
|
Object[][] postReplayResult = CLUSTER.get(1).executeInternal(select);
|
||||||
|
assertRows(postReplayResult, preReplayResult);
|
||||||
|
}
|
||||||
|
|
||||||
|
@AfterClass
|
||||||
|
public static void afterClass()
|
||||||
|
{
|
||||||
|
if (CLUSTER != null)
|
||||||
|
CLUSTER.close();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
@ -94,7 +94,10 @@ public class FQLQueryReader implements ReadMarshallable
|
||||||
values.add(subValues);
|
values.add(subValues);
|
||||||
int numSubValues = in.int32();
|
int numSubValues = in.int32();
|
||||||
for (int zz = 0; zz < numSubValues; zz++)
|
for (int zz = 0; zz < numSubValues; zz++)
|
||||||
subValues.add(ByteBuffer.wrap(in.bytes()));
|
{
|
||||||
|
byte[] valueBytes = in.bytes();
|
||||||
|
subValues.add(valueBytes == null ? null : ByteBuffer.wrap(valueBytes));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
query = new FQLQuery.Batch(keyspace,
|
query = new FQLQuery.Batch(keyspace,
|
||||||
protocolVersion,
|
protocolVersion,
|
||||||
|
|
|
||||||
|
|
@ -126,7 +126,7 @@ public class Dump implements Runnable
|
||||||
break;
|
break;
|
||||||
|
|
||||||
case (FullQueryLogger.BATCH):
|
case (FullQueryLogger.BATCH):
|
||||||
dumpBatch(options, wireIn, sb);
|
dumpBatch(wireIn, sb);
|
||||||
break;
|
break;
|
||||||
|
|
||||||
default:
|
default:
|
||||||
|
|
@ -183,7 +183,7 @@ public class Dump implements Runnable
|
||||||
sb.append(System.lineSeparator());
|
sb.append(System.lineSeparator());
|
||||||
}
|
}
|
||||||
|
|
||||||
private static void dumpBatch(QueryOptions options, WireIn wireIn, StringBuilder sb)
|
private static void dumpBatch(WireIn wireIn, StringBuilder sb)
|
||||||
{
|
{
|
||||||
sb.append("Batch type: ")
|
sb.append("Batch type: ")
|
||||||
.append(wireIn.read(FullQueryLogger.BATCH_TYPE).text())
|
.append(wireIn.read(FullQueryLogger.BATCH_TYPE).text())
|
||||||
|
|
@ -203,7 +203,10 @@ public class Dump implements Runnable
|
||||||
int numSubValues = in.int32();
|
int numSubValues = in.int32();
|
||||||
List<ByteBuffer> subValues = new ArrayList<>(numSubValues);
|
List<ByteBuffer> subValues = new ArrayList<>(numSubValues);
|
||||||
for (int j = 0; j < numSubValues; j++)
|
for (int j = 0; j < numSubValues; j++)
|
||||||
subValues.add(ByteBuffer.wrap(in.bytes()));
|
{
|
||||||
|
byte[] valueBytes = in.bytes();
|
||||||
|
subValues.add(valueBytes == null ? null : ByteBuffer.wrap(valueBytes));
|
||||||
|
}
|
||||||
|
|
||||||
sb.append("Query: ")
|
sb.append("Query: ")
|
||||||
.append(queries.get(i))
|
.append(queries.get(i))
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue