This is an automated email from the ASF dual-hosted git repository.
anton-vinogradov pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ignite.git
The following commit(s) were added to refs/heads/master by this push:
new 08a835b1d9f IGNITE-26630 Separate the stripe key from the cache
partition (#13437)
08a835b1d9f is described below
commit 08a835b1d9ff34cfd1554ac5d7590c251c52418f
Author: Anton Vinogradov <[email protected]>
AuthorDate: Thu Aug 6 21:06:29 2026 +0300
IGNITE-26630 Separate the stripe key from the cache partition (#13437)
---
.../org/apache/ignite/internal/StripedMessage.java | 37 ++++++++++++++++++++++
.../managers/communication/GridIoManager.java | 14 ++++----
.../managers/communication/GridIoMessage.java | 23 +++-----------
.../processors/cache/GridCacheIoManager.java | 24 +++++++-------
.../processors/cache/GridCacheMessage.java | 11 +++----
.../distributed/GridDistributedLockRequest.java | 4 +--
.../GridDistributedTxFinishResponse.java | 2 +-
.../GridDistributedTxPrepareResponse.java | 2 +-
.../cache/distributed/GridNearUnlockRequest.java | 4 +--
.../distributed/dht/GridDhtTxFinishRequest.java | 2 +-
.../distributed/dht/GridDhtTxPrepareRequest.java | 2 +-
.../distributed/dht/atomic/GridDhtAtomicCache.java | 14 ++++----
.../dht/atomic/GridDhtAtomicNearResponse.java | 2 +-
.../atomic/GridDhtAtomicSingleUpdateRequest.java | 2 +-
.../dht/atomic/GridDhtAtomicUpdateRequest.java | 2 +-
.../dht/atomic/GridDhtAtomicUpdateResponse.java | 2 +-
.../atomic/GridNearAtomicAbstractUpdateFuture.java | 6 ++--
.../atomic/GridNearAtomicCheckUpdateRequest.java | 10 ++++--
.../atomic/GridNearAtomicFullUpdateRequest.java | 2 +-
.../atomic/GridNearAtomicSingleUpdateRequest.java | 2 +-
.../dht/atomic/GridNearAtomicUpdateResponse.java | 9 +-----
.../GridDhtPartitionsAbstractMessage.java | 5 ++-
.../cache/distributed/near/GridNearGetRequest.java | 4 +--
.../distributed/near/GridNearSingleGetRequest.java | 2 +-
.../distributed/near/GridNearTxFinishRequest.java | 2 +-
.../distributed/near/GridNearTxPrepareRequest.java | 2 +-
.../cache/query/GridCacheQueryRequest.java | 9 ++++--
.../cache/transactions/IgniteTxHandler.java | 14 ++++----
.../processors/datastreamer/DataStreamerImpl.java | 5 ++-
.../datastreamer/DataStreamerRequest.java | 7 ++--
30 files changed, 126 insertions(+), 100 deletions(-)
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/StripedMessage.java
b/modules/core/src/main/java/org/apache/ignite/internal/StripedMessage.java
new file mode 100644
index 00000000000..acd130d999a
--- /dev/null
+++ b/modules/core/src/main/java/org/apache/ignite/internal/StripedMessage.java
@@ -0,0 +1,37 @@
+/*
+ * 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;
+
+import org.apache.ignite.plugin.extensions.communication.Message;
+
+/** Message that chooses the stripe it is processed in. */
+public interface StripedMessage extends Message {
+ /** Process in any stripe. */
+ public static final int ANY_STRIPE = -1;
+
+ /** Process in the pool itself, not in a stripe. */
+ public static final int NO_STRIPE = Integer.MIN_VALUE;
+
+ /**
+ * The value is an index, not a cache partition: messages sharing it are
processed one after another
+ * by the same thread.
+ *
+ * @return Stripe index, {@link #ANY_STRIPE} or {@link #NO_STRIPE}.
+ */
+ public int stripeIdx();
+}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoManager.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoManager.java
index bba3c4c08bd..320f387a1c6 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoManager.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoManager.java
@@ -142,6 +142,8 @@ import static
org.apache.ignite.events.EventType.EVT_NODE_LEFT;
import static org.apache.ignite.internal.GridTopic.TOPIC_COMM_SYSTEM;
import static org.apache.ignite.internal.GridTopic.TOPIC_COMM_USER;
import static org.apache.ignite.internal.GridTopic.TOPIC_IO_TEST;
+import static org.apache.ignite.internal.StripedMessage.ANY_STRIPE;
+import static org.apache.ignite.internal.StripedMessage.NO_STRIPE;
import static
org.apache.ignite.internal.managers.communication.GridIoPolicy.AFFINITY_POOL;
import static
org.apache.ignite.internal.managers.communication.GridIoPolicy.CALLER_THREAD;
import static
org.apache.ignite.internal.managers.communication.GridIoPolicy.DATA_STREAMER_POOL;
@@ -1384,21 +1386,21 @@ public class GridIoManager extends
GridManagerAdapter<CommunicationSpi<Object>>
if (msg0.processFromNioThread())
c.run();
else
- ctx.pools().getStripedExecutorService().execute(-1, c);
+ ctx.pools().getStripedExecutorService().execute(ANY_STRIPE, c);
return;
}
- final int part = msg.partition(); // Store partition to avoid possible
recalculation.
+ final int stripeIdx = msg.stripeIdx(); // Store to avoid possible
recalculation.
- if (plc == GridIoPolicy.SYSTEM_POOL && part !=
GridIoMessage.STRIPE_DISABLED_PART) {
- ctx.pools().getStripedExecutorService().execute(part, c);
+ if (plc == GridIoPolicy.SYSTEM_POOL && stripeIdx != NO_STRIPE) {
+ ctx.pools().getStripedExecutorService().execute(stripeIdx, c);
return;
}
- if (plc == GridIoPolicy.DATA_STREAMER_POOL && part !=
GridIoMessage.STRIPE_DISABLED_PART) {
- ctx.pools().getDataStreamerExecutorService().execute(part, c);
+ if (plc == GridIoPolicy.DATA_STREAMER_POOL && stripeIdx != NO_STRIPE) {
+ ctx.pools().getDataStreamerExecutorService().execute(stripeIdx, c);
return;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java
index c8e0e220aec..eea2e6336c8 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java
@@ -21,8 +21,7 @@ import org.apache.ignite.internal.ExecutorAwareMessage;
import org.apache.ignite.internal.GridTopicMessage;
import org.apache.ignite.internal.NioField;
import org.apache.ignite.internal.Order;
-import org.apache.ignite.internal.processors.cache.GridCacheMessage;
-import org.apache.ignite.internal.processors.datastreamer.DataStreamerRequest;
+import org.apache.ignite.internal.StripedMessage;
import
org.apache.ignite.internal.thread.context.OperationContextSnapshotMessage;
import org.apache.ignite.internal.util.nio.GridNioServer.MessageWrapper;
import org.apache.ignite.internal.util.tostring.GridToStringInclude;
@@ -33,10 +32,7 @@ import org.jetbrains.annotations.Nullable;
/**
* Wrapper for all grid messages.
*/
-public class GridIoMessage implements Message, MessageWrapper {
- /** */
- public static final Integer STRIPE_DISABLED_PART = Integer.MIN_VALUE;
-
+public class GridIoMessage implements StripedMessage, MessageWrapper {
/** Policy. */
@Order(0)
byte plc;
@@ -175,18 +171,9 @@ public class GridIoMessage implements Message,
MessageWrapper {
throw new AssertionError();
}
- /**
- * Get single partition for this message (if applicable).
- *
- * @return Partition ID.
- */
- public int partition() {
- if (msg instanceof GridCacheMessage)
- return ((GridCacheMessage)msg).partition();
- if (msg instanceof DataStreamerRequest)
- return ((DataStreamerRequest)msg).partition();
- else
- return STRIPE_DISABLED_PART;
+ /** {@inheritDoc} */
+ @Override public int stripeIdx() {
+ return msg instanceof StripedMessage ?
((StripedMessage)msg).stripeIdx() : NO_STRIPE;
}
/**
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheIoManager.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheIoManager.java
index 5f950c022dd..225c143b8db 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheIoManager.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheIoManager.java
@@ -445,7 +445,7 @@ public class GridCacheIoManager extends
GridCacheSharedManagerAdapter {
nearEvicted.add(req.nearKey(i));
GridDhtAtomicUpdateResponse dhtRes = new
GridDhtAtomicUpdateResponse(req.cacheId(),
- req.partition(),
+ req.stripeIdx(),
req.futureId());
dhtRes.nearEvicted(nearEvicted);
@@ -458,7 +458,7 @@ public class GridCacheIoManager extends
GridCacheSharedManagerAdapter {
if (req.nearNodeId() != null) {
GridDhtAtomicNearResponse nearRes = new
GridDhtAtomicNearResponse(req.cacheId(),
- req.partition(),
+ req.stripeIdx(),
req.nearFutureId(),
nodeId,
req.flags());
@@ -791,7 +791,7 @@ public class GridCacheIoManager extends
GridCacheSharedManagerAdapter {
GridDhtTxPrepareRequest req = (GridDhtTxPrepareRequest)msg;
GridDhtTxPrepareResponse res = new GridDhtTxPrepareResponse(
- req.partition(),
+ req.stripeIdx(),
req.version(),
req.futureId(),
req.miniId(),
@@ -806,7 +806,7 @@ public class GridCacheIoManager extends
GridCacheSharedManagerAdapter {
GridDhtAtomicUpdateResponse res = new GridDhtAtomicUpdateResponse(
req.cacheId(),
- req.partition(),
+ req.stripeIdx(),
req.futureId());
res.onError(req.classError());
@@ -815,7 +815,7 @@ public class GridCacheIoManager extends
GridCacheSharedManagerAdapter {
if (req.nearNodeId() != null) {
GridDhtAtomicNearResponse nearRes = new
GridDhtAtomicNearResponse(req.cacheId(),
- req.partition(),
+ req.stripeIdx(),
req.nearFutureId(),
nodeId,
req.flags());
@@ -832,7 +832,7 @@ public class GridCacheIoManager extends
GridCacheSharedManagerAdapter {
req.cacheId(),
nodeId,
req.futureId(),
- req.partition(),
+ req.stripeIdx(),
false);
res.error(req.classError());
@@ -901,7 +901,7 @@ public class GridCacheIoManager extends
GridCacheSharedManagerAdapter {
GridNearTxPrepareRequest req = (GridNearTxPrepareRequest)msg;
GridNearTxPrepareResponse res = new GridNearTxPrepareResponse(
- req.partition(),
+ req.stripeIdx(),
req.version(),
req.futureId(),
req.miniId(),
@@ -980,7 +980,7 @@ public class GridCacheIoManager extends
GridCacheSharedManagerAdapter {
req.cacheId(),
nodeId,
req.futureId(),
- req.partition(),
+ req.stripeIdx(),
false);
res.error(req.classError());
@@ -994,7 +994,7 @@ public class GridCacheIoManager extends
GridCacheSharedManagerAdapter {
req.cacheId(),
nodeId,
req.futureId(),
- req.partition(),
+ req.stripeIdx(),
false);
res.error(req.classError());
@@ -1008,7 +1008,7 @@ public class GridCacheIoManager extends
GridCacheSharedManagerAdapter {
req.cacheId(),
nodeId,
req.futureId(),
- req.partition(),
+ req.stripeIdx(),
false);
res.error(req.classError());
@@ -1020,7 +1020,7 @@ public class GridCacheIoManager extends
GridCacheSharedManagerAdapter {
GridDhtAtomicUpdateResponse res = new GridDhtAtomicUpdateResponse(
req.cacheId(),
- req.partition(),
+ req.stripeIdx(),
req.futureId());
res.onError(req.classError());
@@ -1029,7 +1029,7 @@ public class GridCacheIoManager extends
GridCacheSharedManagerAdapter {
if (req.nearNodeId() != null) {
GridDhtAtomicNearResponse nearRes = new
GridDhtAtomicNearResponse(req.cacheId(),
- req.partition(),
+ req.stripeIdx(),
req.nearFutureId(),
nodeId,
req.flags());
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheMessage.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheMessage.java
index bbbefa0495d..2f09c39eb9f 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheMessage.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/GridCacheMessage.java
@@ -26,6 +26,7 @@ import org.apache.ignite.IgniteCheckedException;
import org.apache.ignite.IgniteLogger;
import org.apache.ignite.internal.DeferredUnmarshalMessage;
import org.apache.ignite.internal.Order;
+import org.apache.ignite.internal.StripedMessage;
import org.apache.ignite.internal.managers.deployment.GridDeployment;
import org.apache.ignite.internal.managers.deployment.GridDeploymentInfo;
import org.apache.ignite.internal.managers.deployment.GridDeploymentInfoBean;
@@ -44,7 +45,7 @@ import org.jetbrains.annotations.Nullable;
*
* @see DeployableMessage
*/
-public abstract class GridCacheMessage implements DeferredUnmarshalMessage {
+public abstract class GridCacheMessage implements DeferredUnmarshalMessage,
StripedMessage {
/** Maximum number of cache lookup indexes. */
public static final int MAX_CACHE_MSG_LOOKUP_INDEX = 7;
@@ -123,11 +124,9 @@ public abstract class GridCacheMessage implements
DeferredUnmarshalMessage {
return -1;
}
- /**
- * @return Partition ID this message is targeted to or {@code -1} if it
cannot be determined.
- */
- public int partition() {
- return -1;
+ /** {@inheritDoc} */
+ @Override public int stripeIdx() {
+ return ANY_STRIPE;
}
/**
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridDistributedLockRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridDistributedLockRequest.java
index 2da8b4dd3cc..57feac70717 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridDistributedLockRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridDistributedLockRequest.java
@@ -366,8 +366,8 @@ public class GridDistributedLockRequest extends
GridDistributedBaseMessage {
}
/** {@inheritDoc} */
- @Override public int partition() {
- return keys != null && !keys.isEmpty() ? keys.get(0).partition() : -1;
+ @Override public int stripeIdx() {
+ return keys != null && !keys.isEmpty() ? keys.get(0).partition() :
ANY_STRIPE;
}
/**
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridDistributedTxFinishResponse.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridDistributedTxFinishResponse.java
index acd4baab4f8..cc92f3badee 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridDistributedTxFinishResponse.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridDistributedTxFinishResponse.java
@@ -64,7 +64,7 @@ public class GridDistributedTxFinishResponse extends
GridCacheMessage {
}
/** {@inheritDoc} */
- @Override public final int partition() {
+ @Override public final int stripeIdx() {
return part;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridDistributedTxPrepareResponse.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridDistributedTxPrepareResponse.java
index 95a8e651c32..43b57697934 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridDistributedTxPrepareResponse.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridDistributedTxPrepareResponse.java
@@ -78,7 +78,7 @@ public class GridDistributedTxPrepareResponse extends
GridDistributedBaseMessage
}
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
return part;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridNearUnlockRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridNearUnlockRequest.java
index bf4a795c85b..9770bf9fe3c 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridNearUnlockRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/GridNearUnlockRequest.java
@@ -90,8 +90,8 @@ public class GridNearUnlockRequest extends
GridDistributedBaseMessage {
}
/** {@inheritDoc} */
- @Override public int partition() {
- return keys != null && !keys.isEmpty() ? keys.get(0).partition() : -1;
+ @Override public int stripeIdx() {
+ return keys != null && !keys.isEmpty() ? keys.get(0).partition() :
ANY_STRIPE;
}
/** {@inheritDoc} */
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTxFinishRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTxFinishRequest.java
index 6d321c3eb41..75d5ba62b2e 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTxFinishRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTxFinishRequest.java
@@ -200,7 +200,7 @@ public class GridDhtTxFinishRequest extends
GridDistributedTxFinishRequest {
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
return U.safeAbs(version().hashCode());
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTxPrepareRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTxPrepareRequest.java
index 03b6f2532db..76848da191c 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTxPrepareRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/GridDhtTxPrepareRequest.java
@@ -320,7 +320,7 @@ public class GridDhtTxPrepareRequest extends
GridDistributedTxPrepareRequest imp
}
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
return U.safeAbs(version().hashCode());
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicCache.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicCache.java
index 2b67f27e8b3..496f579a33f 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicCache.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicCache.java
@@ -1766,7 +1766,7 @@ public class GridDhtAtomicCache<K, V> extends
GridDhtCacheAdapter<K, V> {
GridNearAtomicUpdateResponse res = new
GridNearAtomicUpdateResponse(ctx.cacheId(),
nodeId,
req.futureId(),
- req.partition(),
+ req.stripeIdx(),
false);
res.addFailedKeys(req.keys(), e);
@@ -1789,7 +1789,7 @@ public class GridDhtAtomicCache<K, V> extends
GridDhtCacheAdapter<K, V> {
GridNearAtomicUpdateResponse res = new
GridNearAtomicUpdateResponse(ctx.cacheId(),
node.id(),
req.futureId(),
- req.partition(),
+ req.stripeIdx(),
false);
assert !req.returnValue() || (req.operation() == TRANSFORM ||
req.size() == 1);
@@ -3241,7 +3241,7 @@ public class GridDhtAtomicCache<K, V> extends
GridDhtCacheAdapter<K, V> {
GridNearAtomicUpdateResponse res = new
GridNearAtomicUpdateResponse(ctx.cacheId(),
nodeId,
checkReq.futureId(),
- checkReq.partition(),
+ checkReq.stripeIdx(),
false);
GridCacheReturn ret = new GridCacheReturn(false, true);
@@ -3263,7 +3263,7 @@ public class GridDhtAtomicCache<K, V> extends
GridDhtCacheAdapter<K, V> {
", writeVer=" + req.writeVersion() + ", node=" + nodeId + ']');
}
- assert req.partition() >= 0 : req;
+ assert req.stripeIdx() >= 0 : req;
GridCacheVersion ver = req.writeVersion();
@@ -3273,7 +3273,7 @@ public class GridDhtAtomicCache<K, V> extends
GridDhtCacheAdapter<K, V> {
if (req.nearNodeId() != null) {
nearRes = new GridDhtAtomicNearResponse(ctx.cacheId(),
- req.partition(),
+ req.stripeIdx(),
req.nearFutureId(),
nodeId,
req.flags());
@@ -3407,7 +3407,7 @@ public class GridDhtAtomicCache<K, V> extends
GridDhtCacheAdapter<K, V> {
if (nearEvicted != null) {
dhtRes = new GridDhtAtomicUpdateResponse(ctx.cacheId(),
- req.partition(),
+ req.stripeIdx(),
req.futureId());
dhtRes.nearEvicted(nearEvicted);
@@ -3441,7 +3441,7 @@ public class GridDhtAtomicCache<K, V> extends
GridDhtCacheAdapter<K, V> {
if (dhtRes != null)
sendDhtPrimaryResponse(nodeId, req, dhtRes);
else
- sendDeferredUpdateResponse(req.partition(), nodeId,
req.futureId());
+ sendDeferredUpdateResponse(req.stripeIdx(), nodeId,
req.futureId());
}
/**
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicNearResponse.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicNearResponse.java
index f7944ca5228..d0de42905fb 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicNearResponse.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicNearResponse.java
@@ -112,7 +112,7 @@ public class GridDhtAtomicNearResponse extends
GridCacheIdMessage {
}
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
return partId;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicSingleUpdateRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicSingleUpdateRequest.java
index c39c64ca8bb..40af94b3f72 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicSingleUpdateRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicSingleUpdateRequest.java
@@ -210,7 +210,7 @@ public class GridDhtAtomicSingleUpdateRequest extends
GridDhtAtomicAbstractUpdat
}
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
int p = key.partition();
assert p >= 0;
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicUpdateRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicUpdateRequest.java
index e36874727bf..6a361926a69 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicUpdateRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicUpdateRequest.java
@@ -345,7 +345,7 @@ public class GridDhtAtomicUpdateRequest extends
GridDhtAtomicAbstractUpdateReque
}
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
assert !F.isEmpty(keys) || !F.isEmpty(nearKeys);
int p = !keys.isEmpty() ? keys.get(0).partition() :
nearKeys.get(0).partition();
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicUpdateResponse.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicUpdateResponse.java
index 26ae00ce793..3190bf92459 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicUpdateResponse.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridDhtAtomicUpdateResponse.java
@@ -115,7 +115,7 @@ public class GridDhtAtomicUpdateResponse extends
GridCacheIdMessage implements G
}
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
return partId;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicAbstractUpdateFuture.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicAbstractUpdateFuture.java
index d7638c11c8e..c69b0342637 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicAbstractUpdateFuture.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicAbstractUpdateFuture.java
@@ -481,7 +481,7 @@ public abstract class GridNearAtomicAbstractUpdateFuture
extends GridCacheFuture
GridNearAtomicUpdateResponse res = new
GridNearAtomicUpdateResponse(cctx.cacheId(),
req.nodeId(),
req.futureId(),
- req.partition(),
+ req.stripeIdx(),
true);
ClusterTopologyCheckedException e = new
ClusterTopologyCheckedException("Primary node left grid " +
@@ -502,7 +502,7 @@ public abstract class GridNearAtomicAbstractUpdateFuture
extends GridCacheFuture
GridNearAtomicUpdateResponse res = new
GridNearAtomicUpdateResponse(cctx.cacheId(),
req.nodeId(),
req.futureId(),
- req.partition(),
+ req.stripeIdx(),
e instanceof ClusterTopologyCheckedException);
res.addFailedKeys(req.keys(), e);
@@ -518,7 +518,7 @@ public abstract class GridNearAtomicAbstractUpdateFuture
extends GridCacheFuture
GridNearAtomicUpdateResponse res = new
GridNearAtomicUpdateResponse(cctx.cacheId(),
req.updateRequest().nodeId(),
req.futureId(),
- req.partition(),
+ req.stripeIdx(),
e instanceof ClusterTopologyCheckedException);
res.addFailedKeys(req.updateRequest().keys(), e);
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicCheckUpdateRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicCheckUpdateRequest.java
index 2afb5a21c5b..44e26d854fc 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicCheckUpdateRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicCheckUpdateRequest.java
@@ -54,7 +54,7 @@ public class GridNearAtomicCheckUpdateRequest extends
GridCacheIdMessage {
this.updateReq = updateReq;
this.cacheId = updateReq.cacheId();
- this.partId = updateReq.partition();
+ this.partId = updateReq.stripeIdx();
this.futId = updateReq.futureId();
assert partId >= 0;
@@ -74,8 +74,12 @@ public class GridNearAtomicCheckUpdateRequest extends
GridCacheIdMessage {
return updateReq;
}
- /** {@inheritDoc} */
- @Override public int partition() {
+ /**
+ * The value travels because the primary puts it into the response, it
cannot restore it otherwise.
+ *
+ * {@inheritDoc}
+ */
+ @Override public int stripeIdx() {
return partId;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicFullUpdateRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicFullUpdateRequest.java
index cc868b91bb6..ffe021fbe85 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicFullUpdateRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicFullUpdateRequest.java
@@ -335,7 +335,7 @@ public class GridNearAtomicFullUpdateRequest extends
GridNearAtomicAbstractUpdat
}
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
assert !F.isEmpty(keys);
return keys.get(0).partition();
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicSingleUpdateRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicSingleUpdateRequest.java
index 963df4d0c76..c064d24c0c7 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicSingleUpdateRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicSingleUpdateRequest.java
@@ -96,7 +96,7 @@ public class GridNearAtomicSingleUpdateRequest extends
GridNearAtomicAbstractUpd
}
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
assert key != null;
return key.partition();
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicUpdateResponse.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicUpdateResponse.java
index a4b049a7ab0..942ea04f154 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicUpdateResponse.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/atomic/GridNearAtomicUpdateResponse.java
@@ -355,17 +355,10 @@ public class GridNearAtomicUpdateResponse extends
GridCacheIdMessage implements
}
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
return partId;
}
- /**
- * @param partId New partition ID.
- */
- public void partition(int partId) {
- this.partId = partId;
- }
-
/** {@inheritDoc} */
@Override public boolean addDeploymentInfo() {
return addDepInfo;
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsAbstractMessage.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsAbstractMessage.java
index 0af0f2ec1c0..47b6a58dd18 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsAbstractMessage.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsAbstractMessage.java
@@ -18,7 +18,6 @@
package org.apache.ignite.internal.processors.cache.distributed.dht.preloader;
import org.apache.ignite.internal.Order;
-import org.apache.ignite.internal.managers.communication.GridIoMessage;
import org.apache.ignite.internal.processors.cache.GridCacheMessage;
import org.apache.ignite.internal.processors.cache.version.GridCacheVersion;
import org.apache.ignite.internal.util.typedef.internal.S;
@@ -69,8 +68,8 @@ public abstract class GridDhtPartitionsAbstractMessage
extends GridCacheMessage
}
/** {@inheritDoc} */
- @Override public int partition() {
- return GridIoMessage.STRIPE_DISABLED_PART;
+ @Override public int stripeIdx() {
+ return NO_STRIPE;
}
/** {@inheritDoc} */
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearGetRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearGetRequest.java
index 7e4c199dacb..5d20cd53dea 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearGetRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearGetRequest.java
@@ -251,10 +251,10 @@ public class GridNearGetRequest extends
GridCacheIdMessage implements GridCacheD
}
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
Collection<KeyCacheObject> keys0 = keyMap != null ? keyMap.keySet() :
keys;
- return F.isEmpty(keys0) ? -1 : keys0.iterator().next().partition();
+ return F.isEmpty(keys0) ? ANY_STRIPE :
keys0.iterator().next().partition();
}
/**
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearSingleGetRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearSingleGetRequest.java
index ba8496e2b69..f79fd3cb748 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearSingleGetRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearSingleGetRequest.java
@@ -191,7 +191,7 @@ public class GridNearSingleGetRequest extends
GridCacheIdMessage implements Grid
}
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
assert key != null;
return key.partition();
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearTxFinishRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearTxFinishRequest.java
index 35af30081af..bc31cc65b93 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearTxFinishRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearTxFinishRequest.java
@@ -140,7 +140,7 @@ public class GridNearTxFinishRequest extends
GridDistributedTxFinishRequest {
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
return U.safeAbs(version().hashCode());
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearTxPrepareRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearTxPrepareRequest.java
index fc5074f8941..f3e6504fa61 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearTxPrepareRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/near/GridNearTxPrepareRequest.java
@@ -288,7 +288,7 @@ public class GridNearTxPrepareRequest extends
GridDistributedTxPrepareRequest im
/** {@inheritDoc} */
- @Override public int partition() {
+ @Override public int stripeIdx() {
return U.safeAbs(version().hashCode());
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/GridCacheQueryRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/GridCacheQueryRequest.java
index 11ae812ad4a..15f4e3b5411 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/GridCacheQueryRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/GridCacheQueryRequest.java
@@ -580,9 +580,14 @@ public class GridCacheQueryRequest extends
GridCacheIdMessage implements GridCac
}
/**
- * @return Partition.
+ * @return Partition to scan, {@code -1} to scan all of them.
*/
- @Override public int partition() {
+ public int partition() {
+ return part;
+ }
+
+ /** {@inheritDoc} */
+ @Override public int stripeIdx() {
return part;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/transactions/IgniteTxHandler.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/transactions/IgniteTxHandler.java
index 389bb5f6881..97ca6dfbfda 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/transactions/IgniteTxHandler.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/transactions/IgniteTxHandler.java
@@ -341,7 +341,7 @@ public class IgniteTxHandler {
U.error(log, "Failed to prepare DHT transaction: " +
locTx, e);
return new GridNearTxPrepareResponse(
- req.partition(),
+ req.stripeIdx(),
req.version(),
req.futureId(),
req.miniId(),
@@ -512,7 +512,7 @@ public class IgniteTxHandler {
if (retry) {
GridNearTxPrepareResponse res = new
GridNearTxPrepareResponse(
- req.partition(),
+ req.stripeIdx(),
req.version(),
req.futureId(),
req.miniId(),
@@ -720,7 +720,7 @@ public class IgniteTxHandler {
", req=" + req + ']');
GridNearTxPrepareResponse res = new GridNearTxPrepareResponse(
- req.partition(),
+ req.stripeIdx(),
req.version(),
req.futureId(),
req.miniId(),
@@ -1005,7 +1005,7 @@ public class IgniteTxHandler {
// Always send finish response.
GridCacheMessage res = new GridNearTxFinishResponse(
- req.partition(),
+ req.stripeIdx(),
req.version(),
req.threadId(),
req.futureId(),
@@ -1184,7 +1184,7 @@ public class IgniteTxHandler {
try {
res = new GridDhtTxPrepareResponse(
- req.partition(),
+ req.stripeIdx(),
req.version(),
req.futureId(),
req.miniId(),
@@ -1273,7 +1273,7 @@ public class IgniteTxHandler {
}
res = new GridDhtTxPrepareResponse(
- req.partition(),
+ req.stripeIdx(),
req.version(),
req.futureId(),
req.miniId(),
@@ -1590,7 +1590,7 @@ public class IgniteTxHandler {
private void sendReply(UUID nodeId, GridDhtTxFinishRequest req, boolean
committed, GridCacheVersion nearTxId) {
if (req.replyRequired() || req.checkCommitted()) {
GridDhtTxFinishResponse res = new GridDhtTxFinishResponse(
- req.partition(),
+ req.stripeIdx(),
req.version(),
req.futureId(),
req.miniId());
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java
index c064aa455a9..1040f06ddf8 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java
@@ -66,7 +66,6 @@ import
org.apache.ignite.internal.IgniteInterruptedCheckedException;
import org.apache.ignite.internal.IgniteNodeAttributes;
import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException;
import
org.apache.ignite.internal.cluster.ClusterTopologyServerNotFoundException;
-import org.apache.ignite.internal.managers.communication.GridIoMessage;
import org.apache.ignite.internal.managers.communication.GridMessageListener;
import org.apache.ignite.internal.managers.deployment.GridDeployment;
import org.apache.ignite.internal.managers.eventstorage.GridLocalEventListener;
@@ -125,6 +124,7 @@ import org.jetbrains.annotations.Nullable;
import static org.apache.ignite.events.EventType.EVT_NODE_FAILED;
import static org.apache.ignite.events.EventType.EVT_NODE_LEFT;
import static org.apache.ignite.internal.GridTopic.TOPIC_DATASTREAM;
+import static org.apache.ignite.internal.StripedMessage.NO_STRIPE;
/**
* Data streamer implementation.
@@ -2007,8 +2007,7 @@ public class DataStreamerImpl<K, V> implements
IgniteDataStreamer<K, V>, Delayed
dep != null ? dep.classLoaderId() : null,
dep == null,
topVer,
- (rcvr == ISOLATED_UPDATER) ?
- partId : GridIoMessage.STRIPE_DISABLED_PART);
+ (rcvr == ISOLATED_UPDATER) ? partId : NO_STRIPE);
try {
ctx.io().sendToGridTopic(node, TOPIC_DATASTREAM, req, plc);
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java
index c9e62b0939b..88af958a503 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java
@@ -23,6 +23,7 @@ import java.util.UUID;
import org.apache.ignite.configuration.DeploymentMode;
import org.apache.ignite.internal.DeferredUnmarshalMessage;
import org.apache.ignite.internal.Order;
+import org.apache.ignite.internal.StripedMessage;
import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
import org.apache.ignite.internal.processors.cache.GridCacheUtils;
import org.apache.ignite.internal.util.tostring.GridToStringInclude;
@@ -35,7 +36,7 @@ import org.jetbrains.annotations.Nullable;
import static org.apache.ignite.internal.GridTopic.TOPIC_DATASTREAM;
/** */
-public class DataStreamerRequest implements DeferredUnmarshalMessage,
CacheIdAware {
+public class DataStreamerRequest implements DeferredUnmarshalMessage,
CacheIdAware, StripedMessage {
/** */
@Order(0)
long reqId;
@@ -238,8 +239,8 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
return topVer;
}
- /** @return Partition ID. */
- public int partition() {
+ /** {@inheritDoc} */
+ @Override public int stripeIdx() {
return partId;
}