From f83bd5ac2bbc6755213a6ad0675e7e5400c79670 Mon Sep 17 00:00:00 2001 From: Marcus Eriksson Date: Fri, 24 Aug 2018 14:41:09 +0200 Subject: [PATCH] Add fqltool compare Patch by marcuse; reviewed by Jason Brown and Dinesh Joshi for CASSANDRA-14619 --- CHANGES.txt | 1 + build.xml | 67 ++- ide/idea-iml-file.xml | 2 + ide/idea/workspace.xml | 1 + .../cassandra/tools/fqltool/ResultStore.java | 142 ----- tools/bin/cassandra.in.bat | 2 +- tools/bin/cassandra.in.sh | 2 +- {bin => tools/bin}/fqltool | 2 +- {bin => tools/bin}/fqltool.bat | 2 +- .../cassandra}/fqltool/DriverResultSet.java | 15 +- .../apache/cassandra}/fqltool/FQLQuery.java | 53 +- .../cassandra}/fqltool/FQLQueryIterator.java | 2 +- .../cassandra}/fqltool/FQLQueryReader.java | 2 +- .../cassandra/fqltool}/FullQueryLogTool.java | 12 +- .../cassandra}/fqltool/QueryReplayer.java | 43 +- .../cassandra}/fqltool/ResultComparator.java | 2 +- .../cassandra}/fqltool/ResultHandler.java | 13 +- .../apache/cassandra/fqltool/ResultStore.java | 291 +++++++++++ .../cassandra/fqltool/StoredResultSet.java | 292 +++++++++++ .../cassandra/fqltool/commands/Compare.java | 120 +++++ .../cassandra}/fqltool/commands/Dump.java | 2 +- .../cassandra}/fqltool/commands/Replay.java | 10 +- .../cassandra/fqltool/FQLCompareTest.java | 131 +++++ .../cassandra}/fqltool/FQLReplayTest.java | 493 ++++++++---------- 24 files changed, 1192 insertions(+), 510 deletions(-) delete mode 100644 src/java/org/apache/cassandra/tools/fqltool/ResultStore.java rename {bin => tools/bin}/fqltool (97%) rename {bin => tools/bin}/fqltool.bat (95%) rename {src/java/org/apache/cassandra/tools => tools/fqltool/src/org/apache/cassandra}/fqltool/DriverResultSet.java (94%) rename {src/java/org/apache/cassandra/tools => tools/fqltool/src/org/apache/cassandra}/fqltool/FQLQuery.java (85%) rename {src/java/org/apache/cassandra/tools => tools/fqltool/src/org/apache/cassandra}/fqltool/FQLQueryIterator.java (98%) rename {src/java/org/apache/cassandra/tools => tools/fqltool/src/org/apache/cassandra}/fqltool/FQLQueryReader.java (99%) rename {src/java/org/apache/cassandra/tools => tools/fqltool/src/org/apache/cassandra/fqltool}/FullQueryLogTool.java (89%) rename {src/java/org/apache/cassandra/tools => tools/fqltool/src/org/apache/cassandra}/fqltool/QueryReplayer.java (86%) rename {src/java/org/apache/cassandra/tools => tools/fqltool/src/org/apache/cassandra}/fqltool/ResultComparator.java (98%) rename {src/java/org/apache/cassandra/tools => tools/fqltool/src/org/apache/cassandra}/fqltool/ResultHandler.java (93%) create mode 100644 tools/fqltool/src/org/apache/cassandra/fqltool/ResultStore.java create mode 100644 tools/fqltool/src/org/apache/cassandra/fqltool/StoredResultSet.java create mode 100644 tools/fqltool/src/org/apache/cassandra/fqltool/commands/Compare.java rename {src/java/org/apache/cassandra/tools => tools/fqltool/src/org/apache/cassandra}/fqltool/commands/Dump.java (99%) rename {src/java/org/apache/cassandra/tools => tools/fqltool/src/org/apache/cassandra}/fqltool/commands/Replay.java (94%) create mode 100644 tools/fqltool/test/unit/org/apache/cassandra/fqltool/FQLCompareTest.java rename {test/unit/org/apache/cassandra/tools => tools/fqltool/test/unit/org/apache/cassandra}/fqltool/FQLReplayTest.java (64%) diff --git a/CHANGES.txt b/CHANGES.txt index 1227337073..1ba9975738 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 4.0 + * Add fqltool compare (CASSANDRA-14619) * Add fqltool replay (CASSANDRA-14618) * Log keyspace in full query log (CASSANDRA-14656) * Transient Replication and Cheap Quorums (CASSANDRA-14404) diff --git a/build.xml b/build.xml index b53c47b0d0..1c957d33ed 100644 --- a/build.xml +++ b/build.xml @@ -233,6 +233,7 @@ + @@ -914,6 +915,50 @@ + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + + @@ -950,7 +995,7 @@ @@ -962,6 +1007,16 @@ + + + + + + + + + + - + + + @@ -1392,7 +1449,7 @@ - + @@ -1529,7 +1586,7 @@ @@ -1848,7 +1905,9 @@ + + diff --git a/ide/idea-iml-file.xml b/ide/idea-iml-file.xml index 2512b1d7c1..b83abfa9f1 100644 --- a/ide/idea-iml-file.xml +++ b/ide/idea-iml-file.xml @@ -29,6 +29,8 @@ + + diff --git a/ide/idea/workspace.xml b/ide/idea/workspace.xml index 876ed58d7a..a2dea2a44e 100644 --- a/ide/idea/workspace.xml +++ b/ide/idea/workspace.xml @@ -297,6 +297,7 @@ + diff --git a/src/java/org/apache/cassandra/tools/fqltool/ResultStore.java b/src/java/org/apache/cassandra/tools/fqltool/ResultStore.java deleted file mode 100644 index 6d6aaac94d..0000000000 --- a/src/java/org/apache/cassandra/tools/fqltool/ResultStore.java +++ /dev/null @@ -1,142 +0,0 @@ -/* - * 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.tools.fqltool; - - -import java.io.File; -import java.nio.ByteBuffer; -import java.util.HashSet; -import java.util.List; -import java.util.Set; -import java.util.stream.Collectors; - -import net.openhft.chronicle.bytes.BytesStore; -import net.openhft.chronicle.core.io.Closeable; -import net.openhft.chronicle.queue.ChronicleQueue; -import net.openhft.chronicle.queue.ChronicleQueueBuilder; -import net.openhft.chronicle.queue.ExcerptAppender; -import net.openhft.chronicle.wire.ValueOut; -import org.apache.cassandra.utils.binlog.BinLog; - -/** - * see FQLReplayTest#readResultFile for how to read files produced by this class - */ -public class ResultStore -{ - private final List queues; - private final List appenders; - private final ChronicleQueue queryStoreQueue; - private final ExcerptAppender queryStoreAppender; - private final Set finishedHosts = new HashSet<>(); - - public ResultStore(List resultPaths, File queryFilePath) - { - queues = resultPaths.stream().map(path -> ChronicleQueueBuilder.single(path).build()).collect(Collectors.toList()); - appenders = queues.stream().map(ChronicleQueue::acquireAppender).collect(Collectors.toList()); - queryStoreQueue = queryFilePath != null ? ChronicleQueueBuilder.single(queryFilePath).build() : null; - queryStoreAppender = queryStoreQueue != null ? queryStoreQueue.acquireAppender() : null; - } - - /** - * Store the column definitions in cds - * - * the ColumnDefinitions at position x will get stored by the appender at position x - * - * Calling this method indicates that we are starting a new result set from a query, it must be called before - * calling storeRows. - * - */ - public void storeColumnDefinitions(FQLQuery query, List cds) - { - finishedHosts.clear(); - if (queryStoreAppender != null) - { - BinLog.ReleaseableWriteMarshallable writeMarshallableQuery = query.toMarshallable(); - queryStoreAppender.writeDocument(writeMarshallableQuery); - writeMarshallableQuery.release(); - } - for (int i = 0; i < cds.size(); i++) - { - ResultHandler.ComparableColumnDefinitions cd = cds.get(i); - appenders.get(i).writeDocument(wire -> - { - if (!cd.wasFailed()) - { - wire.write("type").text("column_definitions"); - wire.write("column_count").int32(cd.size()); - for (ResultHandler.ComparableDefinition d : cd.asList()) - { - ValueOut vo = wire.write("column_definition"); - vo.text(d.getName()); - vo.text(d.getType()); - } - } - else - { - wire.write("type").text("query_failed"); - } - }); - } - } - - /** - * Store rows - * - * the row at position x will get stored by appender at position x - * - * Before calling this for a new result set, storeColumnDefinitions must be called. - */ - public void storeRows(List rows) - { - for (int i = 0; i < rows.size(); i++) - { - ResultHandler.ComparableRow row = rows.get(i); - if (row == null && !finishedHosts.contains(i)) - { - appenders.get(i).writeDocument(wire -> wire.write("type").text("end_resultset")); - finishedHosts.add(i); - } - else if (row != null) - { - appenders.get(i).writeDocument(wire -> - { - { - wire.write("type").text("row"); - wire.write("row_column_count").int32(row.getColumnDefinitions().size()); - for (int jj = 0; jj < row.getColumnDefinitions().size(); jj++) - { - ByteBuffer bb = row.getBytesUnsafe(jj); - if (bb != null) - wire.write("column").bytes(BytesStore.wrap(bb)); - else - wire.write("column").bytes("NULL".getBytes()); - } - } - }); - } - } - } - - public void close() - { - queues.forEach(Closeable::close); - if (queryStoreQueue != null) - queryStoreQueue.close(); - } -} diff --git a/tools/bin/cassandra.in.bat b/tools/bin/cassandra.in.bat index 88049216e5..4d395e837e 100644 --- a/tools/bin/cassandra.in.bat +++ b/tools/bin/cassandra.in.bat @@ -39,7 +39,7 @@ goto :eof :okClasspath REM Include the build\classes\main directory so it works in development -set CASSANDRA_CLASSPATH=%CLASSPATH%;%CASSANDRA_CONF%;"%CASSANDRA_HOME%\build\classes\main";"%CASSANDRA_HOME%\build\classes\stress" +set CASSANDRA_CLASSPATH=%CLASSPATH%;%CASSANDRA_CONF%;"%CASSANDRA_HOME%\build\classes\main";"%CASSANDRA_HOME%\build\classes\stress";"%CASSANDRA_HOME%\build\classes\fqltool" REM Add the default storage location. Can be overridden in conf\cassandra.yaml set CASSANDRA_PARAMS=%CASSANDRA_PARAMS% "-Dcassandra.storagedir=%CASSANDRA_HOME%\data" diff --git a/tools/bin/cassandra.in.sh b/tools/bin/cassandra.in.sh index 869eb224fa..bf1ecc414c 100644 --- a/tools/bin/cassandra.in.sh +++ b/tools/bin/cassandra.in.sh @@ -32,7 +32,7 @@ CLASSPATH="$CASSANDRA_CONF" if [ -d $CASSANDRA_HOME/build ] ; then #cassandra_bin="$CASSANDRA_HOME/build/classes/main" cassandra_bin=`ls -1 $CASSANDRA_HOME/build/apache-cassandra*.jar` - cassandra_bin="$cassandra_bin:$CASSANDRA_HOME/build/classes/stress" + cassandra_bin="$cassandra_bin:$CASSANDRA_HOME/build/classes/stress:$CASSANDRA_HOME/build/classes/fqltool" CLASSPATH="$CLASSPATH:$cassandra_bin" fi diff --git a/bin/fqltool b/tools/bin/fqltool similarity index 97% rename from bin/fqltool rename to tools/bin/fqltool index 15a0b200ad..a34128ed83 100755 --- a/bin/fqltool +++ b/tools/bin/fqltool @@ -71,6 +71,6 @@ fi "$JAVA" $JAVA_AGENT -ea -da:net.openhft... -cp "$CLASSPATH" $JVM_OPTS -Xmx$MAX_HEAP_SIZE \ -Dlog4j.configurationFile=log4j2-tools.xml \ $JVM_ARGS \ - org.apache.cassandra.tools.FullQueryLogTool $ARGS + org.apache.cassandra.fqltool.FullQueryLogTool $ARGS # vi:ai sw=4 ts=4 tw=0 et diff --git a/bin/fqltool.bat b/tools/bin/fqltool.bat similarity index 95% rename from bin/fqltool.bat rename to tools/bin/fqltool.bat index f3103d8d6d..acb6d1c74f 100644 --- a/bin/fqltool.bat +++ b/tools/bin/fqltool.bat @@ -24,7 +24,7 @@ if NOT DEFINED JAVA_HOME goto :err set CASSANDRA_PARAMS=%CASSANDRA_PARAMS% -Dcassandra.logdir="%CASSANDRA_HOME%\logs" -"%JAVA_HOME%\bin\java" -cp %CASSANDRA_CLASSPATH% %CASSANDRA_PARAMS% -Dlog4j.configurationFile=log4j2-tools.xml org.apache.cassandra.tools.FullQueryLogTool %* +"%JAVA_HOME%\bin\java" -cp %CASSANDRA_CLASSPATH% %CASSANDRA_PARAMS% -Dlog4j.configurationFile=log4j2-tools.xml org.apache.cassandra.fqltool.FullQueryLogTool %* goto finally :err diff --git a/src/java/org/apache/cassandra/tools/fqltool/DriverResultSet.java b/tools/fqltool/src/org/apache/cassandra/fqltool/DriverResultSet.java similarity index 94% rename from src/java/org/apache/cassandra/tools/fqltool/DriverResultSet.java rename to tools/fqltool/src/org/apache/cassandra/fqltool/DriverResultSet.java index 6c4ee453d4..ccff370df8 100644 --- a/src/java/org/apache/cassandra/tools/fqltool/DriverResultSet.java +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/DriverResultSet.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.cassandra.tools.fqltool; +package org.apache.cassandra.fqltool; import java.nio.ByteBuffer; import java.util.Collections; @@ -60,7 +60,7 @@ public class DriverResultSet implements ResultHandler.ComparableResultSet public ResultHandler.ComparableColumnDefinitions getColumnDefinitions() { if (wasFailed()) - return new DriverColumnDefinitions(null, true); + return new DriverColumnDefinitions(null, true, failureException); return new DriverColumnDefinitions(resultSet.getColumnDefinitions()); } @@ -155,16 +155,18 @@ public class DriverResultSet implements ResultHandler.ComparableResultSet { private final ColumnDefinitions columnDefinitions; private final boolean failed; + private final Throwable failureException; public DriverColumnDefinitions(ColumnDefinitions columnDefinitions) { - this(columnDefinitions, false); + this(columnDefinitions, false, null); } - private DriverColumnDefinitions(ColumnDefinitions columnDefinitions, boolean failed) + private DriverColumnDefinitions(ColumnDefinitions columnDefinitions, boolean failed, Throwable failureException) { this.columnDefinitions = columnDefinitions; this.failed = failed; + this.failureException = failureException; } public List asList() @@ -179,6 +181,11 @@ public class DriverResultSet implements ResultHandler.ComparableResultSet return failed; } + public Throwable getFailureException() + { + return failureException; + } + public int size() { return columnDefinitions.size(); diff --git a/src/java/org/apache/cassandra/tools/fqltool/FQLQuery.java b/tools/fqltool/src/org/apache/cassandra/fqltool/FQLQuery.java similarity index 85% rename from src/java/org/apache/cassandra/tools/fqltool/FQLQuery.java rename to tools/fqltool/src/org/apache/cassandra/fqltool/FQLQuery.java index 6c0a6b9ad0..2862e0f084 100644 --- a/src/java/org/apache/cassandra/tools/fqltool/FQLQuery.java +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/FQLQuery.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.cassandra.tools.fqltool; +package org.apache.cassandra.fqltool; import java.nio.ByteBuffer; import java.util.ArrayList; @@ -43,18 +43,14 @@ public abstract class FQLQuery implements Comparable public final long queryStartTime; public final QueryOptions queryOptions; public final int protocolVersion; - public final String keyspace; - public final long generatedTimestamp; - private final int generatedNowInSeconds; + public final QueryState queryState; public FQLQuery(String keyspace, int protocolVersion, QueryOptions queryOptions, long queryStartTime, long generatedTimestamp, int generatedNowInSeconds) { this.queryStartTime = queryStartTime; this.queryOptions = queryOptions; this.protocolVersion = protocolVersion; - this.keyspace = keyspace; - this.generatedTimestamp = generatedTimestamp; - this.generatedNowInSeconds = generatedNowInSeconds; + this.queryState = queryState(keyspace, generatedTimestamp, generatedNowInSeconds); } public abstract Statement toStatement(); @@ -64,10 +60,14 @@ public abstract class FQLQuery implements Comparable */ public abstract BinLog.ReleaseableWriteMarshallable toMarshallable(); - public QueryState queryState() + public String keyspace() + { + return queryState.getClientState().getRawKeyspace(); + } + + private QueryState queryState(String keyspace, long generatedTimestamp, int generatedNowInSeconds) { ClientState clientState = keyspace != null ? ClientState.forInternalCalls(keyspace) : ClientState.forInternalCalls(); - return new QueryState(clientState, generatedTimestamp, generatedNowInSeconds); } @@ -78,27 +78,19 @@ public abstract class FQLQuery implements Comparable FQLQuery fqlQuery = (FQLQuery) o; return queryStartTime == fqlQuery.queryStartTime && protocolVersion == fqlQuery.protocolVersion && - generatedTimestamp == fqlQuery.generatedTimestamp && - generatedNowInSeconds == fqlQuery.generatedNowInSeconds && - Objects.equals(queryOptions.getValues(), fqlQuery.queryOptions.getValues()) && - Objects.equals(keyspace, fqlQuery.keyspace); + queryState.getTimestamp() == fqlQuery.queryState.getTimestamp() && + Objects.equals(queryState.getClientState().getRawKeyspace(), fqlQuery.queryState.getClientState().getRawKeyspace()) && + Objects.equals(queryOptions.getValues(), fqlQuery.queryOptions.getValues()); } public int hashCode() { - return Objects.hash(queryStartTime, queryOptions, protocolVersion, keyspace, generatedTimestamp, generatedNowInSeconds); + return Objects.hash(queryStartTime, queryOptions, protocolVersion, queryState.getClientState().getRawKeyspace()); } public int compareTo(FQLQuery other) { - int cmp = Longs.compare(queryStartTime, other.queryStartTime); - if (cmp != 0) - return cmp; - cmp = Longs.compare(generatedTimestamp, other.generatedTimestamp); - if (cmp != 0) - return cmp; - - return Longs.compare(generatedNowInSeconds, other.generatedNowInSeconds); + return Longs.compare(queryStartTime, other.queryStartTime); } public String toString() @@ -106,9 +98,7 @@ public abstract class FQLQuery implements Comparable return "FQLQuery{" + "queryStartTime=" + queryStartTime + ", protocolVersion=" + protocolVersion + - ", keyspace='" + keyspace + '\'' + - ", generatedTimestamp=" + generatedTimestamp + - ", generatedNowInSeconds=" + generatedNowInSeconds + + ", queryState='" + queryState + '\'' + '}'; } @@ -137,14 +127,14 @@ public abstract class FQLQuery implements Comparable { SimpleStatement ss = new SimpleStatement(query, values.toArray()); ss.setConsistencyLevel(ConsistencyLevel.valueOf(queryOptions.getConsistency().name())); - ss.setDefaultTimestamp(generatedTimestamp); + ss.setDefaultTimestamp(queryOptions.getTimestamp(queryState)); return ss; } public BinLog.ReleaseableWriteMarshallable toMarshallable() { - return new FullQueryLogger.Query(query, queryOptions, queryState(), queryStartTime); + return new FullQueryLogger.Query(query, queryOptions, queryState, queryStartTime); } public int compareTo(FQLQuery other) @@ -207,13 +197,10 @@ public abstract class FQLQuery implements Comparable public Statement toStatement() { BatchStatement bs = new BatchStatement(batchType); - for (Single query : queries) - { - bs.add(new SimpleStatement(query.query, query.values.toArray())); - } + bs.add(query.toStatement()); bs.setConsistencyLevel(ConsistencyLevel.valueOf(queryOptions.getConsistency().name())); - bs.setDefaultTimestamp(generatedTimestamp); // todo: set actual server side generated time + bs.setDefaultTimestamp(queryOptions.getTimestamp(queryState)); return bs; } @@ -248,7 +235,7 @@ public abstract class FQLQuery implements Comparable queryStrings.add(q.query); values.add(q.values); } - return new FullQueryLogger.Batch(org.apache.cassandra.cql3.statements.BatchStatement.Type.valueOf(batchType.name()), queryStrings, values, queryOptions, queryState(), queryStartTime); + return new FullQueryLogger.Batch(org.apache.cassandra.cql3.statements.BatchStatement.Type.valueOf(batchType.name()), queryStrings, values, queryOptions, queryState, queryStartTime); } public String toString() diff --git a/src/java/org/apache/cassandra/tools/fqltool/FQLQueryIterator.java b/tools/fqltool/src/org/apache/cassandra/fqltool/FQLQueryIterator.java similarity index 98% rename from src/java/org/apache/cassandra/tools/fqltool/FQLQueryIterator.java rename to tools/fqltool/src/org/apache/cassandra/fqltool/FQLQueryIterator.java index 390a52e1a1..ccbb200be5 100644 --- a/src/java/org/apache/cassandra/tools/fqltool/FQLQueryIterator.java +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/FQLQueryIterator.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.cassandra.tools.fqltool; +package org.apache.cassandra.fqltool; import java.util.PriorityQueue; diff --git a/src/java/org/apache/cassandra/tools/fqltool/FQLQueryReader.java b/tools/fqltool/src/org/apache/cassandra/fqltool/FQLQueryReader.java similarity index 99% rename from src/java/org/apache/cassandra/tools/fqltool/FQLQueryReader.java rename to tools/fqltool/src/org/apache/cassandra/fqltool/FQLQueryReader.java index af77c595c3..fd5073cc56 100644 --- a/src/java/org/apache/cassandra/tools/fqltool/FQLQueryReader.java +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/FQLQueryReader.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.cassandra.tools.fqltool; +package org.apache.cassandra.fqltool; import java.nio.ByteBuffer; diff --git a/src/java/org/apache/cassandra/tools/FullQueryLogTool.java b/tools/fqltool/src/org/apache/cassandra/fqltool/FullQueryLogTool.java similarity index 89% rename from src/java/org/apache/cassandra/tools/FullQueryLogTool.java rename to tools/fqltool/src/org/apache/cassandra/fqltool/FullQueryLogTool.java index c1d4713513..97e74872ec 100644 --- a/src/java/org/apache/cassandra/tools/FullQueryLogTool.java +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/FullQueryLogTool.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.cassandra.tools; +package org.apache.cassandra.fqltool; import java.util.List; @@ -31,8 +31,10 @@ import io.airlift.airline.ParseCommandUnrecognizedException; import io.airlift.airline.ParseOptionConversionException; import io.airlift.airline.ParseOptionMissingException; import io.airlift.airline.ParseOptionMissingValueException; -import org.apache.cassandra.tools.fqltool.commands.Dump; -import org.apache.cassandra.tools.fqltool.commands.Replay; +import org.apache.cassandra.config.DatabaseDescriptor; +import org.apache.cassandra.fqltool.commands.Compare; +import org.apache.cassandra.fqltool.commands.Dump; +import org.apache.cassandra.fqltool.commands.Replay; import static com.google.common.base.Throwables.getStackTraceAsString; import static com.google.common.collect.Lists.newArrayList; @@ -41,10 +43,12 @@ public class FullQueryLogTool { public static void main(String... args) { + DatabaseDescriptor.clientInitialization(); List> commands = newArrayList( Help.class, Dump.class, - Replay.class + Replay.class, + Compare.class ); Cli.CliBuilder builder = Cli.builder("fqltool"); diff --git a/src/java/org/apache/cassandra/tools/fqltool/QueryReplayer.java b/tools/fqltool/src/org/apache/cassandra/fqltool/QueryReplayer.java similarity index 86% rename from src/java/org/apache/cassandra/tools/fqltool/QueryReplayer.java rename to tools/fqltool/src/org/apache/cassandra/fqltool/QueryReplayer.java index 0c8382f4fb..d8653e57fe 100644 --- a/src/java/org/apache/cassandra/tools/fqltool/QueryReplayer.java +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/QueryReplayer.java @@ -16,21 +16,21 @@ * limitations under the License. */ -package org.apache.cassandra.tools.fqltool; +package org.apache.cassandra.fqltool; import java.io.Closeable; import java.io.File; +import java.io.IOException; import java.io.PrintStream; import java.util.ArrayList; import java.util.Iterator; import java.util.List; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; import java.util.function.Predicate; import java.util.stream.Collectors; -import javax.annotation.Nullable; +//import javax.annotation.Nullable; import com.google.common.util.concurrent.FluentFuture; import com.google.common.util.concurrent.FutureCallback; @@ -41,7 +41,6 @@ import com.google.common.util.concurrent.MoreExecutors; import com.codahale.metrics.MetricRegistry; import com.codahale.metrics.Timer; import com.datastax.driver.core.Cluster; -import com.datastax.driver.core.ConsistencyLevel; import com.datastax.driver.core.ResultSet; import com.datastax.driver.core.Session; import com.datastax.driver.core.Statement; @@ -93,19 +92,7 @@ public class QueryReplayer implements Closeable Statement statement = query.toStatement(); for (Session session : sessions) { - try - { - if (query.keyspace != null && !query.keyspace.equals(session.getLoggedKeyspace())) - { - if (debug) - out.printf("Switching keyspace from %s to %s%n", session.getLoggedKeyspace(), query.keyspace); - session.execute("USE " + query.keyspace); - } - } - catch (Throwable t) - { - out.printf("USE %s failed: %s%n", query.keyspace, t.getMessage()); - } + maybeSetKeyspace(session, query); if (debug) { out.println("Executing query:"); @@ -119,7 +106,7 @@ public class QueryReplayer implements Closeable Futures.addCallback(resultList, new FutureCallback>() { - public void onSuccess(@Nullable List resultSets) + public void onSuccess(/*@Nullable */List resultSets) { // note that the order of resultSets is signifcant here - resultSets.get(x) should // be the result from a query against targetHosts.get(x) @@ -146,6 +133,23 @@ public class QueryReplayer implements Closeable } } + private void maybeSetKeyspace(Session session, FQLQuery query) + { + try + { + if (query.keyspace() != null && !query.keyspace().equals(session.getLoggedKeyspace())) + { + if (debug) + out.printf("Switching keyspace from %s to %s%n", session.getLoggedKeyspace(), query.keyspace()); + session.execute("USE " + query.keyspace()); + } + } + catch (Throwable t) + { + out.printf("USE %s failed: %s%n", query.keyspace(), t.getMessage()); + } + } + /** * Make sure we catch any query errors * @@ -159,9 +163,10 @@ public class QueryReplayer implements Closeable return fluentFuture.catching(Throwable.class, DriverResultSet::failed, MoreExecutors.directExecutor()); } - public void close() + public void close() throws IOException { sessions.forEach(Session::close); targetClusters.forEach(Cluster::close); + resultHandler.close(); } } diff --git a/src/java/org/apache/cassandra/tools/fqltool/ResultComparator.java b/tools/fqltool/src/org/apache/cassandra/fqltool/ResultComparator.java similarity index 98% rename from src/java/org/apache/cassandra/tools/fqltool/ResultComparator.java rename to tools/fqltool/src/org/apache/cassandra/fqltool/ResultComparator.java index 4bbaf7a756..d8d419a3f8 100644 --- a/src/java/org/apache/cassandra/tools/fqltool/ResultComparator.java +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/ResultComparator.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.cassandra.tools.fqltool; +package org.apache.cassandra.fqltool; import java.util.List; diff --git a/src/java/org/apache/cassandra/tools/fqltool/ResultHandler.java b/tools/fqltool/src/org/apache/cassandra/fqltool/ResultHandler.java similarity index 93% rename from src/java/org/apache/cassandra/tools/fqltool/ResultHandler.java rename to tools/fqltool/src/org/apache/cassandra/fqltool/ResultHandler.java index c7692319d9..8c4c01810c 100644 --- a/src/java/org/apache/cassandra/tools/fqltool/ResultHandler.java +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/ResultHandler.java @@ -16,9 +16,11 @@ * limitations under the License. */ -package org.apache.cassandra.tools.fqltool; +package org.apache.cassandra.fqltool; +import java.io.Closeable; import java.io.File; +import java.io.IOException; import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Iterator; @@ -28,7 +30,7 @@ import java.util.stream.Collectors; import com.google.common.annotations.VisibleForTesting; -public class ResultHandler +public class ResultHandler implements Closeable { private final ResultStore resultStore; private final ResultComparator resultComparator; @@ -95,6 +97,12 @@ public class ResultHandler return rows; } + public void close() throws IOException + { + if (resultStore != null) + resultStore.close(); + } + public interface ComparableResultSet extends Iterable { public ComparableColumnDefinitions getColumnDefinitions(); @@ -106,6 +114,7 @@ public class ResultHandler { public List asList(); public boolean wasFailed(); + public Throwable getFailureException(); public int size(); } diff --git a/tools/fqltool/src/org/apache/cassandra/fqltool/ResultStore.java b/tools/fqltool/src/org/apache/cassandra/fqltool/ResultStore.java new file mode 100644 index 0000000000..d12871727e --- /dev/null +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/ResultStore.java @@ -0,0 +1,291 @@ +/* + * 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.fqltool; + + +import java.io.File; +import java.nio.ByteBuffer; +import java.util.ArrayList; +import java.util.HashSet; +import java.util.List; +import java.util.Set; +import java.util.stream.Collectors; + +import net.openhft.chronicle.bytes.BytesStore; +import net.openhft.chronicle.core.io.Closeable; +import net.openhft.chronicle.core.io.IORuntimeException; +import net.openhft.chronicle.queue.ChronicleQueue; +import net.openhft.chronicle.queue.ChronicleQueueBuilder; +import net.openhft.chronicle.queue.ExcerptAppender; +import net.openhft.chronicle.wire.ReadMarshallable; +import net.openhft.chronicle.wire.ValueIn; +import net.openhft.chronicle.wire.ValueOut; +import net.openhft.chronicle.wire.WireIn; +import net.openhft.chronicle.wire.WireOut; +import net.openhft.chronicle.wire.WriteMarshallable; +import org.apache.cassandra.utils.Pair; +import org.apache.cassandra.utils.binlog.BinLog; + +/** + * note that we store each row as a separate chronicle document to be able to + * avoid reading up the entire result set in memory when comparing + * + * document formats: + * to mark the start of a new result set: + * ------------------- + * version: int16 + * type: column_definitions + * column_count: int32; + * column_definition: text, text + * column_definition: text, text + * .... + * -------------------- + * + * to mark a failed query: + * --------------------- + * version: int16 + * type: query_failed + * message: text + * --------------------- + * + * row: + * -------------------- + * version: int16 + * type: row + * row_column_count: int32 + * column: bytes + * --------------------- + * + * to mark the end of a result set: + * ------------------- + * version: int16 + * type: end_resultset + * ------------------- + * + */ +public class ResultStore +{ + private static final String VERSION = "version"; + private static final String TYPE = "type"; + // types: + private static final String ROW = "row"; + private static final String END = "end_resultset"; + private static final String FAILURE = "query_failed"; + private static final String COLUMN_DEFINITIONS = "column_definitions"; + // fields: + private static final String COLUMN_DEFINITION = "column_definition"; + private static final String COLUMN_COUNT = "column_count"; + private static final String MESSAGE = "message"; + private static final String ROW_COLUMN_COUNT = "row_column_count"; + private static final String COLUMN = "column"; + + private static final int CURRENT_VERSION = 0; + + private final List queues; + private final List appenders; + private final ChronicleQueue queryStoreQueue; + private final ExcerptAppender queryStoreAppender; + private final Set finishedHosts = new HashSet<>(); + + public ResultStore(List resultPaths, File queryFilePath) + { + queues = resultPaths.stream().map(path -> ChronicleQueueBuilder.single(path).build()).collect(Collectors.toList()); + appenders = queues.stream().map(ChronicleQueue::acquireAppender).collect(Collectors.toList()); + queryStoreQueue = queryFilePath != null ? ChronicleQueueBuilder.single(queryFilePath).build() : null; + queryStoreAppender = queryStoreQueue != null ? queryStoreQueue.acquireAppender() : null; + } + + /** + * Store the column definitions in cds + * + * the ColumnDefinitions at position x will get stored by the appender at position x + * + * Calling this method indicates that we are starting a new result set from a query, it must be called before + * calling storeRows. + * + */ + public void storeColumnDefinitions(FQLQuery query, List cds) + { + finishedHosts.clear(); + if (queryStoreAppender != null) + { + BinLog.ReleaseableWriteMarshallable writeMarshallableQuery = query.toMarshallable(); + queryStoreAppender.writeDocument(writeMarshallableQuery); + writeMarshallableQuery.release(); + } + for (int i = 0; i < cds.size(); i++) + { + ResultHandler.ComparableColumnDefinitions cd = cds.get(i); + appenders.get(i).writeDocument(new ColumnDefsWriter(cd)); + } + } + + /** + * Store rows + * + * the row at position x will get stored by appender at position x + * + * Before calling this for a new result set, storeColumnDefinitions must be called. + */ + public void storeRows(List rows) + { + for (int i = 0; i < rows.size(); i++) + { + ResultHandler.ComparableRow row = rows.get(i); + if (row == null && !finishedHosts.contains(i)) + { + appenders.get(i).writeDocument(wire -> { + wire.write(VERSION).int16(CURRENT_VERSION); + wire.write(TYPE).text(END); + }); + finishedHosts.add(i); + } + else if (row != null) + { + appenders.get(i).writeDocument(new RowWriter(row)); + } + } + } + + public void close() + { + queues.forEach(Closeable::close); + if (queryStoreQueue != null) + queryStoreQueue.close(); + } + + static class ColumnDefsWriter implements WriteMarshallable + { + private final ResultHandler.ComparableColumnDefinitions defs; + + ColumnDefsWriter(ResultHandler.ComparableColumnDefinitions defs) + { + this.defs = defs; + } + + public void writeMarshallable(WireOut wire) + { + wire.write(VERSION).int16(CURRENT_VERSION); + if (!defs.wasFailed()) + { + wire.write(TYPE).text(COLUMN_DEFINITIONS); + wire.write(COLUMN_COUNT).int32(defs.size()); + for (ResultHandler.ComparableDefinition d : defs.asList()) + { + ValueOut vo = wire.write(COLUMN_DEFINITION); + vo.text(d.getName()); + vo.text(d.getType()); + } + } + else + { + wire.write(TYPE).text(FAILURE); + wire.write(MESSAGE).text(defs.getFailureException().getMessage()); + } + } + } + + static class ColumnDefsReader implements ReadMarshallable + { + boolean wasFailed; + String failureMessage; + List> columnDefinitions = new ArrayList<>(); + + public void readMarshallable(WireIn wire) throws IORuntimeException + { + int version = wire.read(VERSION).int16(); + String type = wire.read(TYPE).text(); + if (type.equals(FAILURE)) + { + wasFailed = true; + failureMessage = wire.read(MESSAGE).text(); + } + else if (type.equals(COLUMN_DEFINITION)) + { + int columnCount = wire.read(COLUMN_COUNT).int32(); + for (int i = 0; i < columnCount; i++) + { + ValueIn vi = wire.read(COLUMN_DEFINITION); + String name = vi.text(); + String dataType = vi.text(); + columnDefinitions.add(Pair.create(name, dataType)); + } + } + } + } + + /** + * read a single row from the wire, or, marks itself finished if we read "end_resultset" + */ + static class RowReader implements ReadMarshallable + { + boolean isFinished; + List rows = new ArrayList<>(); + + public void readMarshallable(WireIn wire) throws IORuntimeException + { + int version = wire.read(VERSION).int32(); + String type = wire.read(TYPE).text(); + if (!type.equals(END)) + { + isFinished = false; + int rowColumnCount = wire.read(ROW_COLUMN_COUNT).int32(); + + for (int i = 0; i < rowColumnCount; i++) + { + byte[] b = wire.read(COLUMN).bytes(); + rows.add(ByteBuffer.wrap(b)); + } + } + else + { + isFinished = true; + } + } + } + + /** + * Writes a single row to the given wire + */ + static class RowWriter implements WriteMarshallable + { + private final ResultHandler.ComparableRow row; + + RowWriter(ResultHandler.ComparableRow row) + { + this.row = row; + } + + public void writeMarshallable(WireOut wire) + { + wire.write(VERSION).int16(CURRENT_VERSION); + wire.write(TYPE).text(ROW); + wire.write(ROW_COLUMN_COUNT).int32(row.getColumnDefinitions().size()); + for (int jj = 0; jj < row.getColumnDefinitions().size(); jj++) + { + ByteBuffer bb = row.getBytesUnsafe(jj); + if (bb != null) + wire.write(COLUMN).bytes(BytesStore.wrap(bb)); + else + wire.write(COLUMN).bytes("NULL".getBytes()); + } + } + } + +} diff --git a/tools/fqltool/src/org/apache/cassandra/fqltool/StoredResultSet.java b/tools/fqltool/src/org/apache/cassandra/fqltool/StoredResultSet.java new file mode 100644 index 0000000000..b08861dcea --- /dev/null +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/StoredResultSet.java @@ -0,0 +1,292 @@ +/* + * 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.fqltool; + + +import java.nio.ByteBuffer; +import java.util.Collections; +import java.util.Iterator; +import java.util.List; +import java.util.function.Supplier; +import java.util.stream.Collectors; + +import com.google.common.collect.AbstractIterator; + +import net.openhft.chronicle.queue.ExcerptTailer; +import org.apache.cassandra.utils.ByteBufferUtil; +import org.apache.cassandra.utils.Pair; + +/** + * represents a resultset defined by the format in ResultStore on disk + * + * todo: Currently all iterators need to be consumed fully while iterating over result sets + * if this is created from a tailer. This can probably be improved, but for all current uses it is fine. + */ +public class StoredResultSet implements ResultHandler.ComparableResultSet +{ + private final ResultHandler.ComparableColumnDefinitions defs; + public final boolean hasMoreResultSets; + private final Supplier> rowIteratorSupplier; + private final boolean wasFailed; + private final Throwable failureException; + + /** + * create a new StoredResultSet + * + * note that we use an iteratorSupplier to be able to iterate over the same in-memory rows several times *in tests* + */ + public StoredResultSet(ResultHandler.ComparableColumnDefinitions defs, + boolean hasMoreResultSets, + boolean wasFailed, + Throwable failure, + Supplier> iteratorSupplier) + { + this.defs = defs; + this.hasMoreResultSets = hasMoreResultSets; + this.wasFailed = wasFailed; + this.failureException = failure; + this.rowIteratorSupplier = iteratorSupplier; + } + + /** + * creates a ComparableResultSet based on the data in tailer + */ + public static StoredResultSet fromTailer(ExcerptTailer tailer) + { + ResultStore.ColumnDefsReader reader = new ResultStore.ColumnDefsReader(); + boolean hasMoreResultSets = tailer.readDocument(reader); + ResultHandler.ComparableColumnDefinitions defs = new StoredComparableColumnDefinitions(reader.columnDefinitions, + reader.wasFailed, + new RuntimeException(reader.failureMessage)); + + + Iterator rowIterator = new AbstractIterator() + { + protected ResultHandler.ComparableRow computeNext() + { + ResultStore.RowReader rowReader = new ResultStore.RowReader(); + tailer.readDocument(rowReader); + if (rowReader.isFinished) + return endOfData(); + return new StoredComparableRow(rowReader.rows, defs); + } + }; + + return new StoredResultSet(defs, + hasMoreResultSets, + reader.wasFailed, + new RuntimeException(reader.failureMessage), + () -> rowIterator); + } + + public static ResultHandler.ComparableResultSet failed(String failureMessage) + { + return new FailedComparableResultSet(new RuntimeException(failureMessage)); + } + + public Iterator iterator() + { + return rowIteratorSupplier.get(); + } + + public ResultHandler.ComparableColumnDefinitions getColumnDefinitions() + { + return defs; + } + + public boolean wasFailed() + { + return wasFailed; + } + + public Throwable getFailureException() + { + return failureException; + } + + static class StoredComparableRow implements ResultHandler.ComparableRow + { + private final List row; + private final ResultHandler.ComparableColumnDefinitions cds; + + public StoredComparableRow(List row, ResultHandler.ComparableColumnDefinitions cds) + { + this.row = row; + this.cds = cds; + } + + public ByteBuffer getBytesUnsafe(int i) + { + return row.get(i); + } + + public ResultHandler.ComparableColumnDefinitions getColumnDefinitions() + { + return cds; + } + + public boolean equals(Object other) + { + if (!(other instanceof StoredComparableRow)) + return false; + return row.equals(((StoredComparableRow)other).row); + } + + public String toString() + { + return row.stream().map(ByteBufferUtil::bytesToHex).collect(Collectors.joining(",")); + } + } + + static class StoredComparableColumnDefinitions implements ResultHandler.ComparableColumnDefinitions + { + private final List defs; + private final boolean wasFailed; + private final Throwable failureException; + + public StoredComparableColumnDefinitions(List> cds, boolean wasFailed, Throwable failureException) + { + defs = cds != null ? cds.stream().map(StoredComparableDefinition::new).collect(Collectors.toList()) : Collections.emptyList(); + this.wasFailed = wasFailed; + this.failureException = failureException; + } + public List asList() + { + return wasFailed() ? Collections.emptyList() : defs; + } + + public boolean wasFailed() + { + return wasFailed; + } + + public Throwable getFailureException() + { + return failureException; + } + + public int size() + { + return asList().size(); + } + + public Iterator iterator() + { + return defs.iterator(); + } + + public boolean equals(Object other) + { + if (!(other instanceof StoredComparableColumnDefinitions)) + return false; + return defs.equals(((StoredComparableColumnDefinitions)other).defs); + } + + public String toString() + { + return defs.toString(); + } + } + + private static class StoredComparableDefinition implements ResultHandler.ComparableDefinition + { + private final Pair p; + + public StoredComparableDefinition(Pair p) + { + this.p = p; + } + public String getType() + { + return p.right; + } + + public String getName() + { + return p.left; + } + + public boolean equals(Object other) + { + if (!(other instanceof StoredComparableDefinition)) + return false; + return p.equals(((StoredComparableDefinition)other).p); + } + + public String toString() + { + return getName() + ':' + getType(); + } + } + + private static class FailedComparableResultSet implements ResultHandler.ComparableResultSet + { + private final Throwable exception; + + public FailedComparableResultSet(Throwable exception) + { + this.exception = exception; + } + public ResultHandler.ComparableColumnDefinitions getColumnDefinitions() + { + return new ResultHandler.ComparableColumnDefinitions() + { + public List asList() + { + return Collections.emptyList(); + } + + public boolean wasFailed() + { + return true; + } + + public Throwable getFailureException() + { + return exception; + } + + public int size() + { + return 0; + } + + public Iterator iterator() + { + return asList().iterator(); + } + }; + } + + public boolean wasFailed() + { + return true; + } + + public Throwable getFailureException() + { + return new RuntimeException(); + } + + public Iterator iterator() + { + return Collections.emptyListIterator(); + } + } +} diff --git a/tools/fqltool/src/org/apache/cassandra/fqltool/commands/Compare.java b/tools/fqltool/src/org/apache/cassandra/fqltool/commands/Compare.java new file mode 100644 index 0000000000..2375296420 --- /dev/null +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/commands/Compare.java @@ -0,0 +1,120 @@ +/* + * 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.fqltool.commands; + + +import java.util.ArrayList; +import java.util.Iterator; +import java.util.List; +import java.util.stream.Collectors; + +import com.google.common.annotations.VisibleForTesting; +import com.google.common.collect.AbstractIterator; + +import io.airlift.airline.Arguments; +import io.airlift.airline.Command; +import io.airlift.airline.Option; +import net.openhft.chronicle.core.io.Closeable; +import net.openhft.chronicle.queue.ChronicleQueue; +import net.openhft.chronicle.queue.ChronicleQueueBuilder; +import net.openhft.chronicle.queue.ExcerptTailer; +import org.apache.cassandra.fqltool.FQLQueryIterator; +import org.apache.cassandra.fqltool.ResultHandler; +import org.apache.cassandra.fqltool.StoredResultSet; + +/** + */ +@Command(name = "compare", description = "Compare result files generated by fqltool replay") +public class Compare implements Runnable +{ + @Arguments(usage = " [...]", + description = "Directories containing result files to compare.", + required = true) + private List arguments = new ArrayList<>(); + + @Option(title = "queries", + name = { "--queries"}, + description = "Directory to read the queries from. It is produced by the fqltool replay --store-queries option. ", + required = true) + private String querylog; + + @Override + public void run() + { + compare(querylog, arguments); + } + + public static void compare(String querylog, List arguments) + { + List readQueues = null; + try (ResultHandler rh = new ResultHandler(arguments, null, null); + ChronicleQueue queryQ = ChronicleQueueBuilder.single(querylog).readOnly(true).build(); + FQLQueryIterator queries = new FQLQueryIterator(queryQ.createTailer(), 1)) + { + readQueues = arguments.stream().map(s -> ChronicleQueueBuilder.single(s).readOnly(true).build()).collect(Collectors.toList()); + List> its = readQueues.stream().map(q -> new StoredResultSetIterator(q.createTailer())).collect(Collectors.toList()); + while (queries.hasNext()) + rh.handleResults(queries.next(), resultSets(its)); + } + catch (Exception e) + { + throw new RuntimeException(e); + } + finally + { + if (readQueues != null) + readQueues.forEach(Closeable::close); + } + } + + @VisibleForTesting + public static List resultSets(List> its) + { + List resultSets = new ArrayList<>(its.size()); + for (Iterator it : its) + { + if (it.hasNext()) + resultSets.add(it.next()); + else + resultSets.add(null); + } + return resultSets; + } + + @VisibleForTesting + public static class StoredResultSetIterator extends AbstractIterator + { + private final ExcerptTailer tailer; + + public StoredResultSetIterator(ExcerptTailer tailer) + { + this.tailer = tailer; + } + + protected ResultHandler.ComparableResultSet computeNext() + { + StoredResultSet srs = StoredResultSet.fromTailer(tailer); + if (srs.hasMoreResultSets) + return srs; + return endOfData(); + } + } + + +} diff --git a/src/java/org/apache/cassandra/tools/fqltool/commands/Dump.java b/tools/fqltool/src/org/apache/cassandra/fqltool/commands/Dump.java similarity index 99% rename from src/java/org/apache/cassandra/tools/fqltool/commands/Dump.java rename to tools/fqltool/src/org/apache/cassandra/fqltool/commands/Dump.java index 5c23d3e5c4..1263a11283 100644 --- a/src/java/org/apache/cassandra/tools/fqltool/commands/Dump.java +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/commands/Dump.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.cassandra.tools.fqltool.commands; +package org.apache.cassandra.fqltool.commands; import java.io.File; import java.nio.BufferUnderflowException; diff --git a/src/java/org/apache/cassandra/tools/fqltool/commands/Replay.java b/tools/fqltool/src/org/apache/cassandra/fqltool/commands/Replay.java similarity index 94% rename from src/java/org/apache/cassandra/tools/fqltool/commands/Replay.java rename to tools/fqltool/src/org/apache/cassandra/fqltool/commands/Replay.java index 043ead8d41..adea742f77 100644 --- a/src/java/org/apache/cassandra/tools/fqltool/commands/Replay.java +++ b/tools/fqltool/src/org/apache/cassandra/fqltool/commands/Replay.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.cassandra.tools.fqltool.commands; +package org.apache.cassandra.fqltool.commands; import java.io.File; @@ -34,9 +34,9 @@ import net.openhft.chronicle.core.io.Closeable; import net.openhft.chronicle.queue.ChronicleQueue; import net.openhft.chronicle.queue.ChronicleQueueBuilder; -import org.apache.cassandra.tools.fqltool.FQLQuery; -import org.apache.cassandra.tools.fqltool.FQLQueryIterator; -import org.apache.cassandra.tools.fqltool.QueryReplayer; +import org.apache.cassandra.fqltool.FQLQuery; +import org.apache.cassandra.fqltool.FQLQueryIterator; +import org.apache.cassandra.fqltool.QueryReplayer; import org.apache.cassandra.utils.AbstractIterator; import org.apache.cassandra.utils.MergeIterator; @@ -102,7 +102,7 @@ public class Replay implements Runnable List> filters = new ArrayList<>(); if (keyspace != null) - filters.add(fqlQuery -> fqlQuery.keyspace == null || fqlQuery.keyspace.equals(keyspace)); + filters.add(fqlQuery -> fqlQuery.keyspace() == null || fqlQuery.keyspace().equals(keyspace)); try { diff --git a/tools/fqltool/test/unit/org/apache/cassandra/fqltool/FQLCompareTest.java b/tools/fqltool/test/unit/org/apache/cassandra/fqltool/FQLCompareTest.java new file mode 100644 index 0000000000..7990b7ed90 --- /dev/null +++ b/tools/fqltool/test/unit/org/apache/cassandra/fqltool/FQLCompareTest.java @@ -0,0 +1,131 @@ +/* + * 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.fqltool; + +import java.io.File; +import java.io.IOException; +import java.nio.file.Files; +import java.util.ArrayList; +import java.util.Collections; +import java.util.Iterator; +import java.util.List; +import java.util.Objects; +import java.util.stream.Collectors; + +import com.google.common.collect.Lists; +import org.junit.Test; + +import net.openhft.chronicle.core.io.Closeable; +import net.openhft.chronicle.queue.ChronicleQueue; +import net.openhft.chronicle.queue.ChronicleQueueBuilder; +import org.apache.cassandra.cql3.QueryOptions; +import org.apache.cassandra.fqltool.commands.Compare; +import org.apache.cassandra.tools.Util; + + +import static org.psjava.util.AssertStatus.assertTrue; + +public class FQLCompareTest +{ + public FQLCompareTest() + { + Util.initDatabaseDescriptor(); + } + + @Test + public void endToEnd() throws IOException + { + List targetHosts = Lists.newArrayList("hosta", "hostb"); + File tmpDir = Files.createTempDirectory("testresulthandler").toFile(); + File queryDir = Files.createTempDirectory("queries").toFile(); + List resultPaths = generateResultSets(targetHosts, tmpDir, queryDir, true, false); + Compare.compare(queryDir.toString(), resultPaths.stream().map(File::toString).collect(Collectors.toList())); + } + + @Test + public void endToEndQueryFailures() throws IOException + { + List targetHosts = Lists.newArrayList("hosta", "hostb"); + File tmpDir = Files.createTempDirectory("testresulthandler").toFile(); + File queryDir = Files.createTempDirectory("queries").toFile(); + List resultPaths = generateResultSets(targetHosts, tmpDir, queryDir, true,true); + Compare.compare(queryDir.toString(), resultPaths.stream().map(File::toString).collect(Collectors.toList())); + } + + @Test + public void compareEqual() throws IOException + { + List targetHosts = Lists.newArrayList("hosta", "hostb"); + File tmpDir = Files.createTempDirectory("testresulthandler").toFile(); + File queryDir = Files.createTempDirectory("queries").toFile(); + List resultPaths = generateResultSets(targetHosts, tmpDir, queryDir, false,false); + + ResultComparator comparator = new ResultComparator(); + List readQueues = null; + try + { + readQueues = resultPaths.stream().map(s -> ChronicleQueueBuilder.single(s).readOnly(true).build()).collect(Collectors.toList()); + List> its = readQueues.stream().map(q -> new Compare.StoredResultSetIterator(q.createTailer())).collect(Collectors.toList()); + List resultSets = Compare.resultSets(its); + while(resultSets.stream().allMatch(Objects::nonNull)) + { + assertTrue(comparator.compareColumnDefinitions(targetHosts, query(), resultSets.stream().map(ResultHandler.ComparableResultSet::getColumnDefinitions).collect(Collectors.toList()))); + List> rows = resultSets.stream().map(Iterable::iterator).collect(Collectors.toList()); + + List toCompare = ResultHandler.rows(rows); + + while (toCompare.stream().allMatch(Objects::nonNull)) + { + assertTrue(comparator.compareRows(targetHosts, query(), ResultHandler.rows(rows))); + toCompare = ResultHandler.rows(rows); + } + resultSets = Compare.resultSets(its); + } + } + finally + { + if (readQueues != null) + readQueues.forEach(Closeable::close); + } + } + + private List generateResultSets(List targetHosts, File resultDir, File queryDir, boolean random, boolean includeFailures) throws IOException + { + List resultPaths = new ArrayList<>(); + targetHosts.forEach(host -> { File f = new File(resultDir, host); f.mkdir(); resultPaths.add(f);}); + + try (ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir)) + { + for (int i = 0; i < 100; i++) + { + ResultHandler.ComparableResultSet resultSet1 = includeFailures && (i % 10 == 0) + ? StoredResultSet.failed("test failure!") + : FQLReplayTest.createResultSet(10, 10, random); + ResultHandler.ComparableResultSet resultSet2 = FQLReplayTest.createResultSet(10, 10, random); + rh.handleResults(query(), Lists.newArrayList(resultSet1, resultSet2)); + } + } + return resultPaths; + } + + private FQLQuery.Single query() + { + return new FQLQuery.Single("abc", QueryOptions.DEFAULT.getProtocolVersion().asInt(), QueryOptions.DEFAULT, 12345, 5555, 6666, "select * from xyz", Collections.emptyList()); + } +} diff --git a/test/unit/org/apache/cassandra/tools/fqltool/FQLReplayTest.java b/tools/fqltool/test/unit/org/apache/cassandra/fqltool/FQLReplayTest.java similarity index 64% rename from test/unit/org/apache/cassandra/tools/fqltool/FQLReplayTest.java rename to tools/fqltool/test/unit/org/apache/cassandra/fqltool/FQLReplayTest.java index a662699f7f..61c8aa0e92 100644 --- a/test/unit/org/apache/cassandra/tools/fqltool/FQLReplayTest.java +++ b/tools/fqltool/test/unit/org/apache/cassandra/fqltool/FQLReplayTest.java @@ -16,7 +16,7 @@ * limitations under the License. */ -package org.apache.cassandra.tools.fqltool; +package org.apache.cassandra.fqltool; import java.io.File; import java.io.IOException; @@ -28,31 +28,31 @@ import java.util.Iterator; import java.util.List; import java.util.Objects; import java.util.Random; -import java.util.concurrent.atomic.AtomicBoolean; import java.util.stream.Collectors; -import com.google.common.collect.AbstractIterator; -import com.google.common.collect.ImmutableList; -import com.google.common.collect.Iterables; import com.google.common.collect.Lists; import org.junit.Test; +import com.datastax.driver.core.CodecRegistry; +import com.datastax.driver.core.SimpleStatement; +import com.datastax.driver.core.Statement; import net.openhft.chronicle.queue.ChronicleQueue; import net.openhft.chronicle.queue.ChronicleQueueBuilder; import net.openhft.chronicle.queue.ExcerptAppender; import net.openhft.chronicle.queue.ExcerptTailer; -import net.openhft.chronicle.wire.ValueIn; import org.apache.cassandra.audit.FullQueryLogger; import org.apache.cassandra.cql3.QueryOptions; import org.apache.cassandra.cql3.statements.BatchStatement; +import org.apache.cassandra.fqltool.commands.Compare; +import org.apache.cassandra.fqltool.commands.Replay; import org.apache.cassandra.service.ClientState; import org.apache.cassandra.service.QueryState; import org.apache.cassandra.tools.Util; -import org.apache.cassandra.tools.fqltool.commands.Replay; import org.apache.cassandra.utils.ByteBufferUtil; import org.apache.cassandra.utils.MergeIterator; import org.apache.cassandra.utils.Pair; +import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; @@ -84,6 +84,27 @@ public class FQLReplayTest } assertEquals(100, queryCount); } + + @Test + public void testQueryIterator() throws IOException + { + File f = generateQueries(100, false); + int queryCount = 0; + try (ChronicleQueue queue = ChronicleQueueBuilder.single(f).build(); + FQLQueryIterator iter = new FQLQueryIterator(queue.createTailer(), 1)) + { + long last = -1; + while (iter.hasNext()) + { + FQLQuery q = iter.next(); + assertTrue(q.queryStartTime >= last); + last = q.queryStartTime; + queryCount++; + } + } + assertEquals(100, queryCount); + } + @Test public void testMergingIterator() throws IOException { @@ -125,11 +146,11 @@ public class FQLReplayTest assertNotNull(reader.getQuery()); if (reader.getQuery() instanceof FQLQuery.Single) { - assertTrue(reader.getQuery().keyspace == null || reader.getQuery().keyspace.equals("querykeyspace")); + assertTrue(reader.getQuery().keyspace() == null || reader.getQuery().keyspace().equals("querykeyspace")); } else { - assertEquals("someks", reader.getQuery().keyspace); + assertEquals("someks", reader.getQuery().keyspace()); } queryCount++; } @@ -145,9 +166,9 @@ public class FQLReplayTest ResultHandler.ComparableResultSet res = createResultSet(10, 10, true); ResultStore rs = new ResultStore(Collections.singletonList(tmpDir), queryDir); + FQLQuery query = new FQLQuery.Single("abc", QueryOptions.DEFAULT.getProtocolVersion().asInt(), QueryOptions.DEFAULT, 12345, 11111, 22, "select * from abc", Collections.emptyList()); try { - FQLQuery query = new FQLQuery.Single("abc", 3, QueryOptions.DEFAULT, 12345, 11111, 22, "select * from abc", Collections.emptyList()); rs.storeColumnDefinitions(query, Collections.singletonList(res.getColumnDefinitions())); Iterator it = res.iterator(); while (it.hasNext()) @@ -163,12 +184,41 @@ public class FQLReplayTest rs.close(); } - List> resultSets = readResultFile(tmpDir, queryDir); - assertEquals(1, resultSets.size()); - assertEquals(res, resultSets.get(0).right); + compareResults(Collections.singletonList(Pair.create(query, res)), + readResultFile(tmpDir, queryDir)); } + private static void compareResults(List> expected, + List> other) + { + ResultComparator comparator = new ResultComparator(); + assertEquals(expected.size(), other.size()); + for (int i = 0; i < expected.size(); i++) + { + assertEquals(expected.get(i).left, other.get(i).left); + ResultHandler.ComparableResultSet expectedResultSet = expected.get(i).right; + ResultHandler.ComparableResultSet otherResultSet = other.get(i).right; + List hosts = Lists.newArrayList("a", "b"); + comparator.compareColumnDefinitions(hosts, + expected.get(i).left, + Lists.newArrayList(expectedResultSet.getColumnDefinitions(), + otherResultSet.getColumnDefinitions())); + Iterator expectedRowIter = expectedResultSet.iterator(); + Iterator otherRowIter = otherResultSet.iterator(); + + + while(expectedRowIter.hasNext() && otherRowIter.hasNext()) + { + ResultHandler.ComparableRow expectedRow = expectedRowIter.next(); + ResultHandler.ComparableRow otherRow = otherRowIter.next(); + assertTrue(comparator.compareRows(hosts, expected.get(i).left, Lists.newArrayList(expectedRow, otherRow))); + } + assertFalse(expectedRowIter.hasNext()); + assertFalse(otherRowIter.hasNext()); + } + } + @Test public void testCompareColumnDefinitions() { @@ -270,18 +320,22 @@ public class FQLReplayTest File queryDir = Files.createTempDirectory("queries").toFile(); List resultPaths = new ArrayList<>(); targetHosts.forEach(host -> { File f = new File(tmpDir, host); f.mkdir(); resultPaths.add(f);}); - ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir); + ResultHandler.ComparableResultSet res = createResultSet(10, 10, false); ResultHandler.ComparableResultSet res2 = createResultSet(10, 10, false); ResultHandler.ComparableResultSet res3 = createResultSet(10, 10, false); List toCompare = Lists.newArrayList(res, res2, res3); - rh.handleResults(new FQLQuery.Single("abcabc", 3, QueryOptions.DEFAULT, 1111, 2222, 3333, "select * from xyz", Collections.emptyList()), toCompare); + FQLQuery query = new FQLQuery.Single("abcabc", QueryOptions.DEFAULT.getProtocolVersion().asInt(), QueryOptions.DEFAULT, 1111, 2222, 3333, "select * from xyz", Collections.emptyList()); + try (ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir)) + { + rh.handleResults(query, toCompare); + } List> results1 = readResultFile(resultPaths.get(0), queryDir); List> results2 = readResultFile(resultPaths.get(1), queryDir); List> results3 = readResultFile(resultPaths.get(2), queryDir); - assertEquals(results1, results2); - assertEquals(results1, results3); - assertEquals(Iterables.getOnlyElement(results3).right, res); + compareResults(results1, results2); + compareResults(results1, results3); + compareResults(results3, Collections.singletonList(Pair.create(query, res))); } @Test @@ -292,17 +346,21 @@ public class FQLReplayTest File queryDir = Files.createTempDirectory("queries").toFile(); List resultPaths = new ArrayList<>(); targetHosts.forEach(host -> { File f = new File(tmpDir, host); f.mkdir(); resultPaths.add(f);}); - ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir); + ResultHandler.ComparableResultSet res = createResultSet(10, 10, false); ResultHandler.ComparableResultSet res2 = createResultSet(10, 5, false); ResultHandler.ComparableResultSet res3 = createResultSet(10, 10, false); List toCompare = Lists.newArrayList(res, res2, res3); - rh.handleResults(new FQLQuery.Single("aaa", 3, QueryOptions.DEFAULT, 123123, 11111, 22222, "select * from abcabc", Collections.emptyList()), toCompare); + FQLQuery query = new FQLQuery.Single("aaa", QueryOptions.DEFAULT.getProtocolVersion().asInt(), QueryOptions.DEFAULT, 123123, 11111, 22222, "select * from abcabc", Collections.emptyList()); + try (ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir)) + { + rh.handleResults(query, toCompare); + } List> results1 = readResultFile(resultPaths.get(0), queryDir); List> results2 = readResultFile(resultPaths.get(1), queryDir); List> results3 = readResultFile(resultPaths.get(2), queryDir); - assertEquals(results1, results3); - assertEquals(results2.get(0).right, res2); + compareResults(results1, results3); + compareResults(results2, Collections.singletonList(Pair.create(query, res2))); } @Test @@ -313,7 +371,6 @@ public class FQLReplayTest File queryDir = Files.createTempDirectory("queries").toFile(); List resultPaths = new ArrayList<>(); targetHosts.forEach(host -> { File f = new File(tmpDir, host); f.mkdir(); resultPaths.add(f);}); - ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir); List>> resultSets = new ArrayList<>(); Random random = new Random(); for (int i = 0; i < 10; i++) @@ -324,18 +381,32 @@ public class FQLReplayTest { results.add(createResultSet(5, 1 + random.nextInt(10), true)); } - FQLQuery q = new FQLQuery.Single("abc"+i, + FQLQuery q = i % 2 == 0 + ? new FQLQuery.Single("abc"+i, 3, QueryOptions.forInternalCalls(values), i * 1000, 12345, 54321, "select * from xyz where id = "+i, - values); + values) + : new FQLQuery.Batch("abc"+i, + 3, + QueryOptions.forInternalCalls(values), + i * 1000, + i * 54321, + i * 12345, + com.datastax.driver.core.BatchStatement.Type.UNLOGGED, + Lists.newArrayList("select * from aaaa"), + Collections.singletonList(values)); + resultSets.add(Pair.create(q, results)); } - for (int i = 0; i < resultSets.size(); i++) - rh.handleResults(resultSets.get(i).left, resultSets.get(i).right); + try (ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir)) + { + for (int i = 0; i < resultSets.size(); i++) + rh.handleResults(resultSets.get(i).left, resultSets.get(i).right); + } for (int i = 0; i < targetHosts.size(); i++) compareWithFile(resultPaths, queryDir, resultSets, i); @@ -349,7 +420,7 @@ public class FQLReplayTest File queryDir = Files.createTempDirectory("queries").toFile(); List resultPaths = new ArrayList<>(); targetHosts.forEach(host -> { File f = new File(tmpDir, host); f.mkdir(); resultPaths.add(f);}); - ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir); + List>> resultSets = new ArrayList<>(); Random random = new Random(); for (int i = 0; i < 10; i++) @@ -360,8 +431,8 @@ public class FQLReplayTest { results.add(createResultSet(5, 1 + random.nextInt(10), true)); } - results.set(0, FakeResultSet.failed(new RuntimeException("testing abc"))); - results.set(3, FakeResultSet.failed(new RuntimeException("testing abc"))); + results.set(0, StoredResultSet.failed("testing abc")); + results.set(3, StoredResultSet.failed("testing abc")); FQLQuery q = new FQLQuery.Single("abc"+i, 3, QueryOptions.forInternalCalls(values), @@ -372,14 +443,15 @@ public class FQLReplayTest values); resultSets.add(Pair.create(q, results)); } - for (int i = 0; i < resultSets.size(); i++) - rh.handleResults(resultSets.get(i).left, resultSets.get(i).right); - + try (ResultHandler rh = new ResultHandler(targetHosts, resultPaths, queryDir)) + { + for (int i = 0; i < resultSets.size(); i++) + rh.handleResults(resultSets.get(i).left, resultSets.get(i).right); + } for (int i = 0; i < targetHosts.size(); i++) compareWithFile(resultPaths, queryDir, resultSets, i); } - @Test public void testCompare() { @@ -414,6 +486,75 @@ public class FQLReplayTest assertTrue(q9.compareTo(q8) > 0); } + @Test + public void testFQLQuerySingleToStatement() + { + List values = new ArrayList<>(); + for (int i = 0; i < 10; i++) + values.add(ByteBufferUtil.bytes(i)); + FQLQuery.Single single = new FQLQuery.Single("xyz", + QueryOptions.DEFAULT.getProtocolVersion().asInt(), + QueryOptions.forInternalCalls(values), + 1234, + 12345, + 54321, + "select * from aaa", + values); + Statement stmt = single.toStatement(); + assertEquals(stmt.getDefaultTimestamp(), 12345); + assertTrue(stmt instanceof SimpleStatement); + SimpleStatement simpleStmt = (SimpleStatement)stmt; + assertEquals("select * from aaa",simpleStmt.getQueryString(CodecRegistry.DEFAULT_INSTANCE)); + assertArrayEquals(values.toArray(), simpleStmt.getValues(com.datastax.driver.core.ProtocolVersion.fromInt(QueryOptions.DEFAULT.getProtocolVersion().asInt()), CodecRegistry.DEFAULT_INSTANCE)); + } + + + @Test + public void testFQLQueryBatchToStatement() + { + List> values = new ArrayList<>(); + List queries = new ArrayList<>(); + for (int bqCount = 0; bqCount < 10; bqCount++) + { + queries.add("select * from asdf where x = ? and y = " + bqCount); + List queryValues = new ArrayList<>(); + for (int i = 0; i < 10; i++) + queryValues.add(ByteBufferUtil.bytes(i + ":" + bqCount)); + values.add(queryValues); + } + + FQLQuery.Batch batch = new FQLQuery.Batch("xyz", + QueryOptions.DEFAULT.getProtocolVersion().asInt(), + QueryOptions.DEFAULT, + 1234, + 12345, + 54321, + com.datastax.driver.core.BatchStatement.Type.UNLOGGED, + queries, + values); + Statement stmt = batch.toStatement(); + assertEquals(stmt.getDefaultTimestamp(), 12345); + assertTrue(stmt instanceof com.datastax.driver.core.BatchStatement); + com.datastax.driver.core.BatchStatement batchStmt = (com.datastax.driver.core.BatchStatement)stmt; + List statements = Lists.newArrayList(batchStmt.getStatements()); + List fromFQLQueries = batch.queries.stream().map(FQLQuery.Single::toStatement).collect(Collectors.toList()); + assertEquals(statements.size(), fromFQLQueries.size()); + assertEquals(12345, batchStmt.getDefaultTimestamp()); + for (int i = 0; i < statements.size(); i++) + compareStatements(statements.get(i), fromFQLQueries.get(i)); + } + + private void compareStatements(Statement statement1, Statement statement2) + { + assertTrue(statement1 instanceof SimpleStatement && statement2 instanceof SimpleStatement); + SimpleStatement simpleStmt1 = (SimpleStatement)statement1; + SimpleStatement simpleStmt2 = (SimpleStatement)statement2; + assertEquals(simpleStmt1.getQueryString(CodecRegistry.DEFAULT_INSTANCE), simpleStmt2.getQueryString(CodecRegistry.DEFAULT_INSTANCE)); + assertArrayEquals(simpleStmt1.getValues(com.datastax.driver.core.ProtocolVersion.fromInt(QueryOptions.DEFAULT.getProtocolVersion().asInt()), CodecRegistry.DEFAULT_INSTANCE), + simpleStmt2.getValues(com.datastax.driver.core.ProtocolVersion.fromInt(QueryOptions.DEFAULT.getProtocolVersion().asInt()), CodecRegistry.DEFAULT_INSTANCE)); + + } + private File generateQueries(int count, boolean random) throws IOException { Random r = new Random(); @@ -468,7 +609,7 @@ public class FQLReplayTest return new QueryState(clientState); } - private static ResultHandler.ComparableResultSet createResultSet(int columnCount, int rowCount, boolean random) + static ResultHandler.ComparableResultSet createResultSet(int columnCount, int rowCount, boolean random) { List> columnDefs = new ArrayList<>(columnCount); Random r = new Random(); @@ -476,24 +617,27 @@ public class FQLReplayTest { columnDefs.add(Pair.create("a" + i, "int")); } - List> rows = new ArrayList<>(); + ResultHandler.ComparableColumnDefinitions colDefs = new StoredResultSet.StoredComparableColumnDefinitions(columnDefs, false, null); + List rows = new ArrayList<>(); for (int i = 0; i < rowCount; i++) { - List row = new ArrayList<>(columnCount); + List row = new ArrayList<>(columnCount); for (int jj = 0; jj < columnCount; jj++) - row.add(i + " col " + jj + (random ? r.nextInt() : "")); - rows.add(row); + row.add(ByteBufferUtil.bytes(i + " col " + jj + (random ? r.nextInt() : ""))); + + rows.add(new StoredResultSet.StoredComparableRow(row, colDefs)); } - return new FakeResultSet(columnDefs, rows); + return new StoredResultSet(colDefs, true, false, null, rows::iterator); } private static void compareWithFile(List dirs, File resultDir, List>> resultSets, int idx) { List> results1 = readResultFile(dirs.get(idx), resultDir); - for (int i = 0; i < results1.size(); i++) + for (int i = 0; i < resultSets.size(); i++) { - assertEquals(results1.get(i).left, resultSets.get(i).left); - assertEquals(results1.get(i).right, resultSets.get(i).right.get(idx)); + FQLQuery query = resultSets.get(i).left; + List> toCompare = Collections.singletonList(Pair.create(query, resultSets.get(i).right.get(idx))); + compareResults(Collections.singletonList(results1.get(i)), toCompare); } } @@ -503,258 +647,29 @@ public class FQLReplayTest try (ChronicleQueue q = ChronicleQueueBuilder.single(dir).build(); ChronicleQueue queryQ = ChronicleQueueBuilder.single(queryDir).build()) { - ExcerptTailer tailer = q.createTailer(); ExcerptTailer queryTailer = queryQ.createTailer(); - List> columnDefinitions = new ArrayList<>(); - List> rowColumns = new ArrayList<>(); - AtomicBoolean allRowsRead = new AtomicBoolean(false); - AtomicBoolean failedQuery = new AtomicBoolean(false); - while (tailer.readDocument(wire -> { - String type = wire.read("type").text(); - if (type.equals("column_definitions")) - { - int columnCount = wire.read("column_count").int32(); - for (int i = 0; i < columnCount; i++) - { - ValueIn vi = wire.read("column_definition"); - String name = vi.text(); - String dataType = vi.text(); - columnDefinitions.add(Pair.create(name, dataType)); - } - } - else if (type.equals("row")) - { - int rowColumnCount = wire.read("row_column_count").int32(); - List r = new ArrayList<>(rowColumnCount); - for (int i = 0; i < rowColumnCount; i++) - { - byte[] b = wire.read("column").bytes(); - r.add(new String(b)); - } - rowColumns.add(r); - } - else if (type.equals("end_resultset")) - { - allRowsRead.set(true); - } - else if (type.equals("query_failed")) - { - failedQuery.set(true); - } - })) + FQLQueryReader queryReader = new FQLQueryReader(); + Compare.StoredResultSetIterator resultSetIterator = new Compare.StoredResultSetIterator(q.createTailer()); + // we need to materialize the rows in-memory to compare them easier in these tests + while (resultSetIterator.hasNext()) { - if (allRowsRead.get()) + ResultHandler.ComparableResultSet resultSetFromDisk = resultSetIterator.next(); + Iterator rowIterFromDisk = resultSetFromDisk.iterator(); + queryTailer.readDocument(queryReader); + + FQLQuery query = queryReader.getQuery(); + List rows = new ArrayList<>(); + while (rowIterFromDisk.hasNext()) { - FQLQueryReader reader = new FQLQueryReader(); - queryTailer.readDocument(reader); - resultSets.add(Pair.create(reader.getQuery(), failedQuery.get() ? FakeResultSet.failed(new RuntimeException("failure")) - : new FakeResultSet(ImmutableList.copyOf(columnDefinitions), ImmutableList.copyOf(rowColumns)))); - allRowsRead.set(false); - failedQuery.set(false); - columnDefinitions.clear(); - rowColumns.clear(); + rows.add(rowIterFromDisk.next()); } + resultSets.add(Pair.create(query, new StoredResultSet(resultSetFromDisk.getColumnDefinitions(), + resultSetIterator.hasNext(), + resultSetFromDisk.wasFailed(), + resultSetFromDisk.getFailureException(), + rows::iterator))); } } return resultSets; } - - private static class FakeResultSet implements ResultHandler.ComparableResultSet - { - private final List> cdStrings; - private final List> rows; - private final Throwable ex; - - public FakeResultSet(List> cdStrings, List> rows) - { - this(cdStrings, rows, null); - } - - public FakeResultSet(List> cdStrings, List> rows, Throwable ex) - { - this.cdStrings = cdStrings; - this.rows = rows; - this.ex = ex; - } - - public static FakeResultSet failed(Throwable ex) - { - return new FakeResultSet(null, null, ex); - } - - public ResultHandler.ComparableColumnDefinitions getColumnDefinitions() - { - return new FakeComparableColumnDefinitions(cdStrings, wasFailed()); - } - - public boolean wasFailed() - { - return getFailureException() != null; - } - - public Throwable getFailureException() - { - return ex; - } - - public Iterator iterator() - { - if (wasFailed()) - return Collections.emptyListIterator(); - return new AbstractIterator() - { - Iterator> iter = rows.iterator(); - protected ResultHandler.ComparableRow computeNext() - { - if (iter.hasNext()) - return new FakeComparableRow(iter.next(), cdStrings); - return endOfData(); - } - }; - } - - public boolean equals(Object o) - { - if (this == o) return true; - if (!(o instanceof FakeResultSet)) return false; - FakeResultSet that = (FakeResultSet) o; - if (wasFailed() && that.wasFailed()) - return true; - return Objects.equals(cdStrings, that.cdStrings) && - Objects.equals(rows, that.rows); - } - - public int hashCode() - { - return Objects.hash(cdStrings, rows); - } - - public String toString() - { - return "FakeResultSet{" + - "cdStrings=" + cdStrings + - ", rows=" + rows + - '}'; - } - } - - private static class FakeComparableRow implements ResultHandler.ComparableRow - { - private final List row; - private final List> cds; - - public FakeComparableRow(List row, List> cds) - { - this.row = row; - this.cds = cds; - } - - public ByteBuffer getBytesUnsafe(int i) - { - return ByteBufferUtil.bytes(row.get(i)); - } - - public ResultHandler.ComparableColumnDefinitions getColumnDefinitions() - { - return new FakeComparableColumnDefinitions(cds, false); - } - - public boolean equals(Object other) - { - if (!(other instanceof FakeComparableRow)) - return false; - return row.equals(((FakeComparableRow)other).row); - } - - public String toString() - { - return row.toString(); - } - } - - private static class FakeComparableColumnDefinitions implements ResultHandler.ComparableColumnDefinitions - { - private final List defs; - private final boolean failed; - public FakeComparableColumnDefinitions(List> cds, boolean failed) - { - defs = cds != null ? cds.stream().map(FakeComparableDefinition::new).collect(Collectors.toList()) : null; - this.failed = failed; - } - - public List asList() - { - if (wasFailed()) - return Collections.emptyList(); - return defs; - } - - public boolean wasFailed() - { - return failed; - } - - public int size() - { - return defs.size(); - } - - public Iterator iterator() - { - if (wasFailed()) - return Collections.emptyListIterator(); - return new AbstractIterator() - { - Iterator iter = defs.iterator(); - protected ResultHandler.ComparableDefinition computeNext() - { - if (iter.hasNext()) - return iter.next(); - return endOfData(); - } - }; - } - public boolean equals(Object other) - { - if (!(other instanceof FakeComparableColumnDefinitions)) - return false; - return defs.equals(((FakeComparableColumnDefinitions)other).defs); - } - - public String toString() - { - return defs.toString(); - } - } - - private static class FakeComparableDefinition implements ResultHandler.ComparableDefinition - { - private final Pair p; - - public FakeComparableDefinition(Pair p) - { - this.p = p; - } - public String getType() - { - return p.right; - } - - public String getName() - { - return p.left; - } - - public boolean equals(Object other) - { - if (!(other instanceof FakeComparableDefinition)) - return false; - return p.equals(((FakeComparableDefinition)other).p); - } - - public String toString() - { - return getName() + ':' + getType(); - } - } }