From 9bb1f9fa16c32213f33cb533b250eb1f0e436a6c Mon Sep 17 00:00:00 2001 From: Stanislav Varganov Date: Tue, 11 Aug 2026 18:36:46 +0300 Subject: [PATCH] Nested join queries support --- builtin-adapter/BuiltinAdapter.cpp | 15 +- builtin-adapter/BuiltinAdapter.h | 5 +- .../java/ru/rt/restream/reindexer/Query.java | 40 ++- .../reindexer/QueryResultIterator.java | 15 +- .../restream/reindexer/binding/Binding.java | 9 + .../reindexer/binding/QueryResultReader.java | 3 +- .../reindexer/binding/builtin/Builtin.java | 52 ++- .../binding/builtin/BuiltinAdapter.java | 10 +- .../builtin/BuiltinRequestContext.java | 6 +- .../builtin/BuiltinTransactionContext.java | 7 +- .../binding/builtin/server/BuiltinServer.java | 5 + .../reindexer/binding/cproto/Connection.java | 9 + .../binding/cproto/ConnectionPool.java | 4 + .../reindexer/binding/cproto/Cproto.java | 5 + .../binding/cproto/CprotoRequestContext.java | 9 +- .../binding/cproto/PhysicalConnection.java | 123 +++++-- .../reindexer/connector/JoinTest.java | 330 ++++++++++++++++++ 17 files changed, 576 insertions(+), 71 deletions(-) diff --git a/builtin-adapter/BuiltinAdapter.cpp b/builtin-adapter/BuiltinAdapter.cpp index eac37126..e29913ff 100644 --- a/builtin-adapter/BuiltinAdapter.cpp +++ b/builtin-adapter/BuiltinAdapter.cpp @@ -91,6 +91,10 @@ JNIEXPORT jlong JNICALL Java_ru_rt_restream_reindexer_binding_builtin_BuiltinAda return init_reindexer(); } +JNIEXPORT jstring JNICALL Java_ru_rt_restream_reindexer_binding_builtin_BuiltinAdapter_version(JNIEnv *env, jobject) { + return env->NewStringUTF(reindexer_version()); +} + JNIEXPORT void JNICALL Java_ru_rt_restream_reindexer_binding_builtin_BuiltinAdapter_destroy(JNIEnv *, jobject, jlong rx) { destroy_reindexer(rx); @@ -98,13 +102,16 @@ JNIEXPORT void JNICALL Java_ru_rt_restream_reindexer_binding_builtin_BuiltinAdap JNIEXPORT jobject JNICALL Java_ru_rt_restream_reindexer_binding_builtin_BuiltinAdapter_connect(JNIEnv *env, jobject, jlong rx, jstring path, - jstring version) { + jstring version, + jint queryFormatVersion) { reindexer_string dsn = rx_string(env, path); reindexer_string vers = rx_string(env, version); + int64_t capabilities = kBindingCapabilityResultsWithShardIDs | kBindingCapabilityComplexRank; + if (queryFormatVersion == QueryFormatV2) { + capabilities |= kBindingCapabilityQueryFormatV2; + } reindexer_error error = reindexer_connect(rx, dsn, ConnectOpts(), vers, BindingCapabilities( - kBindingCapabilityResultsWithShardIDs - | kBindingCapabilityComplexRank - | kBindingCapabilityQueryFormatV2)); + capabilities)); env->ReleaseStringUTFChars(path, reinterpret_cast(dsn.p)); env->ReleaseStringUTFChars(version, reinterpret_cast(vers.p)); return j_res(env, error); diff --git a/builtin-adapter/BuiltinAdapter.h b/builtin-adapter/BuiltinAdapter.h index b60ec97a..059488e9 100644 --- a/builtin-adapter/BuiltinAdapter.h +++ b/builtin-adapter/BuiltinAdapter.h @@ -24,10 +24,13 @@ extern "C" { JNIEXPORT jlong JNICALL Java_ru_rt_restream_reindexer_binding_builtin_BuiltinAdapter_init(JNIEnv *, jobject); +JNIEXPORT jstring JNICALL Java_ru_rt_restream_reindexer_binding_builtin_BuiltinAdapter_version(JNIEnv *, jobject); + JNIEXPORT void JNICALL Java_ru_rt_restream_reindexer_binding_builtin_BuiltinAdapter_destroy(JNIEnv *, jobject, jlong); JNIEXPORT jobject JNICALL Java_ru_rt_restream_reindexer_binding_builtin_BuiltinAdapter_connect(JNIEnv *, jobject, jlong, - jstring, jstring); + jstring, jstring, + jint); JNIEXPORT jobject JNICALL Java_ru_rt_restream_reindexer_binding_builtin_BuiltinAdapter_openNamespace(JNIEnv *, jobject, jlong, jlong, diff --git a/src/main/java/ru/rt/restream/reindexer/Query.java b/src/main/java/ru/rt/restream/reindexer/Query.java index e2885835..cb21052f 100644 --- a/src/main/java/ru/rt/restream/reindexer/Query.java +++ b/src/main/java/ru/rt/restream/reindexer/Query.java @@ -1317,27 +1317,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(); } @@ -1367,6 +1363,14 @@ private void appendMergeQueries(ByteBuffer target, List> 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> targetNamespaces, int formatVersion) { if (queryJoinType != MERGE) { @@ -1382,7 +1386,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) { @@ -1410,14 +1414,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; } } diff --git a/src/main/java/ru/rt/restream/reindexer/QueryResultIterator.java b/src/main/java/ru/rt/restream/reindexer/QueryResultIterator.java index 134b4ce8..f9e08416 100644 --- a/src/main/java/ru/rt/restream/reindexer/QueryResultIterator.java +++ b/src/main/java/ru/rt/restream/reindexer/QueryResultIterator.java @@ -193,10 +193,19 @@ private void readJoinedItems(Object item, Query queryContext, int nsId) { for (int i = 0; i < itemsCount; i++) { subItems.add(readItem(joinedNamespace, joinedItemReader, joinQuery)); } - subItemsMap.computeIfAbsent(queryContext.getJoinFields().get(joinedField), field -> new ArrayList<>()) - .addAll(subItems); + + String joinField = queryContext.getJoinFields().get(joinedField); + List fieldSubItems = subItemsMap.get(joinField); + if (fieldSubItems == null) { + fieldSubItems = new ArrayList<>(); + subItemsMap.put(joinField, fieldSubItems); + } + fieldSubItems.addAll(subItems); + } + + for (Map.Entry> entry : subItemsMap.entrySet()) { + writeJoinResult(item, entry.getKey(), entry.getValue()); } - subItemsMap.forEach((key, value) -> writeJoinResult(item, key, value)); } private void readJoinedItemsV1(Object item, int nsId) { diff --git a/src/main/java/ru/rt/restream/reindexer/binding/Binding.java b/src/main/java/ru/rt/restream/reindexer/binding/Binding.java index 50c40299..7e1ee3df 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/Binding.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/Binding.java @@ -214,6 +214,15 @@ default int queryFormatVersion() { return Consts.QUERY_FORMAT_V1; } + /** + * Returns true if connected Reindexer server supports nested join queries. + * + * @return true if connected Reindexer server supports nested join queries + */ + default boolean supportsNestedJoinQueries() { + return false; + } + /** * Closes binding to Reindexer instance. */ diff --git a/src/main/java/ru/rt/restream/reindexer/binding/QueryResultReader.java b/src/main/java/ru/rt/restream/reindexer/binding/QueryResultReader.java index 62aa2779..ddb5aa16 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/QueryResultReader.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/QueryResultReader.java @@ -75,8 +75,7 @@ public QueryResult read(byte[] rawQueryResult, int queryFormatVersion) { if (queryFormatVersion == QUERY_FORMAT_V2) { long format = buffer.getVarUInt(); if (format != QUERY_FORMAT_V2) { - String errorMessage = String.format("QueryResults format version='%d' is not supported", format); - throw new RuntimeException(errorMessage); + buffer.rewind(); } } QueryResult queryResult = getQueryResultWithFlags(buffer.getVarUInt()); diff --git a/src/main/java/ru/rt/restream/reindexer/binding/builtin/Builtin.java b/src/main/java/ru/rt/restream/reindexer/binding/builtin/Builtin.java index 5cbaa66e..82dbaa84 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/builtin/Builtin.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/builtin/Builtin.java @@ -45,6 +45,12 @@ public class Builtin implements Binding { private static final Logger LOGGER = LoggerFactory.getLogger(Builtin.class); + private static final int NESTED_JOIN_QUERIES_MIN_MAJOR = 5; + + private static final int NESTED_JOIN_QUERIES_MIN_MINOR = 16; + + private static final int NESTED_JOIN_QUERIES_MIN_PATCH = 0; + private final AtomicLong next = new AtomicLong(0L); private final Gson gson = new GsonBuilder() @@ -57,6 +63,10 @@ public class Builtin implements Binding { private final Duration timeout; + private final boolean supportsNestedJoinQueries; + + private final int queryFormatVersion; + /** * Creates an instance. * @@ -67,9 +77,11 @@ public Builtin(URI uri, Duration requestTimeout) { adapter = new BuiltinAdapter(); timeout = requestTimeout; rx = adapter.init(); + queryFormatVersion = QUERY_FORMAT_V2; + supportsNestedJoinQueries = isNestedJoinQueriesSupported(adapter.version()); String path = uri.getPath(); try { - ReindexerResponse response = adapter.connect(rx, path, REINDEXER_VERSION); + ReindexerResponse response = adapter.connect(rx, path, REINDEXER_VERSION, queryFormatVersion); checkResponse(response); } catch (Exception e) { LOGGER.error("rx: connect error", e); @@ -89,6 +101,8 @@ public Builtin(BuiltinAdapter adapter, long rx, Duration timeout) { this.adapter = adapter; this.rx = rx; this.timeout = timeout; + queryFormatVersion = QUERY_FORMAT_V2; + supportsNestedJoinQueries = isNestedJoinQueriesSupported(adapter.version()); } @Override @@ -152,7 +166,7 @@ public RequestContext select(String query, boolean asJson, int fetchCount, long[ ReindexerResponse response = adapter.select(rx, next.getAndIncrement(), timeout.toMillis(), query, asJson, ptVersions); checkResponse(response); - return new BuiltinRequestContext(response); + return new BuiltinRequestContext(response, queryFormatVersion); } @Override @@ -160,7 +174,7 @@ public RequestContext selectQuery(byte[] queryData, int fetchCount, long[] ptVer ReindexerResponse response = adapter.selectQuery(rx, next.getAndIncrement(), timeout.toMillis(), queryData, ptVersions, asJson); checkResponse(response); - return new BuiltinRequestContext(response); + return new BuiltinRequestContext(response, queryFormatVersion); } @Override @@ -187,7 +201,7 @@ public TransactionContext beginTx(String namespaceName) { txId = (long) arg; } } - return new BuiltinTransactionContext(adapter, rx, txId, next::getAndIncrement, timeout); + return new BuiltinTransactionContext(adapter, rx, txId, next::getAndIncrement, timeout, queryFormatVersion); } @Override @@ -210,7 +224,35 @@ private void checkResponse(ReindexerResponse response) { @Override public int queryFormatVersion() { - return QUERY_FORMAT_V2; + return queryFormatVersion; + } + + @Override + public boolean supportsNestedJoinQueries() { + return supportsNestedJoinQueries; + } + + private boolean isNestedJoinQueriesSupported(String version) { + int[] parsedVersion = parseVersion(version); + if (parsedVersion[0] != NESTED_JOIN_QUERIES_MIN_MAJOR) { + return parsedVersion[0] > NESTED_JOIN_QUERIES_MIN_MAJOR; + } + if (parsedVersion[1] != NESTED_JOIN_QUERIES_MIN_MINOR) { + return parsedVersion[1] > NESTED_JOIN_QUERIES_MIN_MINOR; + } + return parsedVersion[2] >= NESTED_JOIN_QUERIES_MIN_PATCH; + } + + private int[] parseVersion(String version) { + String normalized = version.startsWith("v") ? version.substring(1) : version; + String[] parts = normalized.split("\\D+"); + int[] result = new int[3]; + for (int i = 0; i < result.length && i < parts.length; i++) { + if (!parts[i].isEmpty()) { + result[i] = Integer.parseInt(parts[i]); + } + } + return result; } @Override diff --git a/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinAdapter.java b/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinAdapter.java index 6aa95340..22d5ca96 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinAdapter.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinAdapter.java @@ -72,6 +72,13 @@ private static void loadLibrary(String fileName) throws IOException { */ public native long init(); + /** + * Returns native Reindexer library version. + * + * @return the native Reindexer library version + */ + public native String version(); + /** * Destroys Reindexer instance. * @@ -135,9 +142,10 @@ private static void loadLibrary(String fileName) throws IOException { * @param rx the Reindexer instance pointer * @param path the Reindexer's database path * @param version the Reindexer's version + * @param queryFormatVersion the query serialization format version * @return the {@link ReindexerResponse} to use */ - public native ReindexerResponse connect(long rx, String path, String version); + public native ReindexerResponse connect(long rx, String path, String version, int queryFormatVersion); /** * Opens a namespace. diff --git a/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinRequestContext.java b/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinRequestContext.java index 6bc7cc24..3e71298a 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinRequestContext.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinRequestContext.java @@ -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. @@ -37,7 +35,7 @@ public class BuiltinRequestContext implements RequestContext { * * @param response the {@link ReindexerResponse} to use */ - public BuiltinRequestContext(ReindexerResponse response) { + public BuiltinRequestContext(ReindexerResponse response, int queryFormatVersion) { long resultsPtr = 0L; byte[] rawQueryResult = new byte[0]; Object[] arguments = response.getArguments(); @@ -54,7 +52,7 @@ public BuiltinRequestContext(ReindexerResponse response) { } } QueryResultReader reader = new QueryResultReader(); - queryResult = reader.read(rawQueryResult, QUERY_FORMAT_V2); + queryResult = reader.read(rawQueryResult, queryFormatVersion); queryResult.setResultsPtr(resultsPtr); } diff --git a/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinTransactionContext.java b/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinTransactionContext.java index fe3a9d9d..8dc57adc 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinTransactionContext.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/builtin/BuiltinTransactionContext.java @@ -46,6 +46,8 @@ public class BuiltinTransactionContext implements TransactionContext { private final Duration timeout; + private final int queryFormatVersion; + /** * Creates an instance. * @@ -56,12 +58,13 @@ public class BuiltinTransactionContext implements TransactionContext { * @param timeout the execution timeout */ public BuiltinTransactionContext(BuiltinAdapter adapter, long rx, long transactionId, - Supplier next, Duration timeout) { + Supplier next, Duration timeout, int queryFormatVersion) { this.adapter = adapter; this.rx = rx; this.transactionId = transactionId; this.next = next; this.timeout = timeout; + this.queryFormatVersion = queryFormatVersion; } @Override @@ -91,7 +94,7 @@ private ReindexerResponse modifyItemInternal(byte[] data, int format, int mode, public RequestContext selectQuery(byte[] queryData, int fetchCount, long[] ptVersions, boolean asJson) { ReindexerResponse response = adapter.selectQuery(rx, next.get(), timeout.toMillis(), queryData, ptVersions, asJson); checkResponse(response); - return new BuiltinRequestContext(response); + return new BuiltinRequestContext(response, queryFormatVersion); } @Override diff --git a/src/main/java/ru/rt/restream/reindexer/binding/builtin/server/BuiltinServer.java b/src/main/java/ru/rt/restream/reindexer/binding/builtin/server/BuiltinServer.java index f83cdde1..92398da2 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/builtin/server/BuiltinServer.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/builtin/server/BuiltinServer.java @@ -196,6 +196,11 @@ public int queryFormatVersion() { return builtin.queryFormatVersion(); } + @Override + public boolean supportsNestedJoinQueries() { + return builtin.supportsNestedJoinQueries(); + } + @Override public void close() { ReindexerResponse response = adapter.stopServer(svc); diff --git a/src/main/java/ru/rt/restream/reindexer/binding/cproto/Connection.java b/src/main/java/ru/rt/restream/reindexer/binding/cproto/Connection.java index 1968d95d..436b8e93 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/cproto/Connection.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/cproto/Connection.java @@ -59,6 +59,15 @@ default int queryFormatVersion() { return ru.rt.restream.reindexer.binding.Consts.QUERY_FORMAT_V1; } + /** + * Returns true if connected Reindexer server supports nested join queries. + * + * @return true if connected Reindexer server supports nested join queries + */ + default boolean supportsNestedJoinQueries() { + return false; + } + /** * Closes the connection. */ diff --git a/src/main/java/ru/rt/restream/reindexer/binding/cproto/ConnectionPool.java b/src/main/java/ru/rt/restream/reindexer/binding/cproto/ConnectionPool.java index 984dc842..7c2490b6 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/cproto/ConnectionPool.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/cproto/ConnectionPool.java @@ -149,6 +149,10 @@ public int queryFormatVersion() { return getConnection().queryFormatVersion(); } + public boolean supportsNestedJoinQueries() { + return getConnection().supportsNestedJoinQueries(); + } + private DataSource getDataSource(int connectionPoolSize) { Instant connectionDeadline = Instant.now().plus(timeout); for (; ; ) { diff --git a/src/main/java/ru/rt/restream/reindexer/binding/cproto/Cproto.java b/src/main/java/ru/rt/restream/reindexer/binding/cproto/Cproto.java index 531858bf..44e7406a 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/cproto/Cproto.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/cproto/Cproto.java @@ -196,6 +196,11 @@ public int queryFormatVersion() { return pool.queryFormatVersion(); } + @Override + public boolean supportsNestedJoinQueries() { + return pool.supportsNestedJoinQueries(); + } + /** * Closes the connection pool. */ diff --git a/src/main/java/ru/rt/restream/reindexer/binding/cproto/CprotoRequestContext.java b/src/main/java/ru/rt/restream/reindexer/binding/cproto/CprotoRequestContext.java index 4a4df39f..3baeebd6 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/cproto/CprotoRequestContext.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/cproto/CprotoRequestContext.java @@ -79,7 +79,10 @@ public void fetchResults(int offset, int limit) { : Consts.RESULTS_C_JSON | Consts.RESULTS_WITH_PAYLOAD_TYPES; int fetchCount = limit <= 0 ? Integer.MAX_VALUE : limit; ReindexerResponse rpcResponse = ConnectionUtils.rpcCall(connection, FETCH_RESULTS, requestId, flags, offset, fetchCount); - queryResult = getQueryResult(rpcResponse); + int fetchQueryFormatVersion = connection.supportsNestedJoinQueries() + ? queryFormatVersion + : Consts.QUERY_FORMAT_V1; + queryResult = getQueryResult(rpcResponse, fetchQueryFormatVersion); } /** @@ -97,6 +100,10 @@ public void closeResults() { } private QueryResult getQueryResult(ReindexerResponse rpcResponse) { + return getQueryResult(rpcResponse, queryFormatVersion); + } + + private QueryResult getQueryResult(ReindexerResponse rpcResponse, int queryFormatVersion) { byte[] rawQueryResult = new byte[0]; Object[] responseArguments = rpcResponse.getArguments(); if (responseArguments.length > 0) { diff --git a/src/main/java/ru/rt/restream/reindexer/binding/cproto/PhysicalConnection.java b/src/main/java/ru/rt/restream/reindexer/binding/cproto/PhysicalConnection.java index 1bb0e7f1..d4776be8 100644 --- a/src/main/java/ru/rt/restream/reindexer/binding/cproto/PhysicalConnection.java +++ b/src/main/java/ru/rt/restream/reindexer/binding/cproto/PhysicalConnection.java @@ -29,6 +29,7 @@ import java.io.DataOutputStream; import java.io.IOException; import java.net.Socket; +import java.nio.charset.StandardCharsets; import java.time.Duration; import java.time.Instant; import java.util.ArrayList; @@ -85,6 +86,12 @@ public class PhysicalConnection implements Connection { static final int CPROTO_VERSION_MASK = 0x3FF; + private static final int NESTED_JOIN_QUERIES_MIN_MAJOR = 5; + + private static final int NESTED_JOIN_QUERIES_MIN_MINOR = 16; + + private static final int NESTED_JOIN_QUERIES_MIN_PATCH = 0; + private final ReadWriteLock lock = new ReentrantReadWriteLock(); private final Condition notEmptyBuffer = lock.writeLock().newCondition(); @@ -95,59 +102,40 @@ public class PhysicalConnection implements Connection { private Exception error; - private final Socket clientSocket; + private Socket clientSocket; - private final DataOutputStream output; + private DataOutputStream output; - private final DataInputStream input; + private DataInputStream input; - private final Duration timeout; + private Duration timeout; - private final ScheduledExecutorService scheduler; + private ScheduledExecutorService scheduler; private final BlockingQueue sequences = new ArrayBlockingQueue<>(QUEUE_SIZE); private final List requests = new ArrayList<>(QUEUE_SIZE); - private final ScheduledFuture readTaskFuture; + private ScheduledFuture readTaskFuture; - private final ScheduledFuture writeTaskFuture; + private ScheduledFuture writeTaskFuture; private volatile int queryFormatVersion = QUERY_FORMAT_V1; + private volatile boolean supportsNestedJoinQueries = false; + public PhysicalConnection(String host, int port, String user, String password, String database, SSLSocketFactory sslSocketFactory, Duration requestTimeout, ScheduledExecutorService scheduler) { try { - if (sslSocketFactory != null) { - LOGGER.debug("rx: using SSL/TLS connection to {}:{}", host, port); - SSLSocket sslSocket = (SSLSocket) sslSocketFactory.createSocket(host, port); - // Fail fast if SSL/TLS handshake fails. - sslSocket.startHandshake(); - clientSocket = sslSocket; - } else { - clientSocket = new Socket(host, port); - } - output = new DataOutputStream(clientSocket.getOutputStream()); - input = new DataInputStream(clientSocket.getInputStream()); timeout = requestTimeout; this.scheduler = scheduler; for (int i = 0; i < QUEUE_SIZE; i++) { requests.add(new RpcRequest()); sequences.add(i); } - readTaskFuture = scheduler.scheduleWithFixedDelay(new ReadTask(), 0, 100, TimeUnit.MICROSECONDS); - writeTaskFuture = scheduler.scheduleWithFixedDelay(new WriteTask(), 0, 100, TimeUnit.MICROSECONDS); - ReindexerResponse response = ConnectionUtils.rpcCall(this, Binding.LOGIN, user, password, database, - false, // create DB if missing - false, // checkClusterID - -1, // expectedClusterID - REINDEXER_VERSION, - getAppName(), - BINDING_CAPABILITY_RESULTS_WITH_SHARD_IDS - | BINDING_CAPABILITY_COMPLEX_RANK - | BINDING_CAPABILITY_NAMESPACE_INCARNATIONS - | BINDING_CAPABILITY_QUERY_FORMAT_V2); + openSocket(host, port, sslSocketFactory); + ReindexerResponse response = login(user, password, database, true); updateQueryFormatVersion(response); } catch (Exception e) { onError(e); @@ -155,6 +143,43 @@ public PhysicalConnection(String host, int port, String user, String password, S } } + private void openSocket(String host, int port, SSLSocketFactory sslSocketFactory) throws IOException { + if (sslSocketFactory != null) { + LOGGER.debug("rx: using SSL/TLS connection to {}:{}", host, port); + SSLSocket sslSocket = (SSLSocket) sslSocketFactory.createSocket(host, port); + // Fail fast if SSL/TLS handshake fails. + sslSocket.startHandshake(); + clientSocket = sslSocket; + } else { + clientSocket = new Socket(host, port); + } + output = new DataOutputStream(clientSocket.getOutputStream()); + input = new DataInputStream(clientSocket.getInputStream()); + readTaskFuture = scheduler.scheduleWithFixedDelay(new ReadTask(), 0, 100, TimeUnit.MICROSECONDS); + writeTaskFuture = scheduler.scheduleWithFixedDelay(new WriteTask(), 0, 100, TimeUnit.MICROSECONDS); + } + + private ReindexerResponse login(String user, String password, String database, boolean queryFormatV2) { + return ConnectionUtils.rpcCall(this, Binding.LOGIN, user, password, database, + false, // create DB if missing + false, // checkClusterID + -1, // expectedClusterID + REINDEXER_VERSION, + getAppName(), + getBindingCapabilities(queryFormatV2)); + } + + private long getBindingCapabilities(boolean queryFormatV2) { + if (!queryFormatV2) { + return 0L; + } + long capabilities = BINDING_CAPABILITY_RESULTS_WITH_SHARD_IDS + | BINDING_CAPABILITY_COMPLEX_RANK + | BINDING_CAPABILITY_NAMESPACE_INCARNATIONS + | BINDING_CAPABILITY_QUERY_FORMAT_V2; + return capabilities; + } + private void updateQueryFormatVersion(ReindexerResponse response) { Object[] arguments = response.getArguments(); if (arguments.length > 2 && arguments[2] instanceof Long) { @@ -162,9 +187,38 @@ private void updateQueryFormatVersion(ReindexerResponse response) { queryFormatVersion = (capabilities & BINDING_CAPABILITY_QUERY_FORMAT_V2) != 0 ? QUERY_FORMAT_V2 : QUERY_FORMAT_V1; + supportsNestedJoinQueries = queryFormatVersion == QUERY_FORMAT_V2 + && isNestedJoinQueriesSupportedByServer(arguments[0]); } } + private boolean isNestedJoinQueriesSupportedByServer(Object serverVersionArgument) { + if (!(serverVersionArgument instanceof byte[])) { + return false; + } + String serverVersion = new String((byte[]) serverVersionArgument, StandardCharsets.UTF_8); + int[] version = parseVersion(serverVersion); + if (version[0] != NESTED_JOIN_QUERIES_MIN_MAJOR) { + return version[0] > NESTED_JOIN_QUERIES_MIN_MAJOR; + } + if (version[1] != NESTED_JOIN_QUERIES_MIN_MINOR) { + return version[1] > NESTED_JOIN_QUERIES_MIN_MINOR; + } + return version[2] >= NESTED_JOIN_QUERIES_MIN_PATCH; + } + + private int[] parseVersion(String version) { + String normalized = version.startsWith("v") ? version.substring(1) : version; + String[] parts = normalized.split("\\D+"); + int[] result = new int[3]; + for (int i = 0; i < result.length && i < parts.length; i++) { + if (!parts[i].isEmpty()) { + result[i] = Integer.parseInt(parts[i]); + } + } + return result; + } + private Object getAppName() { return System.getProperty(APP_PROPERTY_NAME, DEF_APP_NAME); } @@ -359,6 +413,11 @@ public int queryFormatVersion() { return queryFormatVersion; } + @Override + public boolean supportsNestedJoinQueries() { + return supportsNestedJoinQueries; + } + private Exception getCurrentError() { lock.readLock().lock(); try { @@ -412,6 +471,10 @@ private void onError(Exception error) { @Override public void close() { + disconnect(); + } + + private void disconnect() { if (readTaskFuture != null) { readTaskFuture.cancel(true); } diff --git a/src/test/java/ru/rt/restream/reindexer/connector/JoinTest.java b/src/test/java/ru/rt/restream/reindexer/connector/JoinTest.java index 0c093eef..1b899612 100644 --- a/src/test/java/ru/rt/restream/reindexer/connector/JoinTest.java +++ b/src/test/java/ru/rt/restream/reindexer/connector/JoinTest.java @@ -28,16 +28,22 @@ import ru.rt.restream.reindexer.ResultIterator; import ru.rt.restream.reindexer.annotations.Reindex; import ru.rt.restream.reindexer.annotations.Transient; +import ru.rt.restream.reindexer.binding.Consts; import ru.rt.restream.reindexer.db.DbBaseTest; import java.util.ArrayList; import java.util.Arrays; import java.util.Collections; +import java.util.HashMap; import java.util.List; +import java.util.Map; import static org.hamcrest.MatcherAssert.assertThat; +import static org.hamcrest.Matchers.containsString; import static org.hamcrest.Matchers.is; import static org.hamcrest.Matchers.nullValue; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assumptions.assumeTrue; import static ru.rt.restream.reindexer.Query.Condition.EQ; import static ru.rt.restream.reindexer.Query.Condition.RANGE; import static ru.rt.restream.reindexer.Query.Condition.SET; @@ -837,6 +843,304 @@ public void testMultipleJoinsSameFieldAccumulateResults() { assertThat(result.joinedActors.get(1).name, is("second")); } + @Test + public void testNestedJoin() { + assumeQueryFormatV2(); + + db.openNamespace("items_with_join", NamespaceOptions.defaultOptions(), ItemWithJoin.class); + db.openNamespace("actors", NamespaceOptions.defaultOptions(), Actor.class); + db.openNamespace("roles", NamespaceOptions.defaultOptions(), Role.class); + + ItemWithJoin item = new ItemWithJoin(); + item.id = 1; + item.name = "item"; + item.actorsIds = Collections.singletonList(10); + db.upsert("items_with_join", item); + + Actor actor = new Actor(); + actor.id = 10; + actor.name = "actor"; + actor.visible = true; + db.upsert("actors", actor); + + Role role = new Role(); + role.id = 100; + role.name = "actor"; + db.upsert("roles", role); + + Query actorQuery = db.query("actors", Actor.class) + .innerJoin(db.query("roles", Role.class), "joinedRoles") + .on("name", EQ, "name"); + + ResultIterator items = db.query("items_with_join", ItemWithJoin.class) + .where("id", EQ, 1) + .innerJoin(actorQuery, "joinedActors") + .on("actorsIds", SET, "id") + .execute(); + + assertThat(items.hasNext(), is(true)); + ItemWithJoin result = items.next(); + assertThat(result.joinedActors.size(), is(1)); + Actor joinedActor = result.joinedActors.get(0); + assertThat(joinedActor.id, is(10)); + assertThat(joinedActor.name, is("actor")); + assertThat(joinedActor.joinedRoles.size(), is(1)); + assertThat(joinedActor.joinedRoles.get(0).id, is(100)); + assertThat(joinedActor.joinedRoles.get(0).name, is("actor")); + } + + @Test + public void testNestedJoinWithLeftAndInnerCombinations() { + assumeQueryFormatV2(); + + db.openNamespace("items_with_join", NamespaceOptions.defaultOptions(), ItemWithJoin.class); + db.openNamespace("actors", NamespaceOptions.defaultOptions(), Actor.class); + db.openNamespace("roles", NamespaceOptions.defaultOptions(), Role.class); + + ItemWithJoin itemWithRole = new ItemWithJoin(); + itemWithRole.id = 1; + itemWithRole.actorsIds = Collections.singletonList(10); + db.upsert("items_with_join", itemWithRole); + + ItemWithJoin itemWithoutRole = new ItemWithJoin(); + itemWithoutRole.id = 2; + itemWithoutRole.actorsIds = Collections.singletonList(20); + db.upsert("items_with_join", itemWithoutRole); + + Actor actorWithRole = new Actor(); + actorWithRole.id = 10; + actorWithRole.name = "actor-with-role"; + db.upsert("actors", actorWithRole); + + Actor actorWithoutRole = new Actor(); + actorWithoutRole.id = 20; + actorWithoutRole.name = "actor-without-role"; + db.upsert("actors", actorWithoutRole); + + Role role = new Role(); + role.id = 100; + role.name = "actor-with-role"; + db.upsert("roles", role); + + Query actorWithInnerRoleQuery = db.query("actors", Actor.class) + .innerJoin(db.query("roles", Role.class), "joinedRoles") + .on("name", EQ, "name"); + + List leftWithInnerNestedItems = db.query("items_with_join", ItemWithJoin.class) + .sort("id", false) + .leftJoin(actorWithInnerRoleQuery, "joinedActors") + .on("actorsIds", SET, "id") + .toList(); + + assertThat(leftWithInnerNestedItems.size(), is(2)); + assertThat(leftWithInnerNestedItems.get(0).joinedActors.size(), is(1)); + assertThat(leftWithInnerNestedItems.get(0).joinedActors.get(0).joinedRoles.size(), is(1)); + assertThat(leftWithInnerNestedItems.get(1).joinedActors.size(), is(0)); + + Query actorWithLeftRoleQuery = db.query("actors", Actor.class) + .leftJoin(db.query("roles", Role.class), "joinedRoles") + .on("name", EQ, "name"); + + List innerWithLeftNestedItems = db.query("items_with_join", ItemWithJoin.class) + .sort("id", false) + .innerJoin(actorWithLeftRoleQuery, "joinedActors") + .on("actorsIds", SET, "id") + .toList(); + + assertThat(innerWithLeftNestedItems.size(), is(2)); + assertThat(innerWithLeftNestedItems.get(0).joinedActors.size(), is(1)); + assertThat(innerWithLeftNestedItems.get(0).joinedActors.get(0).joinedRoles.size(), is(1)); + assertThat(innerWithLeftNestedItems.get(1).joinedActors.size(), is(1)); + assertThat(innerWithLeftNestedItems.get(1).joinedActors.get(0).joinedRoles.size(), is(0)); + } + + @Test + public void testMergeWithNestedJoins() { + assumeQueryFormatV2(); + + db.openNamespace("items_with_join", NamespaceOptions.defaultOptions(), ItemWithJoin.class); + db.openNamespace("actors", NamespaceOptions.defaultOptions(), Actor.class); + db.openNamespace("roles", NamespaceOptions.defaultOptions(), Role.class); + + ItemWithJoin firstItem = new ItemWithJoin(); + firstItem.id = 1; + firstItem.actorsIds = Collections.singletonList(10); + db.upsert("items_with_join", firstItem); + + ItemWithJoin secondItem = new ItemWithJoin(); + secondItem.id = 2; + secondItem.actorsIds = Collections.singletonList(20); + db.upsert("items_with_join", secondItem); + + Actor firstActor = new Actor(); + firstActor.id = 10; + firstActor.name = "first"; + db.upsert("actors", firstActor); + + Actor secondActor = new Actor(); + secondActor.id = 20; + secondActor.name = "second"; + db.upsert("actors", secondActor); + + Role firstRole = new Role(); + firstRole.id = 100; + firstRole.name = "first"; + db.upsert("roles", firstRole); + + Role secondRole = new Role(); + secondRole.id = 200; + secondRole.name = "second"; + db.upsert("roles", secondRole); + + Query firstActorQuery = db.query("actors", Actor.class) + .innerJoin(db.query("roles", Role.class), "joinedRoles") + .on("name", EQ, "name"); + Query secondActorQuery = db.query("actors", Actor.class) + .innerJoin(db.query("roles", Role.class), "joinedRoles") + .on("name", EQ, "name"); + + List items = db.query("items_with_join", ItemWithJoin.class) + .where("id", EQ, 1) + .innerJoin(firstActorQuery, "joinedActors") + .on("actorsIds", SET, "id") + .merge(db.query("items_with_join", ItemWithJoin.class) + .where("id", EQ, 2) + .innerJoin(secondActorQuery, "joinedActors") + .on("actorsIds", SET, "id")) + .toList(); + + Map itemsById = new HashMap<>(); + for (ItemWithJoin item : items) { + itemsById.put(item.id, item); + } + + assertThat(itemsById.size(), is(2)); + assertThat(itemsById.get(1).joinedActors.get(0).joinedRoles.get(0).name, is("first")); + assertThat(itemsById.get(2).joinedActors.get(0).joinedRoles.get(0).name, is("second")); + } + + @Test + public void testNestedJoinWithDepthTwo() { + assumeQueryFormatV2(); + + db.openNamespace("items_with_join", NamespaceOptions.defaultOptions(), ItemWithJoin.class); + db.openNamespace("actors", NamespaceOptions.defaultOptions(), Actor.class); + db.openNamespace("roles", NamespaceOptions.defaultOptions(), Role.class); + db.openNamespace("permissions", NamespaceOptions.defaultOptions(), Permission.class); + + ItemWithJoin item = new ItemWithJoin(); + item.id = 1; + item.actorsIds = Collections.singletonList(10); + db.upsert("items_with_join", item); + + Actor actor = new Actor(); + actor.id = 10; + actor.name = "actor"; + db.upsert("actors", actor); + + Role role = new Role(); + role.id = 100; + role.name = "actor"; + db.upsert("roles", role); + + Permission permission = new Permission(); + permission.id = 1000; + permission.name = "actor"; + db.upsert("permissions", permission); + + Query roleQuery = db.query("roles", Role.class) + .innerJoin(db.query("permissions", Permission.class), "joinedPermissions") + .on("name", EQ, "name"); + Query actorQuery = db.query("actors", Actor.class) + .innerJoin(roleQuery, "joinedRoles") + .on("name", EQ, "name"); + + ItemWithJoin result = db.query("items_with_join", ItemWithJoin.class) + .where("id", EQ, 1) + .innerJoin(actorQuery, "joinedActors") + .on("actorsIds", SET, "id") + .getOne(); + + Actor joinedActor = result.joinedActors.get(0); + Role joinedRole = joinedActor.joinedRoles.get(0); + assertThat(joinedRole.id, is(100)); + assertThat(joinedRole.joinedPermissions.size(), is(1)); + assertThat(joinedRole.joinedPermissions.get(0).id, is(1000)); + } + + @Test + public void testCannotQuerySelectJoinWithMerge() { + assumeQueryFormatV2(); + + db.openNamespace("items_with_join", NamespaceOptions.defaultOptions(), ItemWithJoin.class); + db.openNamespace("actors", NamespaceOptions.defaultOptions(), Actor.class); + + ItemWithJoin item = new ItemWithJoin(); + item.id = 1; + item.actorsIds = Collections.singletonList(10); + db.upsert("items_with_join", item); + + Actor actor = new Actor(); + actor.id = 10; + actor.name = "actor"; + db.upsert("actors", actor); + + Query actorQuery = db.query("actors", Actor.class) + .merge(db.query("actors", Actor.class)); + + RuntimeException exception = assertThrows(RuntimeException.class, + () -> readAll(db.query("items_with_join", ItemWithJoin.class) + .innerJoin(actorQuery, "joinedActors") + .on("actorsIds", SET, "id") + .execute())); + + assertThat(exception.getMessage(), containsString("MERGEs nested into the JOINs are not supported")); + } + + @Test + public void testCannotQuerySelectSubqueryWithJoin() { + assumeQueryFormatV2(); + + db.openNamespace("items_with_join", NamespaceOptions.defaultOptions(), ItemWithJoin.class); + db.openNamespace("actors", NamespaceOptions.defaultOptions(), Actor.class); + db.openNamespace("roles", NamespaceOptions.defaultOptions(), Role.class); + + ItemWithJoin item = new ItemWithJoin(); + item.id = 1; + db.upsert("items_with_join", item); + + Actor actor = new Actor(); + actor.id = 1; + actor.name = "actor"; + db.upsert("actors", actor); + + Query subQuery = db.query("actors", Actor.class) + .select("id") + .innerJoin(db.query("roles", Role.class), "joinedRoles") + .on("name", EQ, "name"); + + RuntimeException exception = assertThrows(RuntimeException.class, + () -> readAll(db.query("items_with_join", ItemWithJoin.class) + .where("id", SET, subQuery) + .execute())); + + assertThat(exception.getMessage(), containsString("Join cannot be in subquery")); + } + + private void assumeQueryFormatV2() { + assumeTrue(db.getBinding().queryFormatVersion() == Consts.QUERY_FORMAT_V2 + && db.getBinding().supportsNestedJoinQueries(), + "Nested joins require QueryFormatV2 and server-side nested join support"); + } + + private void readAll(ResultIterator iterator) { + try (ResultIterator closeableIterator = iterator) { + while (closeableIterator.hasNext()) { + closeableIterator.next(); + } + } + } + @Test public void testJoinInWherePartWithBrackets() { db.openNamespace("items_with_join", NamespaceOptions.defaultOptions(), ItemWithJoin.class); @@ -1038,6 +1342,32 @@ public static class Actor { @Reindex(name = "is_visible") private boolean visible; + + @Transient + private List joinedRoles; + } + + @Setter + @Getter + public static class Role { + @Reindex(name = "id", isPrimaryKey = true) + private Integer id; + + @Reindex(name = "name") + private String name; + + @Transient + private List joinedPermissions; + } + + @Setter + @Getter + public static class Permission { + @Reindex(name = "id", isPrimaryKey = true) + private Integer id; + + @Reindex(name = "name") + private String name; } @Setter