diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/IgniteCacheOffheapManager.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/IgniteCacheOffheapManager.java index 3247e718ad6cf..94eb1155e77cc 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/IgniteCacheOffheapManager.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/IgniteCacheOffheapManager.java @@ -22,6 +22,7 @@ import java.util.Map; import javax.cache.Cache; import org.apache.ignite.IgniteCheckedException; +import org.apache.ignite.internal.binary.BinaryWriterEx; import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion; import org.apache.ignite.internal.processors.cache.distributed.dht.preloader.IgniteDhtDemandedPartitionsMap; import org.apache.ignite.internal.processors.cache.distributed.dht.topology.GridDhtLocalPartition; @@ -123,6 +124,14 @@ public interface IgniteCacheOffheapManager { */ @Nullable public CacheDataRow read(GridCacheContext cctx, KeyCacheObject key) throws IgniteCheckedException; + /** + * @param cctx Cache context. + * @param key Key. + * @return {@code True} if row was found and written, {@code false} otherwise. + * @throws IgniteCheckedException If failed. + */ + public boolean readTo(GridCacheContext cctx, KeyCacheObject key, BinaryWriterEx writer) throws IgniteCheckedException; + /** * @param p Partition. * @return Data store. @@ -601,6 +610,15 @@ void update( */ public CacheDataRow find(GridCacheContext cctx, KeyCacheObject key) throws IgniteCheckedException; + /** + * @param cctx Cache context. + * @param key Key. + * @param writer Writer to copy data to. + * @return {@code True} if key was found, {@code false} otherwise. + * @throws IgniteCheckedException If failed. + */ + public boolean findTo(GridCacheContext cctx, KeyCacheObject key, BinaryWriterEx writer) throws IgniteCheckedException; + /** * @return Data cursor. * @throws IgniteCheckedException If failed. diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/IgniteCacheOffheapManagerImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/IgniteCacheOffheapManagerImpl.java index c97aefd68213f..9703532bd89b7 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/IgniteCacheOffheapManagerImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/IgniteCacheOffheapManagerImpl.java @@ -43,6 +43,7 @@ import org.apache.ignite.internal.GridKernalContext; import org.apache.ignite.internal.GridKernalState; import org.apache.ignite.internal.NodeStoppingException; +import org.apache.ignite.internal.binary.BinaryWriterEx; import org.apache.ignite.internal.pagemem.FullPageId; import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion; import org.apache.ignite.internal.processors.cache.distributed.dht.preloader.CachePartitionPartialCountersMap; @@ -449,6 +450,16 @@ private Iterator cacheData(boolean primary, boolean backup, Affi return row; } + /** {@inheritDoc} */ + @Override public boolean readTo(GridCacheContext cctx, KeyCacheObject key, BinaryWriterEx writer) throws IgniteCheckedException { + CacheDataStore dataStore = dataStore(cctx, key); + + if (dataStore == null) + return false; + + return dataStore.findTo(cctx, key, writer); + } + /** {@inheritDoc} */ @Override public boolean containsKey(GridCacheMapEntry entry) { try { @@ -1807,6 +1818,15 @@ private void clearPendingEntries(GridCacheContext cctx, CacheDataRow oldRow) return row; } + /** {@inheritDoc} */ + @Override public boolean findTo(GridCacheContext cctx, KeyCacheObject key, BinaryWriterEx writer) throws IgniteCheckedException { + int pos = writer.out().position(); + + find(cctx, key); + + return pos != writer.out().position(); + } + /** {@inheritDoc} */ @Override public GridCursor cursor() throws IgniteCheckedException { return dataTree.find(null, null); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridPartitionedSingleGetFuture.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridPartitionedSingleGetFuture.java index 9e36e9691f00f..92d8b44f44aa1 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridPartitionedSingleGetFuture.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridPartitionedSingleGetFuture.java @@ -32,6 +32,7 @@ import org.apache.ignite.internal.IgniteDiagnosticPrepareContext; import org.apache.ignite.internal.IgniteInternalFuture; import org.apache.ignite.internal.NodeStoppingException; +import org.apache.ignite.internal.binary.BinaryWriterEx; import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException; import org.apache.ignite.internal.cluster.ClusterTopologyServerNotFoundException; import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion; @@ -55,6 +56,8 @@ import org.apache.ignite.internal.processors.cache.distributed.near.GridNearSingleGetResponse; import org.apache.ignite.internal.processors.cache.persistence.CacheDataRow; import org.apache.ignite.internal.processors.cache.version.GridCacheVersion; +import org.apache.ignite.internal.processors.platform.client.cache.ClientDirectCacheGetRequest; +import org.apache.ignite.internal.thread.context.OperationContext; import org.apache.ignite.internal.util.lang.GridPlainRunnable; import org.apache.ignite.internal.util.tostring.GridToStringExclude; import org.apache.ignite.internal.util.tostring.GridToStringInclude; @@ -457,28 +460,36 @@ private boolean localGet(AffinityTopologyVersion topVer, KeyCacheObject key, int if (readNoEntry) { KeyCacheObject key0 = (KeyCacheObject)cctx.cacheObjects().prepareForCache(key, cctx); - CacheDataRow row = cctx.offheap().read(cctx, key0); + BinaryWriterEx directWriter = OperationContext.get(ClientDirectCacheGetRequest.DIRECT_WRITER); - if (row != null) { - long expireTime = row.expireTime(); + if (directWriter != null) { + cctx.offheap().readTo(cctx, key0, directWriter); - if (expireTime == 0 || expireTime > U.currentTimeMillis()) { - v = row.value(); + //setResult(found ? FOUND : NOT_FOUND); + } + else { + CacheDataRow row = cctx.offheap().read(cctx, key0); - if (needVer) - ver = row.version(); + if (row != null) { + long expireTime = row.expireTime(); - if (evt) { - cctx.events().readEvent(key, - null, - txLbl, - row.value(), - taskName, - !deserializeBinary); - } + if (expireTime == 0 || expireTime > U.currentTimeMillis()) { + v = row.value(); + + if (needVer) + ver = row.version(); + + if (evt) { + cctx.events().readEvent(key, + null, + txLbl, + row.value(), + taskName, + !deserializeBinary); + } + } else + skipEntry = false; } - else - skipEntry = false; } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/CacheDataRowAdapter.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/CacheDataRowAdapter.java index 810423838a3a8..fd0302c84f480 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/CacheDataRowAdapter.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/CacheDataRowAdapter.java @@ -19,6 +19,7 @@ import java.nio.ByteBuffer; import org.apache.ignite.IgniteCheckedException; +import org.apache.ignite.internal.binary.BinaryWriterEx; import org.apache.ignite.internal.metric.IoStatisticsHolder; import org.apache.ignite.internal.metric.IoStatisticsHolderNoOp; import org.apache.ignite.internal.pagemem.PageIdUtils; @@ -39,6 +40,7 @@ import org.apache.ignite.internal.processors.cache.persistence.tree.io.PageIO; import org.apache.ignite.internal.processors.cache.version.GridCacheVersion; import org.apache.ignite.internal.processors.cacheobject.IgniteCacheObjectProcessor; +import org.apache.ignite.internal.thread.context.OperationContext; import org.apache.ignite.internal.util.GridLongList; import org.apache.ignite.internal.util.GridUnsafe; import org.apache.ignite.internal.util.lang.IgniteThrowableFunction; @@ -52,6 +54,7 @@ import static org.apache.ignite.internal.pagemem.PageIdUtils.pageId; import static org.apache.ignite.internal.processors.cache.persistence.CacheDataRowAdapter.RowData.KEY_ONLY; import static org.apache.ignite.internal.processors.cache.persistence.tree.io.PageIO.T_DATA; +import static org.apache.ignite.internal.processors.platform.client.cache.ClientDirectCacheGetRequest.DIRECT_WRITER; import static org.apache.ignite.internal.util.GridUnsafe.wrapPointer; /** @@ -544,10 +547,27 @@ protected void readFullRow( byte type = PageUtils.getByte(addr, off); off++; - byte[] bytes = PageUtils.getBytes(addr, off, len); - off += len; + BinaryWriterEx writer = OperationContext.get(DIRECT_WRITER); + + if (writer != null && writer.out().hasArray()) { + byte[] data = writer.out().array(); + + writer.out().writeByte((byte)12); + writer.out().writeInt(len); + + int pos = writer.out().position(); + + GridUnsafe.copyMemory(null, addr + off, data, GridUnsafe.BYTE_ARR_OFF + pos, len); + + writer.out().position(pos + len); + } + else { + byte[] bytes = PageUtils.getBytes(addr, off, len); - val = sharedCtx.kernalContext().cacheObjects().toCacheObject(coctx, type, bytes); + val = sharedCtx.kernalContext().cacheObjects().toCacheObject(coctx, type, bytes); + } + + off += len; int verLen; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/GridCacheOffheapManager.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/GridCacheOffheapManager.java index 28911df7c78e1..092da1fc4e4bd 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/GridCacheOffheapManager.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/persistence/GridCacheOffheapManager.java @@ -39,6 +39,7 @@ import org.apache.ignite.IgniteSystemProperties; import org.apache.ignite.SystemProperty; import org.apache.ignite.failure.FailureContext; +import org.apache.ignite.internal.binary.BinaryWriterEx; import org.apache.ignite.internal.managers.encryption.GridEncryptionManager; import org.apache.ignite.internal.managers.encryption.ReencryptStateUtils; import org.apache.ignite.internal.pagemem.FullPageId; @@ -2533,6 +2534,16 @@ private void checkGapsLinkAndPartMetaStorage(PagePartitionMetaIOV3 io, long page return null; } + /** {@inheritDoc} */ + @Override public boolean findTo(GridCacheContext cctx, KeyCacheObject key, BinaryWriterEx writer) throws IgniteCheckedException { + CacheDataStore delegate = init0(true); + + if (delegate != null) + return delegate.findTo(cctx, key, writer); + + return false; + } + /** {@inheritDoc} */ @Override public GridCursor cursor() throws IgniteCheckedException { CacheDataStore delegate = init0(true); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/ClientMessageParser.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/ClientMessageParser.java index d0f54ef2136ba..f7c61c3c73d43 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/ClientMessageParser.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/ClientMessageParser.java @@ -76,6 +76,7 @@ import org.apache.ignite.internal.processors.platform.client.cache.ClientCacheScanQueryRequest; import org.apache.ignite.internal.processors.platform.client.cache.ClientCacheSqlFieldsQueryRequest; import org.apache.ignite.internal.processors.platform.client.cache.ClientCacheSqlQueryRequest; +import org.apache.ignite.internal.processors.platform.client.cache.ClientDirectCacheGetRequest; import org.apache.ignite.internal.processors.platform.client.cluster.ClientClusterChangeStateRequest; import org.apache.ignite.internal.processors.platform.client.cluster.ClientClusterGetDataCenterNodesRequest; import org.apache.ignite.internal.processors.platform.client.cluster.ClientClusterGetStateRequest; @@ -107,6 +108,7 @@ import org.apache.ignite.internal.processors.platform.client.datastructures.ClientIgniteSetValueRemoveAllRequest; import org.apache.ignite.internal.processors.platform.client.datastructures.ClientIgniteSetValueRemoveRequest; import org.apache.ignite.internal.processors.platform.client.datastructures.ClientIgniteSetValueRetainAllRequest; +import org.apache.ignite.internal.processors.platform.client.direct.ClientListenerDirectResponse; import org.apache.ignite.internal.processors.platform.client.service.ClientServiceGetDescriptorRequest; import org.apache.ignite.internal.processors.platform.client.service.ClientServiceGetDescriptorsRequest; import org.apache.ignite.internal.processors.platform.client.service.ClientServiceInvokeRequest; @@ -120,6 +122,9 @@ * Thin client message parser. */ public class ClientMessageParser implements ClientListenerMessageParser { + /** */ + public static boolean USE_DIRECT_READ; + /* General-purpose operations. */ /** */ private static final short OP_RESOURCE_CLOSE = 0; @@ -465,7 +470,10 @@ public ClientListenerRequest decode(BinaryReaderEx reader) { switch (opCode) { case OP_CACHE_GET: - return new ClientCacheGetRequest(reader); + if (USE_DIRECT_READ) + return new ClientDirectCacheGetRequest(reader); + else + return new ClientCacheGetRequest(reader); case OP_BINARY_TYPE_NAME_GET: return new ClientBinaryTypeNameGetRequest(reader); @@ -747,6 +755,9 @@ public ClientListenerRequest decode(BinaryReaderEx reader) { @Override public ClientMessage encode(ClientListenerResponse resp) { assert resp != null; + if (resp instanceof ClientListenerDirectResponse) + return new ClientMessage(((ClientListenerDirectResponse)resp).out()); + BinaryOutputStream outStream = BinaryStreams.createPooledOutputStream(32, false); BinaryWriterEx writer = marsh.writer(outStream); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/ClientRequestHandler.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/ClientRequestHandler.java index 85298183728e5..2b59eca6ff787 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/ClientRequestHandler.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/ClientRequestHandler.java @@ -24,6 +24,9 @@ import org.apache.ignite.internal.IgniteFutureTimeoutCheckedException; import org.apache.ignite.internal.IgniteInternalFuture; import org.apache.ignite.internal.binary.BinaryWriterEx; +import org.apache.ignite.internal.binary.GridBinaryMarshaller; +import org.apache.ignite.internal.binary.streams.BinaryStreams; +import org.apache.ignite.internal.processors.cache.binary.CacheObjectBinaryProcessorImpl; import org.apache.ignite.internal.processors.cache.query.IgniteQueryErrorCode; import org.apache.ignite.internal.processors.odbc.ClientAsyncResponse; import org.apache.ignite.internal.processors.odbc.ClientListenerProtocolVersion; @@ -34,6 +37,7 @@ import org.apache.ignite.internal.processors.platform.client.cache.ClientCacheQueryNextPageRequest; import org.apache.ignite.internal.processors.platform.client.cache.ClientCacheSqlFieldsQueryRequest; import org.apache.ignite.internal.processors.platform.client.cache.ClientCacheSqlQueryRequest; +import org.apache.ignite.internal.processors.platform.client.direct.ClientDirectWriteRequest; import org.apache.ignite.internal.processors.platform.client.tx.ClientTxAwareRequest; import org.apache.ignite.internal.processors.platform.client.tx.ClientTxContext; import org.apache.ignite.internal.util.typedef.X; @@ -63,6 +67,9 @@ public class ClientRequestHandler implements ClientListenerRequestHandler { /** Protocol context. */ private final ClientProtocolContext protocolCtx; + /** Marshaller. */ + private final GridBinaryMarshaller marsh; + /** Logger. */ private final IgniteLogger log; @@ -77,6 +84,7 @@ public class ClientRequestHandler implements ClientListenerRequestHandler { this.ctx = ctx; this.protocolCtx = protocolCtx; + this.marsh = ((CacheObjectBinaryProcessorImpl)ctx.kernalContext().cacheObjects()).marshaller(); log = ctx.kernalContext().log(getClass()); } @@ -124,7 +132,13 @@ private ClientListenerResponse handle0(ClientListenerRequest req) { ClientRequest req0 = (ClientRequest)req; if (req0.isAsync(ctx)) { - IgniteInternalFuture fut = req0.processAsync(ctx); + IgniteInternalFuture fut; + + if (req0 instanceof ClientDirectWriteRequest) { + fut = ((ClientDirectWriteRequest)req).processAsync(ctx, marsh.writer(BinaryStreams.createPooledOutputStream(32, false))); + } + else + fut = req0.processAsync(ctx); if (asyncReqWaitTimeout <= 0) return new ClientAsyncResponse(req0.requestId(), fut); @@ -141,6 +155,9 @@ private ClientListenerResponse handle0(ClientListenerRequest req) { throw new IgniteClientException(ClientStatus.FAILED, e.getMessage(), e); } } + else if (req0 instanceof ClientDirectWriteRequest) { + return ((ClientDirectWriteRequest)req).process(ctx, marsh.writer(BinaryStreams.createPooledOutputStream(32, false))); + } else return req0.process(ctx); } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/ClientResponse.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/ClientResponse.java index 512a90ad091fd..530e7c4bbfc41 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/ClientResponse.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/ClientResponse.java @@ -71,8 +71,28 @@ public ClientResponse(long reqId, int status, String err) { * @param writer Writer. * @param affinityVer Affinity version. */ - public void encode(ClientConnectionContext ctx, BinaryWriterEx writer, - ClientAffinityTopologyVersion affinityVer) { + public void encode(ClientConnectionContext ctx, BinaryWriterEx writer, ClientAffinityTopologyVersion affinityVer) { + encodeHeader(ctx, writer, reqId, status(), error(), affinityVer); + } + + /** + * Encodes the response data. + * @param ctx Connection context. + * @param writer Writer. + */ + @Override public void encode(ClientConnectionContext ctx, BinaryWriterEx writer) { + encode(ctx, writer, ctx.checkAffinityTopologyVersion()); + } + + /** */ + public static void encodeHeader( + ClientConnectionContext ctx, + BinaryWriterEx writer, + long reqId, + int status, + String err, + ClientAffinityTopologyVersion affinityVer + ) { writer.writeLong(reqId); ClientProtocolContext protocolCtx = ctx.currentProtocolContext(); @@ -80,7 +100,7 @@ public void encode(ClientConnectionContext ctx, BinaryWriterEx writer, assert protocolCtx != null; if (protocolCtx.isFeatureSupported(PARTITION_AWARENESS)) { - boolean error = status() != ClientStatus.SUCCESS; + boolean error = status != ClientStatus.SUCCESS; short flags = ClientFlag.makeFlags(error, affinityVer.isChanged()); @@ -94,22 +114,13 @@ public void encode(ClientConnectionContext ctx, BinaryWriterEx writer, return; } - writer.writeInt(status()); + writer.writeInt(status); - if (status() != ClientStatus.SUCCESS) { - writer.writeString(error()); + if (status != ClientStatus.SUCCESS) { + writer.writeString(err); } } - /** - * Encodes the response data. - * @param ctx Connection context. - * @param writer Writer. - */ - @Override public void encode(ClientConnectionContext ctx, BinaryWriterEx writer) { - encode(ctx, writer, ctx.checkAffinityTopologyVersion()); - } - /** * Gets the request id. * diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/cache/ClientDirectCacheGetRequest.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/cache/ClientDirectCacheGetRequest.java new file mode 100644 index 0000000000000..f35f6c4932832 --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/cache/ClientDirectCacheGetRequest.java @@ -0,0 +1,88 @@ +/* + * 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.ignite.internal.processors.platform.client.cache; + +import org.apache.ignite.internal.IgniteInternalFuture; +import org.apache.ignite.internal.binary.BinaryReaderEx; +import org.apache.ignite.internal.binary.BinaryWriterEx; +import org.apache.ignite.internal.processors.cache.GridCacheEntryInfo; +import org.apache.ignite.internal.processors.platform.client.ClientConnectionContext; +import org.apache.ignite.internal.processors.platform.client.ClientObjectResponse; +import org.apache.ignite.internal.processors.platform.client.ClientResponse; +import org.apache.ignite.internal.processors.platform.client.ClientStatus; +import org.apache.ignite.internal.processors.platform.client.direct.ClientDirectWriteRequest; +import org.apache.ignite.internal.processors.platform.client.direct.ClientListenerDirectResponse; +import org.apache.ignite.internal.thread.context.OperationContext; +import org.apache.ignite.internal.thread.context.OperationContextAttribute; +import org.apache.ignite.internal.thread.context.Scope; + +/** + * Cache get request. + */ +public class ClientDirectCacheGetRequest extends ClientCacheKeyRequest implements ClientDirectWriteRequest { + /** */ + public static OperationContextAttribute DIRECT_WRITER = OperationContextAttribute.newInstance(); + + /** */ + public static GridCacheEntryInfo NOT_FOUND = new GridCacheEntryInfo(); + + /** */ + public static GridCacheEntryInfo FOUND = new GridCacheEntryInfo(); + + /** + * Constructor. + * + * @param reader Reader. + */ + public ClientDirectCacheGetRequest(BinaryReaderEx reader) { + super(reader); + } + + /** {@inheritDoc} */ + @Override public ClientResponse process(ClientConnectionContext ctx, BinaryWriterEx writer) { + // TODO: disable fast path in case repair read enabled. + // TODO: check when must disable fast path. + try (Scope ignored = OperationContext.set(DIRECT_WRITER, writer)) { + ClientResponse.encodeHeader(ctx, writer, requestId(), ClientStatus.SUCCESS, null, ctx.checkAffinityTopologyVersion()); + + int pos = writer.out().position(); + + Object val = cache(ctx).get(key()); + + if (pos != writer.out().position()) + return new ClientListenerDirectResponse(requestId(), writer.out()); + + return new ClientObjectResponse(requestId(), val); + } + } + + /** {@inheritDoc} */ + @Override public ClientResponse process0(ClientConnectionContext ctx) { + throw new UnsupportedOperationException("Unsupported!"); + } + + /** {@inheritDoc} */ + @Override protected IgniteInternalFuture processAsync0(ClientConnectionContext ctx) { + throw new UnsupportedOperationException("Unsupported!"); + } + + /** {@inheritDoc} */ + @Override public IgniteInternalFuture processAsync(ClientConnectionContext ctx, BinaryWriterEx writer) { + throw new UnsupportedOperationException("Unsupported!"); + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/direct/ClientDirectWriteRequest.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/direct/ClientDirectWriteRequest.java new file mode 100644 index 0000000000000..d8065a849e7fe --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/direct/ClientDirectWriteRequest.java @@ -0,0 +1,38 @@ +/* + * 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.ignite.internal.processors.platform.client.direct; + + +import org.apache.ignite.binary.BinaryWriter; +import org.apache.ignite.internal.IgniteInternalFuture; +import org.apache.ignite.internal.binary.BinaryWriterEx; +import org.apache.ignite.internal.processors.platform.client.ClientConnectionContext; +import org.apache.ignite.internal.processors.platform.client.ClientResponse; + +/** Interface for requests that writes response directly to {@link BinaryWriter} passed as an argument. */ +public interface ClientDirectWriteRequest { + /** + * Processes the request. + * Reponse writt + */ + public ClientResponse process(ClientConnectionContext ctx, BinaryWriterEx out); + + /** + * Processes the request asynchronously. + */ + public IgniteInternalFuture processAsync(ClientConnectionContext ctx, BinaryWriterEx out); +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/direct/ClientListenerDirectResponse.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/direct/ClientListenerDirectResponse.java new file mode 100644 index 0000000000000..3e537e35d9b9d --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/platform/client/direct/ClientListenerDirectResponse.java @@ -0,0 +1,39 @@ +/* + * 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.ignite.internal.processors.platform.client.direct; + + +import org.apache.ignite.internal.binary.streams.BinaryOutputStream; +import org.apache.ignite.internal.processors.platform.client.ClientResponse; + +/** */ +public class ClientListenerDirectResponse extends ClientResponse { + /** */ + private final BinaryOutputStream out; + + /** */ + public ClientListenerDirectResponse(long reqId, BinaryOutputStream out) { + super(reqId); + + this.out = out; + } + + /** */ + public BinaryOutputStream out() { + return out; + } +} diff --git a/modules/indexing/src/test/java/org/apache/ignite/cache/query/ThinClientSimpleTest.java b/modules/indexing/src/test/java/org/apache/ignite/cache/query/ThinClientSimpleTest.java new file mode 100644 index 0000000000000..96f3b60089369 --- /dev/null +++ b/modules/indexing/src/test/java/org/apache/ignite/cache/query/ThinClientSimpleTest.java @@ -0,0 +1,77 @@ +/* + * 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.ignite.cache.query; + +import java.time.Duration; +import java.time.temporal.ChronoUnit; +import java.util.Arrays; +import java.util.concurrent.ThreadLocalRandom; +import java.util.stream.IntStream; +import org.apache.ignite.IgniteCache; +import org.apache.ignite.Ignition; +import org.apache.ignite.client.ClientCache; +import org.apache.ignite.client.IgniteClient; +import org.apache.ignite.configuration.ClientConfiguration; +import org.apache.ignite.internal.IgniteEx; +import org.apache.ignite.internal.processors.platform.client.ClientMessageParser; +import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest; +import org.junit.Test; + +import static org.apache.ignite.client.Config.SERVER; + +/** */ +public class ThinClientSimpleTest extends GridCommonAbstractTest { + /** */ + @Test + public void testThinClientPerf() throws Exception { + try (IgniteEx srv = startGrid()) { + IgniteCache c = srv.getOrCreateCache(DEFAULT_CACHE_NAME); + + ThreadLocalRandom r = ThreadLocalRandom.current(); + + byte[] val = new byte[100]; + + r.nextBytes(val); + + IntStream.range(0, 1000).forEach(i -> c.put(i, val)); + + try (IgniteClient cln = Ignition.startClient(new ClientConfiguration().setAddresses(SERVER))) { + ClientCache cc = cln.cache(DEFAULT_CACHE_NAME); + + // Warmup + for (int i = 0; i < 100_000; i++) + assertNotNull(cc.get(r.nextInt(1000))); + + for (boolean direct: new boolean[] {true, false}) { + ClientMessageParser.USE_DIRECT_READ = direct; + + long start = System.nanoTime(); + + for (int i = 0; i < 100_000; i++) + assertTrue(Arrays.equals(val, cc.get(r.nextInt(1000)))); + + long finish = System.nanoTime(); + + Duration t = Duration.of(finish - start, ChronoUnit.NANOS); + + System.out.println("time (millis) [direct=" + direct + "] = " + t.toMillis()); + } + } + } + } +}