From 2d8cbb722ed767a22fc887ce451e6b06806c92c2 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Fri, 7 Aug 2026 11:27:36 +0300 Subject: [PATCH 1/8] impl --- .../ignite/internal/CoreMessagesProvider.java | 1 + .../ignite/internal/GridJobSiblingImpl.java | 32 ++++--------------- 2 files changed, 8 insertions(+), 25 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java index 9412a9ab9fa7d..dab42e6f98093 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java @@ -641,6 +641,7 @@ public CoreMessagesProvider(Marshaller dfltMarsh, Marshaller schemaAwareMarsh) { register(GridTaskResultResponse.class); register(JobStealingRequest.class); register(SingleNodeMessage.class); + register(GridJobSiblingImpl.class); // [11500 - 11600]: IO, networking messages. msgIdx = NODE_ID_MSG_TYPE; 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 5004aa946dfa2..a8ac36c444fee 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; @@ -31,6 +27,7 @@ import org.apache.ignite.internal.util.typedef.internal.S; import org.apache.ignite.internal.util.typedef.internal.U; import org.apache.ignite.lang.IgniteUuid; +import org.apache.ignite.plugin.extensions.communication.Message; import static org.apache.ignite.internal.GridTopic.TOPIC_JOB; import static org.apache.ignite.internal.GridTopic.TOPIC_JOB_CANCEL; @@ -40,16 +37,15 @@ /** * This class provides implementation for job sibling. */ -public class GridJobSiblingImpl implements ComputeJobSibling, Externalizable { +public class GridJobSiblingImpl implements ComputeJobSibling, Message { /** */ - private static final long serialVersionUID = 0L; - - /** */ - private IgniteUuid sesId; + @Order(0) + IgniteUuid sesId; /** */ @SuppressWarnings({"FieldAccessedSynchronizedAndUnsynchronized"}) - private IgniteUuid jobId; + @Order(1) + IgniteUuid jobId; /** */ private Object taskTopic; @@ -66,7 +62,7 @@ public class GridJobSiblingImpl implements ComputeJobSibling, Externalizable { /** */ private transient GridKernalContext ctx; - /** */ + /** Empty constructor for serialization purposes. */ public GridJobSiblingImpl() { // No-op. } @@ -173,20 +169,6 @@ public synchronized Object jobTopic() { 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); From 31a73ab7467c181a0dad71ea244775f7c598a6bf Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Fri, 7 Aug 2026 20:14:21 +0300 Subject: [PATCH 2/8] try fix --- modules/core/src/main/resources/META-INF/classnames.properties | 1 - 1 file changed, 1 deletion(-) diff --git a/modules/core/src/main/resources/META-INF/classnames.properties b/modules/core/src/main/resources/META-INF/classnames.properties index fc1966fef4b4e..d1282ad57f523 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 From 94a165545b4af421ef9ab4702dab7fc754201a64 Mon Sep 17 00:00:00 2001 From: Vladimir Steshin Date: Sat, 8 Aug 2026 13:01:26 +0300 Subject: [PATCH 3/8] fixes --- .../ignite/internal/CoreMessagesProvider.java | 1 - .../internal/GridJobExecuteRequest.java | 23 +++++++++++-------- .../ignite/internal/GridJobSiblingImpl.java | 17 ++++---------- .../processors/job/GridJobProcessor.java | 5 +++- 4 files changed, 21 insertions(+), 25 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java index 46d6dc512c57b..a460ed789d24b 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/CoreMessagesProvider.java @@ -641,7 +641,6 @@ public CoreMessagesProvider(Marshaller dfltMarsh, Marshaller schemaAwareMarsh) { register(GridTaskResultResponse.class); register(JobStealingRequest.class); register(SingleNodeMessage.class); - register(GridJobSiblingImpl.class); // [11500 - 11600]: IO, networking messages. msgIdx = NODE_ID_MSG_TYPE; 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 098e8514f681c..4db9931bf0703 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.processors.affinity.AffinityTopologyVersion; import org.apache.ignite.internal.util.tostring.GridToStringExclude; import org.apache.ignite.internal.util.tostring.GridToStringInclude; +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; @@ -107,13 +108,13 @@ public class GridJobExecuteRequest implements ExecutorAwareMessage, DeferredUnma @Order(11) String cpSpi; - /** Left unset for a continuous task: such a job requests its siblings from the task node instead. */ - @Marshalled("siblingsBytes") - Collection siblings; - - /** */ + /** + * Sibling jobs ids. Plain representation of {@link GridJobSiblingImpl#jobId} to reduce the messages number. + * + * @see ComputeJobSibling + */ @Order(12) - byte[] siblingsBytes; + @Nullable Collection sibJobIds; /** Transient since needs to hold local creation time. */ private final long createTime = U.currentTimeMillis(); @@ -253,7 +254,6 @@ public GridJobExecuteRequest( this.top = top; this.topVer = topVer; this.topPred = topPred; - this.siblings = dynamicSiblings ? null : siblings; this.sesAttrs = sesAttrs; this.jobAttrs = jobAttrs; this.clsLdrId = clsLdrId; @@ -269,6 +269,9 @@ public GridJobExecuteRequest( this.execName = execName; this.cpSpi = cpSpi == null || cpSpi.isEmpty() ? null : cpSpi; + + if(!dynamicSiblings && !F.isEmpty(siblings)) + sibJobIds = F.viewReadOnly(siblings, ComputeJobSibling::getJobId); } /** @@ -337,10 +340,10 @@ public long getCreateTime() { } /** - * @return Job siblings. + * @return Sibling job ids. */ - public Collection getSiblings() { - return siblings; + public @Nullable Collection getSiblings() { + return sibJobIds; } /** 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 a8ac36c444fee..eca9a3b18cddc 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 @@ -27,7 +27,6 @@ import org.apache.ignite.internal.util.typedef.internal.S; import org.apache.ignite.internal.util.typedef.internal.U; import org.apache.ignite.lang.IgniteUuid; -import org.apache.ignite.plugin.extensions.communication.Message; import static org.apache.ignite.internal.GridTopic.TOPIC_JOB; import static org.apache.ignite.internal.GridTopic.TOPIC_JOB_CANCEL; @@ -37,15 +36,12 @@ /** * This class provides implementation for job sibling. */ -public class GridJobSiblingImpl implements ComputeJobSibling, Message { +public class GridJobSiblingImpl implements ComputeJobSibling { /** */ - @Order(0) - IgniteUuid sesId; + private IgniteUuid sesId; /** */ - @SuppressWarnings({"FieldAccessedSynchronizedAndUnsynchronized"}) - @Order(1) - IgniteUuid jobId; + private final IgniteUuid jobId; /** */ private Object taskTopic; @@ -60,12 +56,7 @@ public class GridJobSiblingImpl implements ComputeJobSibling, Message { private boolean isJobDone; /** */ - private transient GridKernalContext ctx; - - /** Empty constructor for serialization purposes. */ - public GridJobSiblingImpl() { - // No-op. - } + private GridKernalContext ctx; /** * @param sesId Task session ID. 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 b81a8b43bf21e..b3840d6039a16 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 @@ 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; @@ -1256,6 +1257,8 @@ public void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque U.resolveClassLoader(dep.classLoader(), ctx.config())); } + Collection siblings = Collections.emptyList(); + GridTaskSessionImpl taskSes = ctx.session().createTaskSession( req.sessionId(), node.id(), @@ -1266,7 +1269,7 @@ public void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque req.getTopologyPredicate(), req.startTaskTime(), endTime, - req.getSiblings(), + siblings, req.getSessionAttributes(), req.sessionFullSupport(), req.internal(), From 07341b544d0bbc0a17e4497c9c9c648819778d84 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Sat, 8 Aug 2026 13:15:07 +0300 Subject: [PATCH 4/8] fixes --- .../org/apache/ignite/internal/GridJobExecuteRequest.java | 6 +++--- .../ignite/internal/processors/job/GridJobProcessor.java | 5 +++++ 2 files changed, 8 insertions(+), 3 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 4db9931bf0703..1540acadeeded 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 @@ -270,8 +270,8 @@ public GridJobExecuteRequest( this.cpSpi = cpSpi == null || cpSpi.isEmpty() ? null : cpSpi; - if(!dynamicSiblings && !F.isEmpty(siblings)) - sibJobIds = F.viewReadOnly(siblings, ComputeJobSibling::getJobId); + if (!dynamicSiblings && !F.isEmpty(siblings)) + sibJobIds = siblings.stream().map(sib -> sib.getJobId()).toList(); } /** @@ -342,7 +342,7 @@ public long getCreateTime() { /** * @return Sibling job ids. */ - public @Nullable Collection getSiblings() { + public @Nullable Collection siblingJobsIds() { return sibJobIds; } 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 b3840d6039a16..f34323a5a2176 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 @@ -36,6 +36,7 @@ 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; @@ -55,6 +56,7 @@ 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; @@ -1259,6 +1261,9 @@ public void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque Collection siblings = Collections.emptyList(); + if (!F.isEmpty(req.siblingJobsIds())) + siblings = req.siblingJobsIds().stream().map(sibJobId -> new GridJobSiblingImpl(null, sibJobId, null, null)).collect(Collectors.toList()); + GridTaskSessionImpl taskSes = ctx.session().createTaskSession( req.sessionId(), node.id(), From 4ae70b4f3e559d59744801bffaf128346ad6fafb Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Sat, 8 Aug 2026 13:23:58 +0300 Subject: [PATCH 5/8] fixes --- .../org/apache/ignite/internal/GridJobSiblingImpl.java | 1 - .../org/apache/ignite/internal/GridTaskSessionImpl.java | 2 +- .../ignite/internal/processors/job/GridJobProcessor.java | 9 ++++----- .../processors/session/GridTaskSessionProcessor.java | 2 +- 4 files changed, 6 insertions(+), 8 deletions(-) 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 eca9a3b18cddc..674d6e08666dc 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 @@ -65,7 +65,6 @@ public class GridJobSiblingImpl implements ComputeJobSibling { * @param ctx Managers registry. */ public GridJobSiblingImpl(IgniteUuid sesId, IgniteUuid jobId, UUID nodeId, GridKernalContext ctx) { - assert sesId != null; assert jobId != null; assert nodeId != null; assert ctx != null; 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 27d07e67b10e4..5ccf55e209d50 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 @@ -166,7 +166,7 @@ public GridTaskSessionImpl( @Nullable IgnitePredicate topPred, long startTime, long endTime, - Collection siblings, + @Nullable Collection siblings, @Nullable Map attrs, GridKernalContext ctx, boolean fullSup, 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 f34323a5a2176..db41e7437bc2e 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,7 +20,6 @@ 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; @@ -1259,10 +1258,10 @@ public void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque U.resolveClassLoader(dep.classLoader(), ctx.config())); } - Collection siblings = Collections.emptyList(); - - if (!F.isEmpty(req.siblingJobsIds())) - siblings = req.siblingJobsIds().stream().map(sibJobId -> new GridJobSiblingImpl(null, sibJobId, null, null)).collect(Collectors.toList()); + Collection siblings = F.isEmpty(req.siblingJobsIds()) ? + null + : req.siblingJobsIds().stream().map(sibJobId -> new GridJobSiblingImpl(null, sibJobId, node.id(), ctx)) + .collect(Collectors.toList()); GridTaskSessionImpl taskSes = ctx.session().createTaskSession( req.sessionId(), 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 f0877df777b1f..955865226e20c 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 GridTaskSessionImpl createTaskSession( @Nullable IgnitePredicate topPred, long startTime, long endTime, - Collection siblings, + @Nullable Collection siblings, Map attrs, boolean fullSup, boolean internal, From 21d4b18ed3447ee9518a1d8309a444f91d7081c1 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Sat, 8 Aug 2026 13:53:39 +0300 Subject: [PATCH 6/8] minorities --- .../org/apache/ignite/internal/GridJobExecuteRequest.java | 3 ++- .../java/org/apache/ignite/internal/GridJobSiblingImpl.java | 4 +++- .../ignite/internal/processors/job/GridJobProcessor.java | 5 +++-- 3 files changed, 8 insertions(+), 4 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 1540acadeeded..62b2b2b6bf078 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 @@ -109,7 +109,7 @@ public class GridJobExecuteRequest implements ExecutorAwareMessage, DeferredUnma String cpSpi; /** - * Sibling jobs ids. Plain representation of {@link GridJobSiblingImpl#jobId} to reduce the messages number. + * Sibling jobs ids. Plain representation of {@link GridJobSiblingImpl#getJobId()} to reduce the messages number. * * @see ComputeJobSibling */ @@ -216,6 +216,7 @@ public GridJobExecuteRequest( long timeout, @Nullable Collection top, @Nullable IgnitePredicate topPred, + // TODO : Revise after https://issues.apache.org/jira/browse/IGNITE-28964 Collection siblings, Map sesAttrs, Map jobAttrs, 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 674d6e08666dc..32c350a3bab8c 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 @@ -27,6 +27,7 @@ import org.apache.ignite.internal.util.typedef.internal.S; import org.apache.ignite.internal.util.typedef.internal.U; import org.apache.ignite.lang.IgniteUuid; +import org.jetbrains.annotations.Nullable; import static org.apache.ignite.internal.GridTopic.TOPIC_JOB; import static org.apache.ignite.internal.GridTopic.TOPIC_JOB_CANCEL; @@ -35,6 +36,7 @@ /** * This class provides implementation for job sibling. + * TODO : Revise after https://issues.apache.org/jira/browse/IGNITE-28964 */ public class GridJobSiblingImpl implements ComputeJobSibling { /** */ @@ -64,7 +66,7 @@ public class GridJobSiblingImpl implements ComputeJobSibling { * @param nodeId ID of the node where this sibling was sent for execution. * @param ctx Managers registry. */ - public GridJobSiblingImpl(IgniteUuid sesId, IgniteUuid jobId, UUID nodeId, GridKernalContext ctx) { + public GridJobSiblingImpl(@Nullable IgniteUuid sesId, IgniteUuid jobId, UUID nodeId, GridKernalContext ctx) { assert jobId != null; assert nodeId != null; assert ctx != null; 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 db41e7437bc2e..efa3fdca88584 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 @@ -1258,8 +1258,9 @@ public void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque U.resolveClassLoader(dep.classLoader(), ctx.config())); } - Collection siblings = F.isEmpty(req.siblingJobsIds()) ? - null + // TODO : Revise after https://issues.apache.org/jira/browse/IGNITE-28964 + Collection siblings = F.isEmpty(req.siblingJobsIds()) + ? null : req.siblingJobsIds().stream().map(sibJobId -> new GridJobSiblingImpl(null, sibJobId, node.id(), ctx)) .collect(Collectors.toList()); From d86c00071e3551175524231b0f1b66c71c5b0433 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Sat, 8 Aug 2026 16:29:23 +0300 Subject: [PATCH 7/8] fix --- .../internal/GridJobExecuteRequest.java | 65 ++++++++++++------- .../ignite/internal/GridJobSiblingImpl.java | 4 +- .../processors/job/GridJobProcessor.java | 22 +++++-- .../processors/task/GridTaskWorker.java | 15 ++++- 4 files changed, 73 insertions(+), 33 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 62b2b2b6bf078..f4953a0ebf471 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 @@ -18,12 +18,13 @@ package org.apache.ignite.internal; import java.io.Serializable; +import java.util.ArrayList; import java.util.Collection; +import java.util.List; import java.util.Map; import java.util.UUID; import org.apache.ignite.cluster.ClusterNode; import org.apache.ignite.compute.ComputeJob; -import org.apache.ignite.compute.ComputeJobSibling; import org.apache.ignite.configuration.DeploymentMode; import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion; import org.apache.ignite.internal.util.tostring.GridToStringExclude; @@ -108,43 +109,43 @@ public class GridJobExecuteRequest implements ExecutorAwareMessage, DeferredUnma @Order(11) String cpSpi; - /** - * Sibling jobs ids. Plain representation of {@link GridJobSiblingImpl#getJobId()} to reduce the messages number. - * - * @see ComputeJobSibling - */ + /** Sibling jobs ids. Plain representation of {@link GridJobSiblingImpl#jobId} to reduce the messages number. */ @Order(12) - @Nullable Collection sibJobIds; + @Nullable List sibJobsIds; + + /** Sibling jobs ids. Plain representation of {@link GridJobSiblingImpl#sesId} to reduce the messages number. */ + @Order(13) + @Nullable List sibJobsSesId; /** Transient since needs to hold local creation time. */ private final long createTime = U.currentTimeMillis(); /** */ - @Order(13) + @Order(14) IgniteUuid clsLdrId; /** */ - @Order(14) + @Order(15) DeploymentMode depMode; /** */ - @Order(15) + @Order(16) boolean dynamicSiblings; /** */ - @Order(16) + @Order(17) boolean forceLocDep; /** */ - @Order(17) + @Order(18) boolean sesFullSup; /** */ - @Order(18) + @Order(19) boolean internal; /** */ - @Order(19) + @Order(20) Collection top; /** */ @@ -152,23 +153,23 @@ public class GridJobExecuteRequest implements ExecutorAwareMessage, DeferredUnma IgnitePredicate topPred; /** */ - @Order(20) + @Order(21) byte[] topPredBytes; /** */ - @Order(21) + @Order(22) int[] cacheIds; /** */ - @Order(22) + @Order(23) int part; /** */ - @Order(23) + @Order(24) AffinityTopologyVersion topVer; /** */ - @Order(24) + @Order(25) String execName; /** @@ -216,8 +217,7 @@ public GridJobExecuteRequest( long timeout, @Nullable Collection top, @Nullable IgnitePredicate topPred, - // TODO : Revise after https://issues.apache.org/jira/browse/IGNITE-28964 - Collection siblings, + Collection siblings, Map sesAttrs, Map jobAttrs, String cpSpi, @@ -271,8 +271,16 @@ public GridJobExecuteRequest( this.cpSpi = cpSpi == null || cpSpi.isEmpty() ? null : cpSpi; - if (!dynamicSiblings && !F.isEmpty(siblings)) - sibJobIds = siblings.stream().map(sib -> sib.getJobId()).toList(); + // TODO : Revise after https://issues.apache.org/jira/browse/IGNITE-28964 + if (!dynamicSiblings && !F.isEmpty(siblings)) { + sibJobsIds = new ArrayList<>(siblings.size()); + sibJobsSesId = new ArrayList<>(sibJobsIds.size()); + + siblings.forEach(sibJobImpl -> { + sibJobsIds.add(sibJobImpl.jobId); + sibJobsSesId.add(sibJobImpl.sesId); + }); + } } /** @@ -343,8 +351,15 @@ public long getCreateTime() { /** * @return Sibling job ids. */ - public @Nullable Collection siblingJobsIds() { - return sibJobIds; + public @Nullable List siblingJobsIds() { + return sibJobsIds; + } + + /** + * @return Sibling job session ids. + */ + public @Nullable List siblingJobsSessionIds() { + return sibJobsSesId; } /** 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 32c350a3bab8c..552710234f8d9 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 @@ -40,10 +40,10 @@ */ public class GridJobSiblingImpl implements ComputeJobSibling { /** */ - private IgniteUuid sesId; + IgniteUuid sesId; /** */ - private final IgniteUuid jobId; + final IgniteUuid jobId; /** */ private Object taskTopic; 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 efa3fdca88584..b2fbacd19c290 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 @@ -18,9 +18,11 @@ package org.apache.ignite.internal.processors.job; import java.util.AbstractCollection; +import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.Iterator; +import java.util.List; import java.util.Map; import java.util.NoSuchElementException; import java.util.Objects; @@ -35,7 +37,6 @@ 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; @@ -1259,10 +1260,21 @@ public void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque } // TODO : Revise after https://issues.apache.org/jira/browse/IGNITE-28964 - Collection siblings = F.isEmpty(req.siblingJobsIds()) - ? null - : req.siblingJobsIds().stream().map(sibJobId -> new GridJobSiblingImpl(null, sibJobId, node.id(), ctx)) - .collect(Collectors.toList()); + List siblJobsIds = req.siblingJobsIds(); + List siblJobsSesIds = req.siblingJobsSessionIds(); + + assert F.isEmpty(siblJobsIds) == F.isEmpty(siblJobsSesIds); + + Collection siblings = null; + + if (!F.isEmpty(siblJobsIds)) { + assert siblJobsSesIds.size() == siblJobsIds.size(); + + siblings = new ArrayList<>(siblJobsIds.size()); + + for (int i = 0; i < siblJobsIds.size(); ++i) + siblings.add(new GridJobSiblingImpl(siblJobsSesIds.get(i), siblJobsIds.get(i), node.id(), ctx)); + } GridTaskSessionImpl taskSes = ctx.session().createTaskSession( req.sessionId(), 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 0f262d928f5a9..369023bbe7620 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 @@ -1392,7 +1392,7 @@ private void sendRequest(ComputeJobResult res) { timeout, ses.getTopology(), ses.getTopologyPredicate(), - ses.getJobSiblings(), + downcast(ses.getJobSiblings()), sesAttrs, jobAttrs, ses.getCheckpointSpi(), @@ -1473,6 +1473,19 @@ else if (log.isDebugEnabled()) } } + /** + * TODO : Revise after https://issues.apache.org/jira/browse/IGNITE-28964 + * Downcasts collection type. + * + * @param

Parent type. + * @param Child type. + * @param p Initial collection. + * @return Resulting collection.downcast + */ + private static Collection downcast(Collection

p) { + return (Collection)p; + } + /** * @param nodeId Node ID. */ From 5eae71d9c5c7939be2b18079c32f8c7fafe35ec1 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Sun, 9 Aug 2026 16:33:34 +0300 Subject: [PATCH 8/8] fix --- .../java/org/apache/ignite/internal/GridJobSiblingImpl.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 552710234f8d9..3f4cbb73ca1f1 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 @@ -27,7 +27,6 @@ import org.apache.ignite.internal.util.typedef.internal.S; import org.apache.ignite.internal.util.typedef.internal.U; import org.apache.ignite.lang.IgniteUuid; -import org.jetbrains.annotations.Nullable; import static org.apache.ignite.internal.GridTopic.TOPIC_JOB; import static org.apache.ignite.internal.GridTopic.TOPIC_JOB_CANCEL; @@ -66,7 +65,8 @@ public class GridJobSiblingImpl implements ComputeJobSibling { * @param nodeId ID of the node where this sibling was sent for execution. * @param ctx Managers registry. */ - public GridJobSiblingImpl(@Nullable IgniteUuid sesId, IgniteUuid jobId, UUID nodeId, GridKernalContext ctx) { + public GridJobSiblingImpl(IgniteUuid sesId, IgniteUuid jobId, UUID nodeId, GridKernalContext ctx) { + assert sesId != null; assert jobId != null; assert nodeId != null; assert ctx != null;