This is an automated email from the ASF dual-hosted git repository.
wernerdv 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 9d17faa242f IGNITE-28719 Revise GridTopicMessage usage in messages
(#13438)
9d17faa242f is described below
commit 9d17faa242f455492e5cc9b87610656bcf1274ac
Author: Dmitry Werner <[email protected]>
AuthorDate: Thu Aug 6 16:40:24 2026 +0500
IGNITE-28719 Revise GridTopicMessage usage in messages (#13438)
---
.../deployment/GridDeploymentCommunication.java | 6 +-
.../managers/deployment/GridDeploymentRequest.java | 31 +++-----
.../datastreamer/DataStreamProcessor.java | 14 ----
.../processors/datastreamer/DataStreamerImpl.java | 9 ++-
.../datastreamer/DataStreamerRequest.java | 85 +++++++---------------
...loymentRequestOfUnknownClassProcessingTest.java | 6 +-
.../datastreamer/DataStreamerImplSelfTest.java | 2 +-
7 files changed, 53 insertions(+), 100 deletions(-)
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentCommunication.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentCommunication.java
index a36e98fcdee..dddef831d9f 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentCommunication.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentCommunication.java
@@ -374,9 +374,11 @@ class GridDeploymentCommunication {
", requesters=" + nodeIds + ']');
}
- Object resTopic =
TOPIC_CLASSLOAD.topic(IgniteUuid.fromUuid(ctx.localNodeId()));
+ IgniteUuid resTopicId = IgniteUuid.fromUuid(ctx.localNodeId());
- GridDeploymentRequest req = new GridDeploymentRequest(resTopic,
clsLdrId, rsrcName);
+ Object resTopic = TOPIC_CLASSLOAD.topic(resTopicId);
+
+ GridDeploymentRequest req = new GridDeploymentRequest(resTopicId,
clsLdrId, rsrcName);
// Send node IDs chain with request.
req.nodeIds(nodeIds);
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentRequest.java
index a41baab9ef1..eaa6a8e20b9 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentRequest.java
@@ -19,7 +19,6 @@ package org.apache.ignite.internal.managers.deployment;
import java.util.Collection;
import java.util.UUID;
-import org.apache.ignite.internal.GridTopicMessage;
import org.apache.ignite.internal.Order;
import org.apache.ignite.internal.util.tostring.GridToStringInclude;
import org.apache.ignite.internal.util.typedef.internal.S;
@@ -27,13 +26,13 @@ import org.apache.ignite.lang.IgniteUuid;
import org.apache.ignite.plugin.extensions.communication.Message;
import org.jetbrains.annotations.Nullable;
-/**
- * Deployment request.
- */
+import static org.apache.ignite.internal.GridTopic.TOPIC_CLASSLOAD;
+
+/** Deployment request. */
public class GridDeploymentRequest implements Message {
- /** Response topic message. Response should be sent back to this topic. */
+ /** ID of the node waiting for the response. */
@Order(0)
- @Nullable GridTopicMessage topicMsg;
+ @Nullable IgniteUuid resTopicId;
/** Requested class name. */
@Order(1)
@@ -58,12 +57,12 @@ public class GridDeploymentRequest implements Message {
/**
* Creates deploy request.
*
- * @param topic Response topic.
+ * @param resTopicId ID of the node waiting for the response.
* @param ldrId Class loader ID.
* @param rsrcName Resource name that should be found and sent back.
*/
- GridDeploymentRequest(Object topic, IgniteUuid ldrId, String rsrcName) {
- topicMsg = new GridTopicMessage(topic);
+ GridDeploymentRequest(IgniteUuid resTopicId, IgniteUuid ldrId, String
rsrcName) {
+ this.resTopicId = resTopicId;
this.ldrId = ldrId;
this.rsrcName = rsrcName;
}
@@ -83,7 +82,7 @@ public class GridDeploymentRequest implements Message {
* @return Response topic name.
*/
@Nullable Object responseTopic() {
- return GridTopicMessage.topic(topicMsg);
+ return undeploy() ? null : TOPIC_CLASSLOAD.topic(resTopicId);
}
/**
@@ -110,21 +109,15 @@ public class GridDeploymentRequest implements Message {
* @return Property undeploy.
*/
boolean undeploy() {
- return topicMsg == null;
+ return resTopicId == null;
}
- /**
- * @return Node IDs chain which is updated as request jumps
- * from node to node.
- */
+ /** @return Node IDs chain which is updated as request jumps from node to
node. */
Collection<UUID> nodeIds() {
return nodeIds;
}
- /**
- * @param nodeIds Node IDs chain which is updated as request jumps
- * from node to node.
- */
+ /** @param nodeIds Node IDs chain which is updated as request jumps from
node to node. */
void nodeIds(Collection<UUID> nodeIds) {
this.nodeIds = nodeIds;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamProcessor.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamProcessor.java
index 371b0487d65..464a74d82ee 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamProcessor.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamProcessor.java
@@ -27,7 +27,6 @@ import
org.apache.ignite.internal.IgniteInterruptedCheckedException;
import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException;
import org.apache.ignite.internal.managers.communication.GridIoManager;
import org.apache.ignite.internal.managers.communication.GridMessageListener;
-import org.apache.ignite.internal.managers.communication.MessageMarshalling;
import org.apache.ignite.internal.managers.deployment.GridDeployment;
import org.apache.ignite.internal.processors.GridProcessorAdapter;
import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
@@ -208,19 +207,6 @@ public class DataStreamProcessor extends
GridProcessorAdapter {
}
}
- // The generic receive pass skips this request (entries must stay
serialized for the update job and its
- // peer-deployment loader), so the response topic is restored
here; it carries only internal classes.
- try {
- if (req.resTopicMsg != null)
- MessageMarshalling.unmarshal(req.resTopicMsg, ctx);
- }
- catch (IgniteCheckedException e) {
- U.error(log, "Failed to unmarshal response topic (no response
will be sent) [nodeId=" + nodeId +
- ", req=" + req + ']', e);
-
- return;
- }
-
Object topic = req.responseTopic();
ClassLoader clsLdr;
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 9a030b8e335..c064aa455a9 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
@@ -205,6 +205,9 @@ public class DataStreamerImpl<K, V> implements
IgniteDataStreamer<K, V>, Delayed
/** Communication topic for responses. */
private final Object topic;
+ /** Topic ID for responses. */
+ private final IgniteUuid topicId;
+
/** {@code True} if data loader has been cancelled. */
private volatile boolean cancelled;
@@ -349,8 +352,10 @@ public class DataStreamerImpl<K, V> implements
IgniteDataStreamer<K, V>, Delayed
ctx.event().addLocalEventListener(discoLsnr, EVT_NODE_FAILED,
EVT_NODE_LEFT);
+ topicId = IgniteUuid.fromUuid(ctx.localNodeId());
+
// Generate unique topic for this loader.
- topic = TOPIC_DATASTREAM.topic(IgniteUuid.fromUuid(ctx.localNodeId()));
+ topic = TOPIC_DATASTREAM.topic(topicId);
ctx.io().addMessageListener(topic, new GridMessageListener() {
@Override public void onMessage(UUID nodeId, Object msg, byte plc)
{
@@ -1988,7 +1993,7 @@ public class DataStreamerImpl<K, V> implements
IgniteDataStreamer<K, V>, Delayed
DataStreamerRequest req = new DataStreamerRequest(
reqId,
- topic,
+ topicId,
cacheName,
updaterBytes,
entries,
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 0167489c254..c9e62b0939b 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
@@ -22,7 +22,6 @@ import java.util.Map;
import java.util.UUID;
import org.apache.ignite.configuration.DeploymentMode;
import org.apache.ignite.internal.DeferredUnmarshalMessage;
-import org.apache.ignite.internal.GridTopicMessage;
import org.apache.ignite.internal.Order;
import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
import org.apache.ignite.internal.processors.cache.GridCacheUtils;
@@ -33,9 +32,9 @@ import
org.apache.ignite.plugin.extensions.communication.CacheIdAware;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
-/**
- *
- */
+import static org.apache.ignite.internal.GridTopic.TOPIC_DATASTREAM;
+
+/** */
public class DataStreamerRequest implements DeferredUnmarshalMessage,
CacheIdAware {
/** */
@Order(0)
@@ -43,7 +42,7 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
/** */
@Order(1)
- GridTopicMessage resTopicMsg;
+ IgniteUuid resTopicId;
/** Cache name. */
@Order(2)
@@ -103,16 +102,14 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
@Order(15)
int partId;
- /**
- * Empty constructor.
- */
+ /** Empty constructor. */
public DataStreamerRequest() {
// No-op.
}
/**
* @param reqId Request ID.
- * @param resTopic Response topic.
+ * @param resTopicId Response topic ID.
* @param cacheName Cache name.
* @param updaterBytes Cache receiver.
* @param entries Entries to put.
@@ -130,7 +127,7 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
*/
public DataStreamerRequest(
long reqId,
- Object resTopic,
+ IgniteUuid resTopicId,
@Nullable String cacheName,
byte[] updaterBytes,
Collection<DataStreamerEntry> entries,
@@ -149,7 +146,7 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
assert topVer != null;
this.reqId = reqId;
- resTopicMsg = new GridTopicMessage(resTopic);
+ this.resTopicId = resTopicId;
this.cacheName = cacheName;
this.updaterBytes = updaterBytes;
this.entries = entries;
@@ -166,114 +163,82 @@ public class DataStreamerRequest implements
DeferredUnmarshalMessage, CacheIdAwa
this.partId = partId;
}
- /**
- * @return Request ID.
- */
+ /** @return Request ID. */
long requestId() {
return reqId;
}
- /**
- * @return Response topic.
- */
+ /** @return Response topic. */
Object responseTopic() {
- return GridTopicMessage.topic(resTopicMsg);
+ return TOPIC_DATASTREAM.topic(resTopicId);
}
- /**
- * @return Cache name.
- */
+ /** @return Cache name. */
String cacheName() {
return cacheName;
}
- /**
- * @return Updater.
- */
+ /** @return Updater. */
byte[] updaterBytes() {
return updaterBytes;
}
- /**
- * @return Entries to update.
- */
+ /** @return Entries to update. */
Collection<DataStreamerEntry> entries() {
return entries;
}
- /**
- * @return {@code True} to ignore ownership.
- */
+ /** @return {@code True} to ignore ownership. */
boolean ignoreDeploymentOwnership() {
return ignoreDepOwnership;
}
- /**
- * @return Skip store flag.
- */
+ /** @return Skip store flag. */
boolean skipStore() {
return skipStore;
}
- /**
- * @return Keep binary flag.
- */
+ /** @return Keep binary flag. */
boolean keepBinary() {
return keepBinary;
}
- /**
- * @return Deployment mode.
- */
+ /** @return Deployment mode. */
DeploymentMode deploymentMode() {
return depMode;
}
- /**
- * @return Sample class name.
- */
+ /** @return Sample class name. */
String sampleClassName() {
return sampleClsName;
}
- /**
- * @return User version.
- */
+ /** @return User version. */
String userVersion() {
return userVer;
}
- /**
- * @return Participants.
- */
+ /** @return Participants. */
Map<UUID, IgniteUuid> participants() {
return ldrParticipants;
}
- /**
- * @return Class loader ID.
- */
+ /** @return Class loader ID. */
IgniteUuid classLoaderId() {
return clsLdrId;
}
- /**
- * @return {@code True} to force local deployment.
- */
+ /** @return {@code True} to force local deployment. */
boolean forceLocalDeployment() {
return forceLocDep;
}
- /**
- * @return Topology version.
- */
+ /** @return Topology version. */
AffinityTopologyVersion topologyVersion() {
return topVer;
}
- /**
- * @return Partition ID.
- */
+ /** @return Partition ID. */
public int partition() {
return partId;
}
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/managers/deployment/DeploymentRequestOfUnknownClassProcessingTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/managers/deployment/DeploymentRequestOfUnknownClassProcessingTest.java
index 5ab4b6b4b22..5ef17a1ab13 100644
---
a/modules/core/src/test/java/org/apache/ignite/internal/managers/deployment/DeploymentRequestOfUnknownClassProcessingTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/internal/managers/deployment/DeploymentRequestOfUnknownClassProcessingTest.java
@@ -98,7 +98,9 @@ public class DeploymentRequestOfUnknownClassProcessingTest
extends GridCommonAbs
remNodeLog.registerListener(remNodeLogLsnr);
- Object topic =
TOPIC_CLASSLOAD.topic(IgniteUuid.fromUuid(locNode.localNode().id()));
+ IgniteUuid topicId = IgniteUuid.fromUuid(locNode.localNode().id());
+
+ Object topic = TOPIC_CLASSLOAD.topic(topicId);
locNode.context().io().addMessageListener(topic, new
GridMessageListener() {
@Override public void onMessage(UUID nodeId, Object msg, byte plc)
{
@@ -124,7 +126,7 @@ public class DeploymentRequestOfUnknownClassProcessingTest
extends GridCommonAbs
}
});
- GridDeploymentRequest req = new GridDeploymentRequest(topic,
locDep.classLoaderId(), UNKNOWN_CLASS_NAME);
+ GridDeploymentRequest req = new GridDeploymentRequest(topicId,
locDep.classLoaderId(), UNKNOWN_CLASS_NAME);
locNode.context().io().sendToGridTopic(remNode.localNode(),
TOPIC_CLASSLOAD, req, GridIoPolicy.P2P_POOL);
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java
index f99462429d5..2bc777180e0 100644
---
a/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java
@@ -690,7 +690,7 @@ public class DataStreamerImplSelfTest extends
GridCommonAbstractTest {
appMsg = new DataStreamerRequest(
req.requestId(),
- req.responseTopic(),
+ req.resTopicId,
req.cacheName(),
req.updaterBytes(),
req.entries(),