From c84e59aab807d373bfc4b670600351a5a09ce21d Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Sat, 8 Aug 2026 15:05:47 +0300 Subject: [PATCH 1/3] impl --- .../internal/GridJobExecuteRequest.java | 15 +++---- .../ignite/internal/GridJobSiblingImpl.java | 23 +--------- .../ignite/internal/GridTaskSessionImpl.java | 43 ++++++++++++++++--- .../processors/job/GridJobProcessor.java | 2 +- .../session/GridTaskSessionProcessor.java | 9 ++-- .../processors/task/GridTaskWorker.java | 6 ++- 6 files changed, 56 insertions(+), 42 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 098e8514f681c..408fd9b684554 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 @@ -23,7 +23,6 @@ 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; @@ -109,7 +108,7 @@ public class GridJobExecuteRequest implements ExecutorAwareMessage, DeferredUnma /** Left unset for a continuous task: such a job requests its siblings from the task node instead. */ @Marshalled("siblingsBytes") - Collection siblings; + Collection siblingJobsIds; /** */ @Order(12) @@ -188,7 +187,7 @@ public GridJobExecuteRequest() { * @param timeout Task execution timeout. * @param top Topology. * @param topPred Topology predicate. - * @param siblings Collection of split siblings. + * @param siblingJobsIds Collection sibling jobs ids. * @param sesAttrs Session attributes. * @param jobAttrs Job attributes. * @param cpSpi Collision SPI. @@ -215,7 +214,7 @@ public GridJobExecuteRequest( long timeout, @Nullable Collection top, @Nullable IgnitePredicate topPred, - Collection siblings, + @Nullable Collection siblingJobsIds, Map sesAttrs, Map jobAttrs, String cpSpi, @@ -253,7 +252,7 @@ public GridJobExecuteRequest( this.top = top; this.topVer = topVer; this.topPred = topPred; - this.siblings = dynamicSiblings ? null : siblings; + this.siblingJobsIds = dynamicSiblings ? null : siblingJobsIds; this.sesAttrs = sesAttrs; this.jobAttrs = jobAttrs; this.clsLdrId = clsLdrId; @@ -337,10 +336,10 @@ public long getCreateTime() { } /** - * @return Job siblings. + * @return Siblings jobs ids. */ - public Collection getSiblings() { - return siblings; + public Collection 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 5004aa946dfa2..d79ed77f5c266 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,10 +36,7 @@ /** * This class provides implementation for job sibling. */ -public class GridJobSiblingImpl implements ComputeJobSibling, Externalizable { - /** */ - private static final long serialVersionUID = 0L; - +public class GridJobSiblingImpl implements ComputeJobSibling { /** */ private IgniteUuid sesId; @@ -173,20 +166,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); 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..3b658956cab7d 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 @@ -25,6 +25,7 @@ import java.util.Map; import java.util.UUID; import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; import org.apache.ignite.IgniteCheckedException; import org.apache.ignite.IgniteException; import org.apache.ignite.cluster.ClusterNode; @@ -80,7 +81,7 @@ public class GridTaskSessionImpl implements GridTaskSessionInternal { private final GridKernalContext ctx; /** */ - private Collection siblings; + private @Nullable Collection siblings; /** Guarded by {@link #mux}. */ private Map attrs; @@ -148,7 +149,7 @@ public class GridTaskSessionImpl implements GridTaskSessionInternal { * @param topPred Topology predicate. * @param startTime Task execution start time. * @param endTime Task execution end time. - * @param siblings Collection of siblings. + * @param siblingJobsIds Collection of sibling jobs ids. * @param attrs Session attributes. * @param ctx Grid Kernal Context. * @param fullSup Session full support enabled flag. @@ -166,7 +167,7 @@ public GridTaskSessionImpl( @Nullable IgnitePredicate topPred, long startTime, long endTime, - Collection siblings, + @Nullable Collection siblingJobsIds, @Nullable Map attrs, GridKernalContext ctx, boolean fullSup, @@ -191,7 +192,7 @@ public GridTaskSessionImpl( this.sesId = sesId; this.startTime = startTime; this.endTime = endTime; - this.siblings = siblings != null ? unmodifiableCollection(siblings) : null; + this.siblings = localSiblingsWrap(siblingJobsIds); this.ctx = ctx; if (attrs != null && !attrs.isEmpty()) { @@ -209,6 +210,16 @@ public GridTaskSessionImpl( this.secCtx = secCtx; } + /** + * Creates local representation of {@link ComputeJobSibling}s. + * + * @see LocalComputeJobSiblingWrap + */ + private static @Nullable Collection localSiblingsWrap(@Nullable Collection siblingJobsIds) { + return F.isEmpty(siblingJobsIds) ? null : siblingJobsIds.stream().map(LocalComputeJobSiblingWrap::new) + .collect(Collectors.toList()); + } + /** {@inheritDoc} */ @Override public boolean isFullSupport() { return fullSup; @@ -535,7 +546,7 @@ public void setClassLoader(ClassLoader clsLdr) { } /** {@inheritDoc} */ - @Override public Collection getJobSiblings() { + @Override public @Nullable Collection getJobSiblings() { synchronized (mux) { return siblings; } @@ -989,4 +1000,26 @@ public Object login() { @Override public String toString() { return S.toString(GridTaskSessionImpl.class, this); } + + /** A {@link ComputeJobSibling} providing {@link ComputeJobSibling#getJobId()} and unable to {@link ComputeJobSibling#cancel()}. */ + private static class LocalComputeJobSiblingWrap implements ComputeJobSibling { + /** */ + private final IgniteUuid jobId; + + /** */ + private LocalComputeJobSiblingWrap(IgniteUuid id) { + jobId = id; + } + + /** {@inheritDoc} */ + @Override public IgniteUuid getJobId() { + return jobId; + } + + /** {@inheritDoc} */ + @Override public void cancel() throws IgniteException { + throw new IgniteException(new UnsupportedOperationException("Cancelation of sibling compute jobs is allowed " + + "only where it started.")); + } + } } 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..2c15d3bc063fc 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 @@ -1266,7 +1266,7 @@ public void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque req.getTopologyPredicate(), req.startTaskTime(), endTime, - req.getSiblings(), + req.siblingJobsIds(), req.getSessionAttributes(), req.sessionFullSupport(), req.internal(), 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..4ea8c4d2c5fae 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 @@ -24,7 +24,6 @@ import java.util.concurrent.ConcurrentMap; import org.apache.ignite.IgniteCheckedException; import org.apache.ignite.cluster.ClusterNode; -import org.apache.ignite.compute.ComputeJobSibling; import org.apache.ignite.internal.GridKernalContext; import org.apache.ignite.internal.GridTaskSessionImpl; import org.apache.ignite.internal.managers.deployment.GridDeployment; @@ -77,7 +76,7 @@ public GridTaskSessionProcessor(GridKernalContext ctx) { * @param topPred Topology predicate. * @param startTime Execution start time. * @param endTime Execution end time. - * @param siblings Collection of siblings. + * @param siblingJobsIds Collection of sibling jons ids. * @param attrs Map of attributes. * @param fullSup {@code True} to enable distributed session attributes and checkpoints. * @param internal {@code True} in case of internal task. @@ -95,7 +94,7 @@ public GridTaskSessionImpl createTaskSession( @Nullable IgnitePredicate topPred, long startTime, long endTime, - Collection siblings, + @Nullable Collection siblingJobsIds, Map attrs, boolean fullSup, boolean internal, @@ -113,7 +112,7 @@ public GridTaskSessionImpl createTaskSession( topPred, startTime, endTime, - siblings, + siblingJobsIds, attrs, ctx, false, @@ -139,7 +138,7 @@ public GridTaskSessionImpl createTaskSession( topPred, startTime, endTime, - siblings, + siblingJobsIds, attrs, ctx, true, 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..5ec5e4726b11e 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 @@ -1381,6 +1381,10 @@ private void sendRequest(ComputeJobResult res) { boolean forceLocDep = internal || !ctx.deploy().enabled(); + Collection siblingJobsIds = F.isEmpty(ses.getJobSiblings()) + ? null + : ses.getJobSiblings().stream().map(ComputeJobSibling::getJobId).toList(); + req = new GridJobExecuteRequest( ses.getId(), res.getJobContext().getJobId(), @@ -1392,7 +1396,7 @@ private void sendRequest(ComputeJobResult res) { timeout, ses.getTopology(), ses.getTopologyPredicate(), - ses.getJobSiblings(), + siblingJobsIds, sesAttrs, jobAttrs, ses.getCheckpointSpi(), From 5fccee3540b4ecd23c236bc74832b14460edf245 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Sat, 8 Aug 2026 19:22:30 +0300 Subject: [PATCH 2/3] fix --- .../org/apache/ignite/internal/GridTaskSessionImpl.java | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) 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 3b658956cab7d..9a193e1ea7bb6 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 @@ -81,7 +81,7 @@ public class GridTaskSessionImpl implements GridTaskSessionInternal { private final GridKernalContext ctx; /** */ - private @Nullable Collection siblings; + private Collection siblings; /** Guarded by {@link #mux}. */ private Map attrs; @@ -215,8 +215,8 @@ public GridTaskSessionImpl( * * @see LocalComputeJobSiblingWrap */ - private static @Nullable Collection localSiblingsWrap(@Nullable Collection siblingJobsIds) { - return F.isEmpty(siblingJobsIds) ? null : siblingJobsIds.stream().map(LocalComputeJobSiblingWrap::new) + private static Collection localSiblingsWrap(@Nullable Collection siblingJobsIds) { + return F.isEmpty(siblingJobsIds) ? Collections.emptyList() : siblingJobsIds.stream().map(LocalComputeJobSiblingWrap::new) .collect(Collectors.toList()); } @@ -546,7 +546,7 @@ public void setClassLoader(ClassLoader clsLdr) { } /** {@inheritDoc} */ - @Override public @Nullable Collection getJobSiblings() { + @Override public Collection getJobSiblings() { synchronized (mux) { return siblings; } From 8529126d80a3a3a2911ad13a064e445c3c617937 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Sun, 9 Aug 2026 16:50:03 +0300 Subject: [PATCH 3/3] reimpl research --- .../internal/GridJobExecuteRequest.java | 24 ++++++------ .../ignite/internal/GridJobSiblingImpl.java | 12 ++---- .../ignite/internal/GridTaskSessionImpl.java | 39 ++----------------- .../processors/job/GridJobProcessor.java | 30 +++++++++++++- .../session/GridTaskSessionProcessor.java | 9 +++-- .../processors/task/GridTaskWorker.java | 18 ++++++--- .../resources/META-INF/classnames.properties | 1 - 7 files changed, 65 insertions(+), 68 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 408fd9b684554..6ba1b7f777db4 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 @@ -19,6 +19,7 @@ import java.io.Serializable; import java.util.Collection; +import java.util.List; import java.util.Map; import java.util.UUID; import org.apache.ignite.cluster.ClusterNode; @@ -27,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; @@ -106,13 +108,9 @@ 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 siblingJobsIds; - - /** */ + /** Sibling jobs ids. Plain representation of {@link GridJobSiblingImpl#jobId} to reduce the messages number. */ @Order(12) - byte[] siblingsBytes; + @Nullable List sibJobsIds; /** Transient since needs to hold local creation time. */ private final long createTime = U.currentTimeMillis(); @@ -187,7 +185,7 @@ public GridJobExecuteRequest() { * @param timeout Task execution timeout. * @param top Topology. * @param topPred Topology predicate. - * @param siblingJobsIds Collection sibling jobs ids. + * @param siblings Collection of split siblings. * @param sesAttrs Session attributes. * @param jobAttrs Job attributes. * @param cpSpi Collision SPI. @@ -214,7 +212,7 @@ public GridJobExecuteRequest( long timeout, @Nullable Collection top, @Nullable IgnitePredicate topPred, - @Nullable Collection siblingJobsIds, + Collection siblings, Map sesAttrs, Map jobAttrs, String cpSpi, @@ -252,7 +250,6 @@ public GridJobExecuteRequest( this.top = top; this.topVer = topVer; this.topPred = topPred; - this.siblingJobsIds = dynamicSiblings ? null : siblingJobsIds; this.sesAttrs = sesAttrs; this.jobAttrs = jobAttrs; this.clsLdrId = clsLdrId; @@ -268,6 +265,9 @@ public GridJobExecuteRequest( this.execName = execName; this.cpSpi = cpSpi == null || cpSpi.isEmpty() ? null : cpSpi; + + if (!dynamicSiblings && !F.isEmpty(siblings)) + sibJobsIds = siblings.stream().map(GridJobSiblingImpl::getJobId).toList(); } /** @@ -336,10 +336,10 @@ public long getCreateTime() { } /** - * @return Siblings jobs ids. + * @return Sibling job ids. */ - public Collection siblingJobsIds() { - return siblingJobsIds; + public @Nullable List siblingJobsIds() { + return sibJobsIds; } /** 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 d79ed77f5c266..3c20f5c84ba66 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 @@ -38,11 +38,10 @@ */ public class GridJobSiblingImpl implements ComputeJobSibling { /** */ - private IgniteUuid sesId; + IgniteUuid sesId; /** */ - @SuppressWarnings({"FieldAccessedSynchronizedAndUnsynchronized"}) - private IgniteUuid jobId; + final IgniteUuid jobId; /** */ private Object taskTopic; @@ -57,12 +56,7 @@ public class GridJobSiblingImpl implements ComputeJobSibling { private boolean isJobDone; /** */ - private transient GridKernalContext ctx; - - /** */ - public GridJobSiblingImpl() { - // No-op. - } + private GridKernalContext ctx; /** * @param sesId Task session ID. 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 9a193e1ea7bb6..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 @@ -25,7 +25,6 @@ import java.util.Map; import java.util.UUID; import java.util.concurrent.atomic.AtomicInteger; -import java.util.stream.Collectors; import org.apache.ignite.IgniteCheckedException; import org.apache.ignite.IgniteException; import org.apache.ignite.cluster.ClusterNode; @@ -149,7 +148,7 @@ public class GridTaskSessionImpl implements GridTaskSessionInternal { * @param topPred Topology predicate. * @param startTime Task execution start time. * @param endTime Task execution end time. - * @param siblingJobsIds Collection of sibling jobs ids. + * @param siblings Collection of siblings. * @param attrs Session attributes. * @param ctx Grid Kernal Context. * @param fullSup Session full support enabled flag. @@ -167,7 +166,7 @@ public GridTaskSessionImpl( @Nullable IgnitePredicate topPred, long startTime, long endTime, - @Nullable Collection siblingJobsIds, + @Nullable Collection siblings, @Nullable Map attrs, GridKernalContext ctx, boolean fullSup, @@ -192,7 +191,7 @@ public GridTaskSessionImpl( this.sesId = sesId; this.startTime = startTime; this.endTime = endTime; - this.siblings = localSiblingsWrap(siblingJobsIds); + this.siblings = siblings != null ? unmodifiableCollection(siblings) : null; this.ctx = ctx; if (attrs != null && !attrs.isEmpty()) { @@ -210,16 +209,6 @@ public GridTaskSessionImpl( this.secCtx = secCtx; } - /** - * Creates local representation of {@link ComputeJobSibling}s. - * - * @see LocalComputeJobSiblingWrap - */ - private static Collection localSiblingsWrap(@Nullable Collection siblingJobsIds) { - return F.isEmpty(siblingJobsIds) ? Collections.emptyList() : siblingJobsIds.stream().map(LocalComputeJobSiblingWrap::new) - .collect(Collectors.toList()); - } - /** {@inheritDoc} */ @Override public boolean isFullSupport() { return fullSup; @@ -1000,26 +989,4 @@ public Object login() { @Override public String toString() { return S.toString(GridTaskSessionImpl.class, this); } - - /** A {@link ComputeJobSibling} providing {@link ComputeJobSibling#getJobId()} and unable to {@link ComputeJobSibling#cancel()}. */ - private static class LocalComputeJobSiblingWrap implements ComputeJobSibling { - /** */ - private final IgniteUuid jobId; - - /** */ - private LocalComputeJobSiblingWrap(IgniteUuid id) { - jobId = id; - } - - /** {@inheritDoc} */ - @Override public IgniteUuid getJobId() { - return jobId; - } - - /** {@inheritDoc} */ - @Override public void cancel() throws IgniteException { - throw new IgniteException(new UnsupportedOperationException("Cancelation of sibling compute jobs is allowed " + - "only where it started.")); - } - } } 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 2c15d3bc063fc..ce51761387470 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 @@ -35,6 +35,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; @@ -1256,6 +1257,11 @@ public void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque U.resolveClassLoader(dep.classLoader(), ctx.config())); } + Collection siblings = null; + + if (!F.isEmpty(req.siblingJobsIds())) + siblings = req.siblingJobsIds().stream().map(LocalComputeJobSibling::new).collect(Collectors.toList()); + GridTaskSessionImpl taskSes = ctx.session().createTaskSession( req.sessionId(), node.id(), @@ -1266,7 +1272,7 @@ public void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque req.getTopologyPredicate(), req.startTaskTime(), endTime, - req.siblingJobsIds(), + siblings, req.getSessionAttributes(), req.sessionFullSupport(), req.internal(), @@ -2507,4 +2513,26 @@ else if (jobId == null) public long computeJobWorkerInterruptTimeout() { return computeJobWorkerInterruptTimeout.getOrDefault(ctx.config().getFailureDetectionTimeout()); } + + /** */ + private static class LocalComputeJobSibling implements ComputeJobSibling { + /** */ + private final IgniteUuid sibJobId; + + /** */ + private LocalComputeJobSibling(IgniteUuid sibJobId) { + this.sibJobId = sibJobId; + } + + /** {@inheritDoc} */ + @Override public IgniteUuid getJobId() { + return sibJobId; + } + + /** {@inheritDoc} */ + @Override public void cancel() throws IgniteException { + throw new IgniteException(new UnsupportedOperationException("Cancelation of sibling compute jobs are allowed " + + "only where it started.")); + } + } } 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 4ea8c4d2c5fae..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 @@ -24,6 +24,7 @@ import java.util.concurrent.ConcurrentMap; import org.apache.ignite.IgniteCheckedException; import org.apache.ignite.cluster.ClusterNode; +import org.apache.ignite.compute.ComputeJobSibling; import org.apache.ignite.internal.GridKernalContext; import org.apache.ignite.internal.GridTaskSessionImpl; import org.apache.ignite.internal.managers.deployment.GridDeployment; @@ -76,7 +77,7 @@ public GridTaskSessionProcessor(GridKernalContext ctx) { * @param topPred Topology predicate. * @param startTime Execution start time. * @param endTime Execution end time. - * @param siblingJobsIds Collection of sibling jons ids. + * @param siblings Collection of siblings. * @param attrs Map of attributes. * @param fullSup {@code True} to enable distributed session attributes and checkpoints. * @param internal {@code True} in case of internal task. @@ -94,7 +95,7 @@ public GridTaskSessionImpl createTaskSession( @Nullable IgnitePredicate topPred, long startTime, long endTime, - @Nullable Collection siblingJobsIds, + @Nullable Collection siblings, Map attrs, boolean fullSup, boolean internal, @@ -112,7 +113,7 @@ public GridTaskSessionImpl createTaskSession( topPred, startTime, endTime, - siblingJobsIds, + siblings, attrs, ctx, false, @@ -138,7 +139,7 @@ public GridTaskSessionImpl createTaskSession( topPred, startTime, endTime, - siblingJobsIds, + siblings, attrs, ctx, true, 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 5ec5e4726b11e..de70559c16aab 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 @@ -1381,10 +1381,6 @@ private void sendRequest(ComputeJobResult res) { boolean forceLocDep = internal || !ctx.deploy().enabled(); - Collection siblingJobsIds = F.isEmpty(ses.getJobSiblings()) - ? null - : ses.getJobSiblings().stream().map(ComputeJobSibling::getJobId).toList(); - req = new GridJobExecuteRequest( ses.getId(), res.getJobContext().getJobId(), @@ -1396,7 +1392,7 @@ private void sendRequest(ComputeJobResult res) { timeout, ses.getTopology(), ses.getTopologyPredicate(), - siblingJobsIds, + downcast(ses.getJobSiblings()), sesAttrs, jobAttrs, ses.getCheckpointSpi(), @@ -1477,6 +1473,18 @@ else if (log.isDebugEnabled()) } } + /** + * 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. */ 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