Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 3 additions & 1 deletion builtin-adapter/BuiltinAdapter.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -102,7 +102,9 @@ JNIEXPORT jobject JNICALL Java_ru_rt_restream_reindexer_binding_builtin_BuiltinA
reindexer_string dsn = rx_string(env, path);
reindexer_string vers = rx_string(env, version);
reindexer_error error = reindexer_connect(rx, dsn, ConnectOpts(), vers, BindingCapabilities(
kBindingCapabilityResultsWithShardIDs
kBindingCapabilityQrIdleTimeouts
| kBindingCapabilityResultsWithShardIDs
| kBindingCapabilityIncarnationTags
| kBindingCapabilityComplexRank
| kBindingCapabilityQueryFormatV2));
env->ReleaseStringUTFChars(path, reinterpret_cast<const char *>(dsn.p));
Expand Down
61 changes: 37 additions & 24 deletions src/main/java/ru/rt/restream/reindexer/Query.java
Original file line number Diff line number Diff line change
Expand Up @@ -199,6 +199,8 @@ public enum Condition {

private final List<String> joinFields = new ArrayList<>();

private Query<?> lastJoinQuery;

private final List<Query<?>> mergeQueries = new ArrayList<>();

private final List<ReindexerNamespace<?>> namespaces = new ArrayList<>();
Expand Down Expand Up @@ -296,6 +298,11 @@ public <J> Query<T> leftJoin(Query<J> joinQuery, String field) {
}

private <J> Query<T> join(Query<J> joinQuery, String field, int joinType) {
if (root != null) {
@SuppressWarnings("unchecked")
Query<T> rootQuery = (Query<T>) root;
return rootQuery.join(joinQuery, field, joinType);
}
logBuilder.join(joinQuery.logBuilder, joinType);
if (joinQuery.root != null) {
throw new IllegalStateException("query.join call on already joined query. You should create new Query");
Expand All @@ -311,6 +318,7 @@ private <J> Query<T> join(Query<J> joinQuery, String field, int joinType) {
joinQuery.root = this;
joinQueries.add(joinQuery);
joinFields.add(field);
lastJoinQuery = joinQuery;
return this;
}

Expand All @@ -323,12 +331,13 @@ private <J> Query<T> join(Query<J> joinQuery, String field, int joinType) {
* @return the {@link Query} for further customizations
*/
public Query<T> on(String joinField, Condition condition, String joinIndex) {
logBuilder.on(nextOperation, joinField, condition.code, joinIndex);
buffer.putVarUInt32(QUERY_JOIN_ON);
buffer.putVarUInt32(nextOperation);
buffer.putVarUInt32(condition.code);
buffer.putVString(joinField);
buffer.putVString(joinIndex);
Query<?> target = lastJoinQuery != null ? lastJoinQuery : this;
target.logBuilder.on(nextOperation, joinField, condition.code, joinIndex);
target.buffer.putVarUInt32(QUERY_JOIN_ON);
target.buffer.putVarUInt32(nextOperation);
target.buffer.putVarUInt32(condition.code);
target.buffer.putVString(joinField);
target.buffer.putVString(joinIndex);
nextOperation = OP_AND;
return this;
}
Expand Down Expand Up @@ -1317,27 +1326,23 @@ public byte[] bytes() {
}

private byte[] toSubQueryBytes(int formatVersion) {
byte[] queryBytes = getQueryBytes(formatVersion);
if (formatVersion == QUERY_FORMAT_V2 || hasNestedJoins()) {
ByteBuffer copy = new ByteBuffer(queryBytes);
copy.putVarUInt32(QUERY_END);
copy.putVarUInt32(0);
copy.putVarUInt32(0);
return copy.bytes();
if (formatVersion == QUERY_FORMAT_V2) {
return serializeQuery(formatVersion);
}
if (!joinQueries.isEmpty() || !mergeQueries.isEmpty()) {
throw new IllegalStateException("Join and merge queries in subquery are not supported by QueryFormatV1");
}
return queryBytes;
return getQueryBytes(formatVersion);
}

private byte[] toExecutableBytes() {
int formatVersion = reindexer.getBinding().queryFormatVersion();
ByteBuffer queryBuffer = new ByteBuffer(getQueryBytes(formatVersion));
queryBuffer.putVarUInt32(QUERY_END);
if (formatVersion == QUERY_FORMAT_V2) {
appendJoinQueries(queryBuffer, new ArrayList<>(), formatVersion);
appendMergeQueries(queryBuffer, new ArrayList<>(), formatVersion);
} else {
appendJoinQueriesV1(queryBuffer, false);
return serializeQuery(formatVersion);
}
ByteBuffer queryBuffer = new ByteBuffer(getQueryBytes(formatVersion));
queryBuffer.putVarUInt32(QUERY_END);
appendJoinQueriesV1(queryBuffer, false);
return queryBuffer.bytes();
}

Expand Down Expand Up @@ -1367,6 +1372,14 @@ private void appendMergeQueries(ByteBuffer target, List<ReindexerNamespace<?>> t
}
}

private byte[] serializeQuery(int formatVersion) {
ByteBuffer queryBuffer = new ByteBuffer(getQueryBytes(formatVersion));
queryBuffer.putVarUInt32(QUERY_END);
appendJoinQueries(queryBuffer, new ArrayList<>(), formatVersion);
appendMergeQueries(queryBuffer, new ArrayList<>(), formatVersion);
return queryBuffer.bytes();
}

private void appendQuery(ByteBuffer target, Query<?> query, int queryJoinType,
List<ReindexerNamespace<?>> targetNamespaces, int formatVersion) {
if (queryJoinType != MERGE) {
Expand All @@ -1382,7 +1395,7 @@ private void appendQuery(ByteBuffer target, Query<?> query, int queryJoinType,
}

private void appendJoinQueriesV1(ByteBuffer target, boolean appendNamespaces) {
if (hasNestedJoins()) {
if (hasNestedQueries()) {
throw new IllegalStateException("Nested joins are not supported by QueryFormatV1");
}
for (Query<?> joinQuery : joinQueries) {
Expand Down Expand Up @@ -1410,14 +1423,14 @@ private void appendMergeQueriesV1(ByteBuffer target) {
}
}

private boolean hasNestedJoins() {
private boolean hasNestedQueries() {
for (Query<?> joinQuery : joinQueries) {
if (!joinQuery.joinQueries.isEmpty() || joinQuery.hasNestedJoins()) {
if (!joinQuery.joinQueries.isEmpty() || !joinQuery.mergeQueries.isEmpty() || joinQuery.hasNestedQueries()) {
return true;
}
}
for (Query<?> mergeQuery : mergeQueries) {
if (mergeQuery.hasNestedJoins()) {
if (!mergeQuery.mergeQueries.isEmpty() || mergeQuery.hasNestedQueries()) {
return true;
}
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@
import static ru.rt.restream.reindexer.binding.Consts.QUERY_RESULT_RANK_FORMAT;
import static ru.rt.restream.reindexer.binding.Consts.QUERY_RESULT_SHARDING_VERSION;
import static ru.rt.restream.reindexer.binding.Consts.QUERY_RESULT_SHARD_ID;
import static ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V1;
import static ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V2;
import static ru.rt.restream.reindexer.binding.Consts.RANK_FORMAT_SINGLE_FLOAT;
import static ru.rt.restream.reindexer.binding.Consts.RESULTS_FORMAT_MASK;
Expand Down Expand Up @@ -60,7 +61,11 @@ public class QueryResultReader {
* @return the {@link QueryResult} to use
*/
public QueryResult read(byte[] rawQueryResult) {
return read(rawQueryResult, Consts.QUERY_FORMAT_V1);
ByteBuffer buffer = new ByteBuffer(rawQueryResult).rewind();
int queryFormatVersion = buffer.getVarUInt() == QUERY_FORMAT_V2
? QUERY_FORMAT_V2
: QUERY_FORMAT_V1;
return read(rawQueryResult, queryFormatVersion);
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,6 @@
import ru.rt.restream.reindexer.binding.RequestContext;
import ru.rt.restream.reindexer.util.NativeUtils;

import static ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V2;

/**
* A request context which is holds a {@link QueryResult},
* the {@link #fetchResults(int, int)} method is NOOP since Builtin does not support it.
Expand Down Expand Up @@ -54,7 +52,7 @@ public BuiltinRequestContext(ReindexerResponse response) {
}
}
QueryResultReader reader = new QueryResultReader();
queryResult = reader.read(rawQueryResult, QUERY_FORMAT_V2);
queryResult = reader.read(rawQueryResult);
queryResult.setResultsPtr(resultsPtr);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@
import static ru.rt.restream.reindexer.binding.Consts.BINDING_CAPABILITY_COMPLEX_RANK;
import static ru.rt.restream.reindexer.binding.Consts.BINDING_CAPABILITY_NAMESPACE_INCARNATIONS;
import static ru.rt.restream.reindexer.binding.Consts.BINDING_CAPABILITY_QUERY_FORMAT_V2;
import static ru.rt.restream.reindexer.binding.Consts.BINDING_CAPABILITY_QR_IDLE_TIMEOUTS;
import static ru.rt.restream.reindexer.binding.Consts.BINDING_CAPABILITY_RESULTS_WITH_SHARD_IDS;
import static ru.rt.restream.reindexer.binding.Consts.DEF_APP_NAME;
import static ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V1;
Expand Down Expand Up @@ -144,7 +145,8 @@ public PhysicalConnection(String host, int port, String user, String password, S
-1, // expectedClusterID
REINDEXER_VERSION,
getAppName(),
BINDING_CAPABILITY_RESULTS_WITH_SHARD_IDS
BINDING_CAPABILITY_QR_IDLE_TIMEOUTS
| BINDING_CAPABILITY_RESULTS_WITH_SHARD_IDS
| BINDING_CAPABILITY_COMPLEX_RANK
| BINDING_CAPABILITY_NAMESPACE_INCARNATIONS
| BINDING_CAPABILITY_QUERY_FORMAT_V2);
Expand Down
Loading
Loading