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

Reply via email to