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..29bc0721735bc 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,16 +18,18 @@ 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; 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 +109,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(); @@ -188,7 +186,7 @@ public GridJobExecuteRequest() { * @param timeout Task execution timeout. * @param top Topology. * @param topPred Topology predicate. - * @param siblings Collection of split siblings. + * @param siblingJobsIds Collection of sibling jobs ids. * @param sesAttrs Session attributes. * @param jobAttrs Job attributes. * @param cpSpi Collision SPI. @@ -215,7 +213,7 @@ public GridJobExecuteRequest( long timeout, @Nullable Collection top, @Nullable IgnitePredicate topPred, - Collection siblings, + @Nullable Collection siblingJobsIds, Map sesAttrs, Map jobAttrs, String cpSpi, @@ -253,7 +251,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 +266,9 @@ public GridJobExecuteRequest( this.execName = execName; this.cpSpi = cpSpi == null || cpSpi.isEmpty() ? null : cpSpi; + + if (!dynamicSiblings && !F.isEmpty(siblingJobsIds)) + sibJobsIds = new ArrayList<>(siblingJobsIds); } /** @@ -337,10 +337,10 @@ public long getCreateTime() { } /** - * @return Job siblings. + * @return Sibling jobs 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..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 @@ -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; @@ -39,17 +35,14 @@ /** * This class provides implementation for job sibling. + * TODO : Revise after https://issues.apache.org/jira/browse/IGNITE-28964 */ -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 +57,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 +161,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..223def03c5d15 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 @@ -1176,11 +1176,15 @@ private void updateJobMetrics0() { /** * @param node Node. * @param req Request. + * @param locJobSiblings Local siblings of the job. 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 locJobSiblings + ) { 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 +1199,7 @@ public void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque 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; } @@ -1252,6 +1256,7 @@ public void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque try { // The job payload waits for this point: only now is there a deployment to unmarshal it with. if (!loc) { + // TODO : Revise in https://issues.apache.org/jira/browse/IGNITE-28964 MessageMarshalling.unmarshal(req, ctx, null, U.resolveClassLoader(dep.classLoader(), ctx.config())); } @@ -1266,7 +1271,7 @@ public void processJobExecuteRequest(ClusterNode node, final GridJobExecuteReque req.getTopologyPredicate(), req.startTaskTime(), endTime, - req.getSiblings(), + locJobSiblings, req.getSessionAttributes(), req.sessionFullSupport(), req.internal(), @@ -2208,7 +2213,7 @@ private class JobExecutionListener implements GridMessageListener { assert node != null; - processJobExecuteRequest(node, (GridJobExecuteRequest)msg); + processJobExecuteRequest(node, (GridJobExecuteRequest)msg, null); } } 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..4b0cd31dd0a6e 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(), + F.isEmpty(ses.getJobSiblings()) ? null : ses.getJobSiblings().stream().map(ComputeJobSibling::getJobId).toList(), sesAttrs, jobAttrs, ses.getCheckpointSpi(), @@ -1409,7 +1409,7 @@ private void sendRequest(ComputeJobResult res) { 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 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