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..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,15 +19,16 @@ 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; 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; 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,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 siblings; - - /** */ + /** 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(); @@ -215,7 +212,7 @@ public GridJobExecuteRequest( long timeout, @Nullable Collection top, @Nullable IgnitePredicate topPred, - Collection siblings, + Collection siblings, Map sesAttrs, Map jobAttrs, String cpSpi, @@ -253,7 +250,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 +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(); } /** @@ -337,10 +336,10 @@ public long getCreateTime() { } /** - * @return Job siblings. + * @return Sibling job ids. */ - public Collection getSiblings() { - return siblings; + 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 5004aa946dfa2..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 @@ -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 @@ /** * 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 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..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 b81a8b43bf21e..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.getSiblings(), + 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 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, 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..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 @@ -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,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