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 270a71f03b3 IGNITE-28270 Investigate possibity to use
MessageSerializer for ComputeJobSibling (#13445)
270a71f03b3 is described below
commit 270a71f03b3b0019ea599d58df5e2a2c79109003
Author: Vladimir Steshin <[email protected]>
AuthorDate: Thu Aug 13 19:35:23 2026 +0300
IGNITE-28270 Investigate possibity to use MessageSerializer for
ComputeJobSibling (#13445)
---
.../ignite/internal/GridJobExecuteRequest.java | 38 +++++++++-------------
.../apache/ignite/internal/GridJobSiblingImpl.java | 35 +++-----------------
.../ignite/internal/GridJobSiblingsResponse.java | 19 ++++-------
.../ignite/internal/GridTaskSessionImpl.java | 6 ++--
.../internal/processors/job/GridJobProcessor.java | 35 +++++++++++++++-----
.../session/GridTaskSessionProcessor.java | 2 +-
.../internal/processors/task/GridTaskWorker.java | 2 +-
.../main/resources/META-INF/classnames.properties | 1 -
8 files changed, 58 insertions(+), 80 deletions(-)
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/GridJobExecuteRequest.java
b/modules/core/src/main/java/org/apache/ignite/internal/GridJobExecuteRequest.java
index c431610e1ea..a90c99a1d2f 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/GridJobExecuteRequest.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/GridJobExecuteRequest.java
@@ -28,6 +28,7 @@ import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfo;
import
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
import org.apache.ignite.internal.util.tostring.GridToStringExclude;
+import org.apache.ignite.internal.util.typedef.F;
import org.apache.ignite.internal.util.typedef.internal.S;
import org.apache.ignite.internal.util.typedef.internal.U;
import org.apache.ignite.lang.IgnitePredicate;
@@ -102,35 +103,27 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
@Order(10)
String cpSpi;
- /** Left unset for a continuous task: such a job requests its siblings
from the task node instead. */
- @Marshalled("siblingsBytes")
- Collection<ComputeJobSibling> siblings;
-
/** */
@Order(11)
- byte[] siblingsBytes;
+ @Nullable IgniteUuid[] siblingJobsIds;
/** Transient since needs to hold local creation time. */
private final long createTime = U.currentTimeMillis();
/** */
@Order(12)
- boolean dynamicSiblings;
-
- /** */
- @Order(13)
boolean forceLocDep;
/** */
- @Order(14)
+ @Order(13)
boolean sesFullSup;
/** */
- @Order(15)
+ @Order(14)
boolean internal;
/** */
- @Order(16)
+ @Order(15)
Collection<UUID> top;
/** */
@@ -138,23 +131,23 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
IgnitePredicate<ClusterNode> topPred;
/** */
- @Order(17)
+ @Order(16)
byte[] topPredBytes;
/** */
- @Order(18)
+ @Order(17)
int[] cacheIds;
/** */
- @Order(19)
+ @Order(18)
int part;
/** */
- @Order(20)
+ @Order(19)
AffinityTopologyVersion topVer;
/** */
- @Order(21)
+ @Order(20)
String execName;
/**
@@ -232,10 +225,8 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
this.top = top;
this.topVer = topVer;
this.topPred = topPred;
- this.siblings = dynamicSiblings ? null : siblings;
this.sesAttrs = sesAttrs;
this.jobAttrs = jobAttrs;
- this.dynamicSiblings = dynamicSiblings;
this.forceLocDep = forceLocDep;
this.sesFullSup = sesFullSup;
this.internal = internal;
@@ -245,6 +236,9 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
this.execName = execName;
this.cpSpi = cpSpi == null || cpSpi.isEmpty() ? null : cpSpi;
+
+ if (!dynamicSiblings && !F.isEmpty(siblings))
+ siblingJobsIds =
siblings.stream().map(ComputeJobSibling::getJobId).toArray(IgniteUuid[]::new);
}
/**
@@ -311,10 +305,10 @@ public class GridJobExecuteRequest implements
ExecutorAwareMessage, DeferredUnma
}
/**
- * @return Job siblings.
+ * @return Siblings jobs ids.
*/
- public Collection<ComputeJobSibling> getSiblings() {
- return siblings;
+ public @Nullable IgniteUuid[] siblingJobsIds() {
+ return siblingJobsIds;
}
/**
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/GridJobSiblingImpl.java
b/modules/core/src/main/java/org/apache/ignite/internal/GridJobSiblingImpl.java
index 5004aa946df..3c20f5c84ba 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/GridJobSiblingImpl.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/GridJobSiblingImpl.java
@@ -17,10 +17,6 @@
package org.apache.ignite.internal;
-import java.io.Externalizable;
-import java.io.IOException;
-import java.io.ObjectInput;
-import java.io.ObjectOutput;
import java.util.Collection;
import java.util.UUID;
import org.apache.ignite.IgniteCheckedException;
@@ -40,16 +36,12 @@ import static
org.apache.ignite.internal.managers.communication.GridIoPolicy.SYS
/**
* This class provides implementation for job sibling.
*/
-public class GridJobSiblingImpl implements ComputeJobSibling, Externalizable {
+public class GridJobSiblingImpl implements ComputeJobSibling {
/** */
- private static final long serialVersionUID = 0L;
+ IgniteUuid sesId;
/** */
- private IgniteUuid sesId;
-
- /** */
- @SuppressWarnings({"FieldAccessedSynchronizedAndUnsynchronized"})
- private IgniteUuid jobId;
+ final IgniteUuid jobId;
/** */
private Object taskTopic;
@@ -64,12 +56,7 @@ public class GridJobSiblingImpl implements
ComputeJobSibling, Externalizable {
private boolean isJobDone;
/** */
- private transient GridKernalContext ctx;
-
- /** */
- public GridJobSiblingImpl() {
- // No-op.
- }
+ private GridKernalContext ctx;
/**
* @param sesId Task session ID.
@@ -173,20 +160,6 @@ public class GridJobSiblingImpl implements
ComputeJobSibling, Externalizable {
ctx.job().cancelJob(sesId, jobId, false);
}
- /** {@inheritDoc} */
- @Override public void writeExternal(ObjectOutput out) throws IOException {
- // Don't serialize node ID.
- U.writeIgniteUuid(out, sesId);
- U.writeIgniteUuid(out, jobId);
- }
-
- /** {@inheritDoc} */
- @Override public void readExternal(ObjectInput in) throws IOException,
ClassNotFoundException {
- // Don't serialize node ID.
- sesId = U.readIgniteUuid(in);
- jobId = U.readIgniteUuid(in);
- }
-
/** {@inheritDoc} */
@Override public String toString() {
return S.toString(GridJobSiblingImpl.class, this);
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/GridJobSiblingsResponse.java
b/modules/core/src/main/java/org/apache/ignite/internal/GridJobSiblingsResponse.java
index d998cdcc156..1470f687a52 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/GridJobSiblingsResponse.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/GridJobSiblingsResponse.java
@@ -19,22 +19,19 @@ package org.apache.ignite.internal;
import java.util.Collection;
import org.apache.ignite.compute.ComputeJobSibling;
+import org.apache.ignite.internal.util.typedef.F;
import org.apache.ignite.internal.util.typedef.internal.S;
+import org.apache.ignite.lang.IgniteUuid;
import org.apache.ignite.plugin.extensions.communication.Message;
import org.jetbrains.annotations.Nullable;
/**
* Job siblings response.
*/
-@UseBinaryMarshaller
public class GridJobSiblingsResponse implements Message {
- /** */
- @Marshalled("siblingsBytes")
- @Nullable Collection<ComputeJobSibling> siblings;
-
/** */
@Order(0)
- byte[] siblingsBytes;
+ public @Nullable IgniteUuid[] siblingJobsIds;
/**
* Empty constructor.
@@ -47,14 +44,10 @@ public class GridJobSiblingsResponse implements Message {
* @param siblings Siblings.
*/
public GridJobSiblingsResponse(@Nullable Collection<ComputeJobSibling>
siblings) {
- this.siblings = siblings;
- }
+ if (F.isEmpty(siblings))
+ return;
- /**
- * @return Job siblings.
- */
- public @Nullable Collection<ComputeJobSibling> jobSiblings() {
- return siblings;
+ siblingJobsIds =
siblings.stream().map(ComputeJobSibling::getJobId).toArray(IgniteUuid[]::new);
}
/** {@inheritDoc} */
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/GridTaskSessionImpl.java
b/modules/core/src/main/java/org/apache/ignite/internal/GridTaskSessionImpl.java
index 27d07e67b10..53fd6591f7e 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/GridTaskSessionImpl.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/GridTaskSessionImpl.java
@@ -80,7 +80,7 @@ public class GridTaskSessionImpl implements
GridTaskSessionInternal {
private final GridKernalContext ctx;
/** */
- private Collection<ComputeJobSibling> siblings;
+ private @Nullable Collection<ComputeJobSibling> siblings;
/** Guarded by {@link #mux}. */
private Map<Object, Object> attrs;
@@ -166,7 +166,7 @@ public class GridTaskSessionImpl implements
GridTaskSessionInternal {
@Nullable IgnitePredicate<ClusterNode> topPred,
long startTime,
long endTime,
- Collection<ComputeJobSibling> siblings,
+ @Nullable Collection<ComputeJobSibling> siblings,
@Nullable Map<Object, Object> attrs,
GridKernalContext ctx,
boolean fullSup,
@@ -535,7 +535,7 @@ public class GridTaskSessionImpl implements
GridTaskSessionInternal {
}
/** {@inheritDoc} */
- @Override public Collection<ComputeJobSibling> getJobSiblings() {
+ @Override public @Nullable Collection<ComputeJobSibling> getJobSiblings() {
synchronized (mux) {
return siblings;
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/job/GridJobProcessor.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/job/GridJobProcessor.java
index 096206004f8..e0975f49b13 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/job/GridJobProcessor.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/job/GridJobProcessor.java
@@ -20,6 +20,7 @@ package org.apache.ignite.internal.processors.job;
import java.util.AbstractCollection;
import java.util.Arrays;
import java.util.Collection;
+import java.util.Collections;
import java.util.Iterator;
import java.util.Map;
import java.util.NoSuchElementException;
@@ -35,6 +36,7 @@ import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Predicate;
+import java.util.stream.Collectors;
import java.util.stream.Stream;
import org.apache.ignite.IgniteCheckedException;
import org.apache.ignite.IgniteDeploymentException;
@@ -54,6 +56,7 @@ import org.apache.ignite.internal.GridJobContextImpl;
import org.apache.ignite.internal.GridJobExecuteRequest;
import org.apache.ignite.internal.GridJobExecuteResponse;
import org.apache.ignite.internal.GridJobSessionImpl;
+import org.apache.ignite.internal.GridJobSiblingImpl;
import org.apache.ignite.internal.GridJobSiblingsRequest;
import org.apache.ignite.internal.GridJobSiblingsResponse;
import org.apache.ignite.internal.GridKernalContext;
@@ -735,9 +738,14 @@ public class GridJobProcessor extends GridProcessorAdapter
{
// Error is set?
if (t.get1() != null)
throw new IgniteCheckedException(t.get1());
- else
- // Return result
- return t.get2().jobSiblings();
+ else {
+ IgniteUuid[] siblingJobsIds = t.get2().siblingJobsIds;
+
+ return F.isEmpty(siblingJobsIds)
+ ? Collections.emptyList()
+ : Stream.of(siblingJobsIds).map(sibJobId -> new
GridJobSiblingImpl(ses.getId(), sibJobId, taskNodeId, ctx))
+ .collect(Collectors.toList());
+ }
}
catch (InterruptedException e) {
throw new IgniteCheckedException("Interrupted while waiting
for job siblings response: " + ses, e);
@@ -1176,11 +1184,15 @@ public class GridJobProcessor extends
GridProcessorAdapter {
/**
* @param node Node.
* @param req Request.
+ * @param siblingJobs Siblings jobs. TODO : Revise in
https://issues.apache.org/jira/browse/IGNITE-28964
*/
@SuppressWarnings("TooBroadScope")
- public void processJobExecuteRequest(ClusterNode node, final
GridJobExecuteRequest req) {
+ public void processJobExecuteRequest(
+ ClusterNode node, GridJobExecuteRequest req,
+ @Nullable Collection<ComputeJobSibling> siblingJobs
+ ) {
if (log.isDebugEnabled())
- log.debug("Received job request message [req=" + req + ", nodeId="
+ node.id() + ']');
+ log.debug("Processing job request message [req=" + req + ",
nodeId=" + node.id() + ']');
PartitionsReservation partsReservation = null;
@@ -1195,7 +1207,7 @@ public class GridJobProcessor extends
GridProcessorAdapter {
if (!rwLock.tryReadLock()) {
if (log.isDebugEnabled())
- log.debug("Received job execution request while stopping this
node (will ignore): " + req);
+ log.debug("Processing job execution request while stopping
this node (will ignore): " + req);
return;
}
@@ -1258,7 +1270,7 @@ public class GridJobProcessor extends
GridProcessorAdapter {
req.getTopologyPredicate(),
req.startTaskTime(),
endTime,
- req.getSiblings(),
+ siblingJobs,
req.getSessionAttributes(),
req.sessionFullSupport(),
req.internal(),
@@ -2198,7 +2210,14 @@ public class GridJobProcessor extends
GridProcessorAdapter {
assert node != null;
- processJobExecuteRequest(node, (GridJobExecuteRequest)msg);
+ GridJobExecuteRequest req = (GridJobExecuteRequest)msg;
+
+ Collection<ComputeJobSibling> siblingJobs =
F.isEmpty(req.siblingJobsIds())
+ ? null
+ : Stream.of(req.siblingJobsIds()).map(sibJobId -> new
GridJobSiblingImpl(req.sessionId(), sibJobId, nodeId, ctx))
+ .collect(Collectors.toList());
+
+ processJobExecuteRequest(node, (GridJobExecuteRequest)msg,
siblingJobs);
}
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/session/GridTaskSessionProcessor.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/session/GridTaskSessionProcessor.java
index f0877df777b..955865226e2 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/session/GridTaskSessionProcessor.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/session/GridTaskSessionProcessor.java
@@ -95,7 +95,7 @@ public class GridTaskSessionProcessor extends
GridProcessorAdapter {
@Nullable IgnitePredicate<ClusterNode> topPred,
long startTime,
long endTime,
- Collection<ComputeJobSibling> siblings,
+ @Nullable Collection<ComputeJobSibling> siblings,
Map<Object, Object> attrs,
boolean fullSup,
boolean internal,
diff --git
a/modules/core/src/main/java/org/apache/ignite/internal/processors/task/GridTaskWorker.java
b/modules/core/src/main/java/org/apache/ignite/internal/processors/task/GridTaskWorker.java
index 5f62ee2b9f0..32df3c6ca5f 100644
---
a/modules/core/src/main/java/org/apache/ignite/internal/processors/task/GridTaskWorker.java
+++
b/modules/core/src/main/java/org/apache/ignite/internal/processors/task/GridTaskWorker.java
@@ -1406,7 +1406,7 @@ public class GridTaskWorker<T, R> extends GridWorker
implements GridTimeoutObjec
ses.executorName());
if (loc)
-
ctx.job().processJobExecuteRequest(ctx.discovery().localNode(), req);
+
ctx.job().processJobExecuteRequest(ctx.discovery().localNode(), req,
ses.getJobSiblings());
else {
byte plc;
diff --git a/modules/core/src/main/resources/META-INF/classnames.properties
b/modules/core/src/main/resources/META-INF/classnames.properties
index b9e5b18ebca..65f0244da47 100644
--- a/modules/core/src/main/resources/META-INF/classnames.properties
+++ b/modules/core/src/main/resources/META-INF/classnames.properties
@@ -220,7 +220,6 @@ org.apache.ignite.internal.GridJobCancelRequest
org.apache.ignite.internal.GridJobContextImpl
org.apache.ignite.internal.GridJobExecuteRequest
org.apache.ignite.internal.GridJobExecuteResponse
-org.apache.ignite.internal.GridJobSiblingImpl
org.apache.ignite.internal.GridJobSiblingsRequest
org.apache.ignite.internal.GridJobSiblingsResponse
org.apache.ignite.internal.GridKernalContextImpl