This is an automated email from the ASF dual-hosted git repository.

shishkovilja 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 f83f8a55305 IGNITE-28790 Use MessageSerializer to transfer continuous 
query DTO's in DataBag. (#13426)
f83f8a55305 is described below

commit f83f8a55305d36c7a25dd5cdeb88621763611b7b
Author: Vladimir Steshin <[email protected]>
AuthorDate: Wed Aug 12 17:02:08 2026 +0300

    IGNITE-28790 Use MessageSerializer to transfer continuous query DTO's in 
DataBag. (#13426)
---
 .../ignite/internal/CoreMessagesProvider.java      |  18 +
 .../ignite/internal/GridEventConsumeHandler.java   |  76 ++--
 .../ignite/internal/GridMessageListenHandler.java  | 113 +++---
 .../deployment/GridDeploymentInfoMessage.java      |   2 +-
 .../processors/affinity/GridAffinityUtils.java     |   2 +-
 .../CacheContinuousQueryDeployableObject.java      |  37 +-
 .../continuous/CacheContinuousQueryHandler.java    | 333 ++++++++----------
 .../continuous/ContinousRoutineDiscoveryData.java  |  72 ++++
 .../ContinousRoutineDiscoveryDataItem.java         |  96 +++++
 .../continuous/ContinousRoutineLocalInfo.java      | 133 +++++++
 .../continuous/ContinuousRoutineInfo.java          |  44 ++-
 .../ContinuousRoutinesCommonDiscoveryData.java     |  14 +-
 ...ContinuousRoutinesJoiningNodeDiscoveryData.java |  14 +-
 .../continuous/GridContinuousHandler.java          |   4 +-
 .../continuous/GridContinuousProcessor.java        | 387 +++------------------
 .../processors/continuous/StartRequestData.java    |   5 +-
 .../org/apache/ignite/spi/IgniteSpiAdapter.java    |   4 +-
 .../ignite/spi/discovery/DiscoveryDataBag.java     |   8 -
 .../ignite/spi/discovery/tcp/ClientImpl.java       |   2 +-
 .../ignite/spi/discovery/tcp/ServerImpl.java       |   4 +-
 .../ignite/spi/discovery/tcp/TcpDiscoverySpi.java  |  13 +
 .../CacheContinuousQueryEntriesExpireTest.java     |   4 +-
 ...ueryRemoteFilterMissingInClassPathSelfTest.java |  53 +--
 ...coveryDataDeserializationFailureHanderTest.java |   1 +
 ...GridCacheContinuousQueryNodesFilteringTest.java |  15 +-
 .../continuous/GridEventConsumeSelfTest.java       |   9 +-
 .../zk/internal/ZookeeperDiscoveryImpl.java        |   2 +-
 27 files changed, 729 insertions(+), 736 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 000cbd47f02..d95ce2dab4e 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
@@ -178,7 +178,9 @@ import 
org.apache.ignite.internal.processors.cache.query.GridCacheQueryRequest;
 import 
org.apache.ignite.internal.processors.cache.query.GridCacheQueryResponse;
 import org.apache.ignite.internal.processors.cache.query.GridCacheSqlQuery;
 import 
org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryBatchAck;
+import 
org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryDeployableObject;
 import 
org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryEntry;
+import 
org.apache.ignite.internal.processors.cache.query.continuous.CacheContinuousQueryHandler;
 import org.apache.ignite.internal.processors.cache.transactions.IgniteTxEntry;
 import org.apache.ignite.internal.processors.cache.transactions.IgniteTxKey;
 import 
org.apache.ignite.internal.processors.cache.transactions.TxEntryValueHolder;
@@ -199,7 +201,13 @@ import 
org.apache.ignite.internal.processors.cluster.ClusterMetricsUpdateMessage
 import 
org.apache.ignite.internal.processors.cluster.ClusterUpdateNotifierDataBagItem;
 import org.apache.ignite.internal.processors.cluster.DiscoveryDataClusterState;
 import org.apache.ignite.internal.processors.cluster.NodeFullMetricsMessage;
+import 
org.apache.ignite.internal.processors.continuous.ContinousRoutineDiscoveryData;
+import 
org.apache.ignite.internal.processors.continuous.ContinousRoutineDiscoveryDataItem;
+import 
org.apache.ignite.internal.processors.continuous.ContinousRoutineLocalInfo;
+import org.apache.ignite.internal.processors.continuous.ContinuousRoutineInfo;
 import 
org.apache.ignite.internal.processors.continuous.ContinuousRoutineStartResultMessage;
+import 
org.apache.ignite.internal.processors.continuous.ContinuousRoutinesCommonDiscoveryData;
+import 
org.apache.ignite.internal.processors.continuous.ContinuousRoutinesJoiningNodeDiscoveryData;
 import org.apache.ignite.internal.processors.continuous.GridContinuousMessage;
 import org.apache.ignite.internal.processors.continuous.StartRequestData;
 import 
org.apache.ignite.internal.processors.continuous.StartRoutineAckDiscoveryMessage;
@@ -627,6 +635,16 @@ public class CoreMessagesProvider extends 
AbstractMessageFactoryProvider {
         register(QueryProposalsDataBagItem.class);
         register(QueryEntityMessage.class);
         register(QueryEntityExMessage.class);
+        register(ContinuousRoutineInfo.class);
+        register(ContinuousRoutinesJoiningNodeDiscoveryData.class);
+        register(CacheContinuousQueryDeployableObject.class);
+        register(CacheContinuousQueryHandler.class);
+        register(GridEventConsumeHandler.class);
+        register(GridMessageListenHandler.class);
+        register(ContinousRoutineLocalInfo.class);
+        register(ContinousRoutineDiscoveryDataItem.class);
+        register(ContinousRoutineDiscoveryData.class);
+        register(ContinuousRoutinesCommonDiscoveryData.class);
 
         // [11200 - 11300]: Compute, distributed process messages.
         msgIdx = 11200;
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java
index 8ebea159666..af4d85c204f 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/GridEventConsumeHandler.java
@@ -61,12 +61,9 @@ import static org.apache.ignite.events.EventType.EVTS_ALL;
 /**
  * Continuous routine handler for remote event listening.
  */
-class GridEventConsumeHandler implements GridContinuousHandler {
-    /** */
-    private static final long serialVersionUID = 0L;
-
+public final class GridEventConsumeHandler implements GridContinuousHandler, 
MarshallableMessage {
     /** Default callback. */
-    private static final IgniteBiPredicate<UUID, Event> DFLT_CALLBACK = new 
P2<UUID, Event>() {
+    private static final IgniteBiPredicate<UUID, Event> DFLT_CALLBACK = new 
P2<>() {
         @Override public boolean apply(UUID uuid, Event e) {
             return true;
         }
@@ -76,28 +73,32 @@ class GridEventConsumeHandler implements 
GridContinuousHandler {
     private IgniteBiPredicate<UUID, Event> cb;
 
     /** Filter. */
-    private IgnitePredicate<Event> filter;
+    @Nullable volatile IgnitePredicate<Event> filter;
 
-    /** Serialized filter. */
-    private byte[] filterBytes;
+    /** Marshaled {@link #filter}. */
+    @Order(0)
+    @Nullable volatile byte[] filterBytes;
 
-    /** Deployment class name. */
-    private String clsName;
+    /** Deployment class name. Is {@code null} if P2P deployment is disabled. 
*/
+    @Order(1)
+    @Nullable volatile String clsName;
 
-    /** Deployment info. */
-    private GridDeploymentInfo depInfo;
+    /** Deployment info. Is {@code null} if P2P deployment is disabled. */
+    @Order(2)
+    @Nullable volatile GridDeploymentInfoMessage depInfo;
 
     /** Types. */
-    private int[] types;
+    @Order(3)
+    int[] types;
 
     /** Listener. */
     private GridLocalEventListener lsnr;
 
     /** P2P unmarshalling future. */
-    private IgniteInternalFuture<Void> p2pUnmarshalFut = new 
GridFinishedFuture<>();
+    private volatile IgniteInternalFuture<Void> p2pUnmarshalFut = new 
GridFinishedFuture<>();
 
     /**
-     * Required by {@link Externalizable}.
+     * Empty constructor for serialization purposes.
      */
     public GridEventConsumeHandler() {
         // No-op.
@@ -225,8 +226,6 @@ class GridEventConsumeHandler implements 
GridContinuousHandler {
                                                 EventWrapper wrapper = new 
EventWrapper(evt);
 
                                                 if (evt instanceof CacheEvent) 
{
-                                                    String cacheName = 
((CacheEvent)evt).cacheName();
-
                                                     ClusterNode node = 
ctx.discovery().node(t3.get1());
 
                                                     if (node == null)
@@ -393,6 +392,10 @@ class GridEventConsumeHandler implements 
GridContinuousHandler {
         assert ctx != null;
         assert ctx.config().isPeerClassLoadingEnabled();
 
+        // TODO : Remove this check after 
https://issues.apache.org/jira/browse/IGNITE-28945
+        if (filterBytes != null)
+            return;
+
         if (filter != null) {
             Class cls = U.detectClass(filter);
 
@@ -415,6 +418,10 @@ class GridEventConsumeHandler implements 
GridContinuousHandler {
         assert ctx != null;
         assert ctx.config().isPeerClassLoadingEnabled();
 
+        // TODO : Remove this check after 
https://issues.apache.org/jira/browse/IGNITE-28945
+        if (filter != null)
+            return;
+
         if (filterBytes != null) {
             try {
                 GridDeployment dep = ctx.deploy().globalDeployment(depInfo, 
clsName, nodeId);
@@ -467,36 +474,27 @@ class GridEventConsumeHandler implements 
GridContinuousHandler {
     }
 
     /** {@inheritDoc} */
-    @Override public void writeExternal(ObjectOutput out) throws IOException {
-        boolean b = filterBytes != null;
+    @Override public void marshal(Marshaller marsh) throws 
IgniteCheckedException {
+        assert (clsName == null) == (depInfo == null);
 
-        out.writeBoolean(b);
-
-        if (b) {
-            U.writeByteArray(out, filterBytes);
-            U.writeString(out, clsName);
-            out.writeObject(depInfo);
-        }
-        else
-            out.writeObject(filter);
-
-        out.writeObject(types);
+        /** Are marshaled in {@link #p2pUnmarshal(UUID, GridKernalContext)}. */
+        if (filter != null && depInfo == null)
+            filterBytes = marsh.marshal(filter);
     }
 
     /** {@inheritDoc} */
-    @Override public void readExternal(ObjectInput in) throws IOException, 
ClassNotFoundException {
-        boolean b = in.readBoolean();
+    @Override public void unmarshal(Marshaller marsh, ClassLoader clsLdr) 
throws IgniteCheckedException {
+        assert (clsName == null) == (depInfo == null);
 
-        if (b) {
+        /** Are unmarshaled in {@link #p2pUnmarshal(UUID, GridKernalContext)}. 
*/
+        if (depInfo != null) {
             p2pUnmarshalFut = new GridFutureAdapter<>();
-            filterBytes = U.readByteArray(in);
-            clsName = U.readString(in);
-            depInfo = (GridDeploymentInfo)in.readObject();
+
+            return;
         }
-        else
-            filter = (IgnitePredicate<Event>)in.readObject();
 
-        types = (int[])in.readObject();
+        if (filterBytes != null)
+            filter = marsh.unmarshal(filterBytes, clsLdr);
     }
 
     /**
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java
index f3397e5d214..c1bd4a683c1 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/GridMessageListenHandler.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.Collections;
 import java.util.Map;
@@ -39,41 +35,40 @@ import 
org.apache.ignite.internal.util.lang.GridPeerDeployAware;
 import org.apache.ignite.internal.util.typedef.internal.S;
 import org.apache.ignite.internal.util.typedef.internal.U;
 import org.apache.ignite.lang.IgniteBiPredicate;
+import org.apache.ignite.marshaller.Marshaller;
 import org.jetbrains.annotations.Nullable;
 
 /**
  * Continuous handler for message subscription.
  */
-public class GridMessageListenHandler implements GridContinuousHandler {
+public final class GridMessageListenHandler implements GridContinuousHandler, 
MarshallableMessage {
     /** */
-    private static final long serialVersionUID = 0L;
+    private volatile @Nullable Object topic;
 
-    /** */
-    private Object topic;
+    /** Marshalled {@link #topic}. */
+    @Order(0)
+    @Nullable volatile byte[] topicBytes;
 
     /** */
-    private IgniteBiPredicate<UUID, Object> pred;
+    private volatile IgniteBiPredicate<UUID, Object> pred;
 
-    /** */
-    private byte[] topicBytes;
+    /** Marshalled {@link #pred}. */
+    @Order(1)
+    volatile byte[] predBytes;
 
-    /** */
-    private byte[] predBytes;
+    /** Class name of {@link #pred}. Is {@code null} if the P2P deployment is 
disabled. */
+    @Order(2)
+    @Nullable volatile String clsName;
 
-    /** */
-    private String clsName;
-
-    /** */
-    private GridDeploymentInfoMessage depInfo;
-
-    /** */
-    private boolean depEnabled;
+    /** P2P deploy info of {@link #pred}. Is {@code null} if the P2P 
deployment is disabled. */
+    @Order(3)
+    @Nullable volatile GridDeploymentInfoMessage predDepInfo;
 
     /** P2P unmarshalling future. */
-    private IgniteInternalFuture<Void> p2pUnmarshalFut = new 
GridFinishedFuture<>();
+    private volatile IgniteInternalFuture<Void> p2pUnmarshalFut = new 
GridFinishedFuture<>();
 
     /**
-     * Required by {@link Externalizable}.
+     * Empty constructor for serialization purposes
      */
     public GridMessageListenHandler() {
         // No-op.
@@ -151,6 +146,10 @@ public class GridMessageListenHandler implements 
GridContinuousHandler {
         assert ctx != null;
         assert ctx.config().isPeerClassLoadingEnabled();
 
+        // TODO : Remove this check after 
https://issues.apache.org/jira/browse/IGNITE-28945
+        if (predDepInfo != null)
+            return;
+
         if (topic != null)
             topicBytes = U.marshal(ctx.marshaller(), topic);
 
@@ -166,9 +165,19 @@ public class GridMessageListenHandler implements 
GridContinuousHandler {
         if (dep == null)
             throw new IgniteDeploymentCheckedException("Failed to deploy 
message listener.");
 
-        depInfo = new GridDeploymentInfoMessage(dep);
+        predDepInfo = new GridDeploymentInfoMessage(dep);
+    }
+
+    /** {@inheritDoc} */
+    @Override public void marshal(Marshaller marsh) throws 
IgniteCheckedException {
+        /** Are marshaled in {@link #p2pMarshal(GridKernalContext)}. */
+        if (predDepInfo != null)
+            return;
+
+        if (topic != null)
+            topicBytes = marsh.marshal(topic);
 
-        depEnabled = true;
+        predBytes = marsh.marshal(pred);
     }
 
     /** {@inheritDoc} */
@@ -177,8 +186,12 @@ public class GridMessageListenHandler implements 
GridContinuousHandler {
         assert ctx != null;
         assert ctx.config().isPeerClassLoadingEnabled();
 
+        // TODO : Remove this check after 
https://issues.apache.org/jira/browse/IGNITE-28945
+        if (pred != null)
+            return;
+
         try {
-            ClassLoader ldr = ctx.deploy().globalDeployment(depInfo, clsName, 
nodeId).classLoader();
+            ClassLoader ldr = ctx.deploy().globalDeployment(predDepInfo, 
clsName, nodeId).classLoader();
 
             if (topicBytes != null)
                 topic = U.unmarshal(ctx, topicBytes, U.resolveClassLoader(ldr, 
ctx.config()));
@@ -199,6 +212,21 @@ public class GridMessageListenHandler implements 
GridContinuousHandler {
         ((GridFutureAdapter)p2pUnmarshalFut).onDone();
     }
 
+    /** {@inheritDoc} */
+    @Override public void unmarshal(Marshaller marsh, ClassLoader clsLdr) 
throws IgniteCheckedException {
+        /** Are unmarshaled in {@link #p2pUnmarshal(UUID, GridKernalContext)}. 
*/
+        if (predDepInfo != null) {
+            p2pUnmarshalFut = new GridFutureAdapter<>();
+
+            return;
+        }
+
+        if (topicBytes != null)
+            topic = marsh.unmarshal(topicBytes, clsLdr);
+
+        pred = marsh.unmarshal(predBytes, clsLdr);
+    }
+
     /** {@inheritDoc} */
     @Override public GridContinuousBatch createBatch() {
         return new GridContinuousBatchAdapter();
@@ -229,39 +257,6 @@ public class GridMessageListenHandler implements 
GridContinuousHandler {
         }
     }
 
-    /** {@inheritDoc} */
-    @Override public void writeExternal(ObjectOutput out) throws IOException {
-        out.writeBoolean(depEnabled);
-
-        if (depEnabled) {
-            U.writeByteArray(out, topicBytes);
-            U.writeByteArray(out, predBytes);
-            U.writeString(out, clsName);
-            out.writeObject(depInfo);
-        }
-        else {
-            out.writeObject(topic);
-            out.writeObject(pred);
-        }
-    }
-
-    /** {@inheritDoc} */
-    @Override public void readExternal(ObjectInput in) throws IOException, 
ClassNotFoundException {
-        depEnabled = in.readBoolean();
-
-        if (depEnabled) {
-            p2pUnmarshalFut = new GridFutureAdapter<>();
-            topicBytes = U.readByteArray(in);
-            predBytes = U.readByteArray(in);
-            clsName = U.readString(in);
-            depInfo = (GridDeploymentInfoMessage)in.readObject();
-        }
-        else {
-            topic = in.readObject();
-            pred = (IgniteBiPredicate<UUID, Object>)in.readObject();
-        }
-    }
-
     /** {@inheritDoc} */
     @Override public String toString() {
         return S.toString(GridMessageListenHandler.class, this);
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoMessage.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoMessage.java
index 4a4f4e4a6a9..3c00a13799a 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoMessage.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/managers/deployment/GridDeploymentInfoMessage.java
@@ -30,7 +30,7 @@ import 
org.apache.ignite.plugin.extensions.communication.Message;
 /**
  * Deployment of classes, as it travels inside the messages carrying them.
  */
-public class GridDeploymentInfoMessage implements Message, GridDeploymentInfo, 
Serializable {
+public final class GridDeploymentInfoMessage implements Message, 
GridDeploymentInfo, Serializable {
     /** */
     private static final long serialVersionUID = 0L;
 
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/affinity/GridAffinityUtils.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/affinity/GridAffinityUtils.java
index 2eda9688193..7e0246b7897 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/affinity/GridAffinityUtils.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/affinity/GridAffinityUtils.java
@@ -86,7 +86,7 @@ class GridAffinityUtils {
     }
 
     /**
-     * Unmarshalls transfer object from remote node within a given context.
+     * Unmarshals transfer object from remote node within a given context.
      *
      * @param ctx Grid kernal context that provides deployment and marshalling 
services.
      * @param sndNodeId {@link UUID} of the sender node.
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryDeployableObject.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryDeployableObject.java
index 5f31963ba21..57843e86d1b 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryDeployableObject.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryDeployableObject.java
@@ -17,40 +17,37 @@
 
 package org.apache.ignite.internal.processors.cache.query.continuous;
 
-import java.io.Externalizable;
-import java.io.IOException;
-import java.io.ObjectInput;
-import java.io.ObjectOutput;
 import java.util.UUID;
 import org.apache.ignite.IgniteCheckedException;
 import org.apache.ignite.internal.GridKernalContext;
 import org.apache.ignite.internal.IgniteDeploymentCheckedException;
+import org.apache.ignite.internal.Order;
 import org.apache.ignite.internal.managers.deployment.GridDeployment;
-import org.apache.ignite.internal.managers.deployment.GridDeploymentInfo;
 import 
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
 import org.apache.ignite.internal.util.tostring.GridToStringExclude;
 import org.apache.ignite.internal.util.typedef.internal.S;
 import org.apache.ignite.internal.util.typedef.internal.U;
+import org.apache.ignite.plugin.extensions.communication.Message;
 
 /**
  * Deployable object.
  */
-class CacheContinuousQueryDeployableObject implements Externalizable {
-    /** */
-    private static final long serialVersionUID = 0L;
-
+public final class CacheContinuousQueryDeployableObject implements Message {
     /** Serialized object. */
     @GridToStringExclude
-    private byte[] bytes;
+    @Order(0)
+    byte[] bytes;
 
     /** Deployment class name. */
-    private String clsName;
+    @Order(1)
+    String clsName;
 
     /** Deployment info. */
-    private GridDeploymentInfo depInfo;
+    @Order(2)
+    GridDeploymentInfoMessage depInfo;
 
     /**
-     * Required by {@link Externalizable}.
+     * Empty constructor for serialization purposes.
      */
     public CacheContinuousQueryDeployableObject() {
         // No-op.
@@ -93,20 +90,6 @@ class CacheContinuousQueryDeployableObject implements 
Externalizable {
         return U.unmarshal(ctx, bytes, U.resolveClassLoader(dep.classLoader(), 
ctx.config()));
     }
 
-    /** {@inheritDoc} */
-    @Override public void writeExternal(ObjectOutput out) throws IOException {
-        U.writeByteArray(out, bytes);
-        U.writeString(out, clsName);
-        out.writeObject(depInfo);
-    }
-
-    /** {@inheritDoc} */
-    @Override public void readExternal(ObjectInput in) throws IOException, 
ClassNotFoundException {
-        bytes = U.readByteArray(in);
-        clsName = U.readString(in);
-        depInfo = (GridDeploymentInfo)in.readObject();
-    }
-
     /** {@inheritDoc} */
     @Override public String toString() {
         return S.toString(CacheContinuousQueryDeployableObject.class, this);
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java
index 805ebda1e00..ac541a10a41 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryHandler.java
@@ -17,10 +17,6 @@
 
 package org.apache.ignite.internal.processors.cache.query.continuous;
 
-import java.io.Externalizable;
-import java.io.IOException;
-import java.io.ObjectInput;
-import java.io.ObjectOutput;
 import java.util.ArrayList;
 import java.util.Collection;
 import java.util.Collections;
@@ -51,6 +47,9 @@ import org.apache.ignite.events.CacheQueryExecutedEvent;
 import org.apache.ignite.events.CacheQueryReadEvent;
 import org.apache.ignite.internal.GridKernalContext;
 import org.apache.ignite.internal.IgniteInternalFuture;
+import org.apache.ignite.internal.MarshallableMessage;
+import org.apache.ignite.internal.Marshalled;
+import org.apache.ignite.internal.Order;
 import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException;
 import org.apache.ignite.internal.managers.communication.GridIoPolicy;
 import org.apache.ignite.internal.managers.communication.MessageMarshalling;
@@ -83,6 +82,7 @@ import org.apache.ignite.internal.util.typedef.internal.U;
 import org.apache.ignite.lang.IgniteAsyncCallback;
 import org.apache.ignite.lang.IgniteBiTuple;
 import org.apache.ignite.lang.IgniteClosure;
+import org.apache.ignite.marshaller.Marshaller;
 import org.jetbrains.annotations.NotNull;
 import org.jetbrains.annotations.Nullable;
 
@@ -96,10 +96,7 @@ import static 
org.apache.ignite.internal.processors.cache.query.continuous.Cache
 /**
  * Continuous query handler.
  */
-public class CacheContinuousQueryHandler<K, V> implements 
GridContinuousHandler {
-    /** */
-    private static final long serialVersionUID = 0L;
-
+public final class CacheContinuousQueryHandler<K, V> implements 
GridContinuousHandler, MarshallableMessage {
     /** @see #IGNITE_CONTINUOUS_QUERY_BACKUP_ACK_THRESHOLD */
     public static final int DFLT_CONTINUOUS_QUERY_BACKUP_ACK_THRESHOLD = 100;
 
@@ -131,8 +128,8 @@ public class CacheContinuousQueryHandler<K, V> implements 
GridContinuousHandler
      * Transformer implementation for processing received remote events.
      * They are already transformed so we simply return transformed value for 
event.
      */
-    private transient IgniteClosure<CacheEntryEvent<? extends K, ? extends V>, 
?> returnValTrans =
-        new IgniteClosure<CacheEntryEvent<? extends K, ? extends V>, Object>() 
{
+    private IgniteClosure<CacheEntryEvent<? extends K, ? extends V>, ?> 
returnValTrans =
+        new IgniteClosure<>() {
             @Override public Object apply(CacheEntryEvent<? extends K, ? 
extends V> evt) {
                 assert evt.getKey() == null;
 
@@ -141,124 +138,153 @@ public class CacheContinuousQueryHandler<K, V> 
implements GridContinuousHandler
         };
 
     /** Cache name. */
-    private String cacheName;
+    @Order(0)
+    String cacheName;
 
     /** Topic for ordered messages. */
-    private Object topic;
+    @Marshalled("topicBytes")
+    Object topic;
+
+    /** Marshalled {@link #topic}. */
+    @Order(1)
+    byte[] topicBytes;
 
     /** P2P unmarshalling future. */
-    protected transient IgniteInternalFuture<Void> p2pUnmarshalFut = new 
GridFinishedFuture<>();
+    protected volatile IgniteInternalFuture<Void> p2pUnmarshalFut = new 
GridFinishedFuture<>();
 
     /** Initialization future. */
-    protected transient IgniteInternalFuture<Void> initFut;
+    protected IgniteInternalFuture<Void> initFut;
 
     /** Local listener. */
-    private transient CacheEntryUpdatedListener<K, V> locLsnr;
+    private CacheEntryUpdatedListener<K, V> locLsnr;
 
     /** Remote filter. */
-    private CacheEntryEventSerializableFilter<K, V> rmtFilter;
+    volatile CacheEntryEventSerializableFilter<K, V> rmtFilter;
 
-    /** Deployable object for filter. */
-    private CacheContinuousQueryDeployableObject rmtFilterDep;
+    /** Deployable object for {@link #rmtFilter}. Is {@code null} if no 
external marsshalling used. */
+    @Order(2)
+    @Nullable volatile CacheContinuousQueryDeployableObject rmtFilterDep;
+
+    /** Marshalled {@link #rmtFilter} if {@link #rmtFilterDep} is {@code 
null}. */
+    @Order(3)
+    @Nullable volatile byte[] rmtFilterBytes;
 
     /** Remote filter factory. */
-    private Factory<? extends CacheEntryEventFilter> rmtFilterFactory;
+    @Nullable volatile Factory<? extends CacheEntryEventFilter> 
rmtFilterFactory;
+
+    /** Deployable object for {@link #rmtFilterFactory}. Is {@code null} if no 
external marsshalling used. */
+    @Order(4)
+    volatile CacheContinuousQueryDeployableObject rmtFilterFactoryDep;
 
-    /** Deployable object for filter factory. */
-    private CacheContinuousQueryDeployableObject rmtFilterFactoryDep;
+    /** Marshalled {@link #rmtFilterFactory} if {@link #rmtFilterFactoryDep} 
is {@code null}. */
+    @Order(5)
+    @Nullable volatile byte[] rmtFilterFactoryBytes;
 
     /** Remote filter created by {@link #rmtFilterFactory}. */
-    private transient CacheEntryEventFilter rmtFilterFromFactory;
+    private CacheEntryEventFilter rmtFilterFromFactory;
 
     /** Event types for JCache API. */
-    private byte types;
+    @Order(6)
+    byte types;
 
     /** Remote transformer factory. */
-    private Factory<? extends IgniteClosure<CacheEntryEvent<? extends K, ? 
extends V>, ?>> rmtTransFactory;
+    volatile Factory<? extends IgniteClosure<CacheEntryEvent<? extends K, ? 
extends V>, ?>> rmtTransFactory;
+
+    /** Deployable object for {@link #rmtTransFactory}. Is {@code null} if no 
external marsshalling used. */
+    @Order(7)
+    volatile CacheContinuousQueryDeployableObject rmtTransFactoryDep;
 
-    /** Deployable object for transformer factory. */
-    private CacheContinuousQueryDeployableObject rmtTransFactoryDep;
+    /** Marshalled {@link #rmtTransFactory} if {@link #rmtTransFactoryDep} is 
{@code null}. */
+    @Order(8)
+    @Nullable volatile byte[] rmtTransFactoryBytes;
 
     /** Remote transformer created by {@link #rmtTransFactory}. */
-    private transient IgniteClosure<CacheEntryEvent<? extends K, ? extends V>, 
?> rmtTrans;
+    private IgniteClosure<CacheEntryEvent<? extends K, ? extends V>, ?> 
rmtTrans;
 
     /** Local listener for transformed events. */
-    private transient EventListener<?> locTransLsnr;
+    private EventListener<?> locTransLsnr;
 
     /** Internal flag. */
-    private boolean internal;
+    @Order(9)
+    boolean internal;
 
     /** Notify existing flag. */
-    private boolean notifyExisting;
+    @Order(10)
+    boolean notifyExisting;
 
     /** Old value required flag. */
-    private boolean oldValRequired;
+    @Order(11)
+    boolean oldValRequired;
 
     /** Synchronous flag. */
-    private boolean sync;
+    @Order(12)
+    boolean sync;
 
     /** Ignore expired events flag. */
-    private boolean ignoreExpired;
+    @Order(13)
+    boolean ignoreExpired;
 
     /** Task name hash code. */
-    private int taskHash;
+    @Order(14)
+    int taskHash;
 
     /** Whether to skip primary check for REPLICATED cache. */
-    private transient boolean skipPrimaryCheck;
+    boolean skipPrimaryCheck;
 
     /** */
-    private transient boolean locOnly;
+    private boolean locOnly;
 
     /** */
-    private boolean keepBinary;
+    @Order(15)
+    boolean keepBinary;
 
     /** */
-    private transient ConcurrentMap<Integer, 
CacheContinuousQueryPartitionRecovery> rcvs;
+    private ConcurrentMap<Integer, CacheContinuousQueryPartitionRecovery> rcvs;
 
     /** */
-    private transient ConcurrentMap<Integer, CacheContinuousQueryEventBuffer> 
entryBufs;
+    private ConcurrentMap<Integer, CacheContinuousQueryEventBuffer> entryBufs;
 
     /** */
-    private transient CacheContinuousQueryAcknowledgeBuffer ackBuf;
+    private CacheContinuousQueryAcknowledgeBuffer ackBuf;
 
     /** */
-    private transient int cacheId;
+    private int cacheId;
 
     /** */
-    private transient volatile Map<Integer, Long> initUpdCntrs;
+    private volatile Map<Integer, Long> initUpdCntrs;
 
     /** */
-    private transient volatile Map<UUID, Map<Integer, Long>> 
initUpdCntrsPerNode;
+    private volatile Map<UUID, Map<Integer, Long>> initUpdCntrsPerNode;
 
     /** */
-    private transient volatile AffinityTopologyVersion initTopVer;
+    private volatile AffinityTopologyVersion initTopVer;
 
     /** */
-    private transient volatile boolean nodeLeft;
+    private volatile boolean nodeLeft;
 
     /** */
-    private transient boolean ignoreClsNotFound;
+    private boolean ignoreClsNotFound;
 
     /** */
-    transient boolean asyncCb;
+    boolean asyncCb;
 
     /** */
-    private transient UUID nodeId;
+    private UUID nodeId;
 
     /** */
-    private transient UUID routineId;
+    private UUID routineId;
 
     /** Local update counters values on listener start. Used for skipping 
events fired before the listener start. */
-    private transient volatile Map<Integer, Long> locInitUpdCntrs;
+    private volatile Map<Integer, Long> locInitUpdCntrs;
 
     /** */
-    private transient GridKernalContext ctx;
+    private GridKernalContext ctx;
 
     /** */
-    private transient IgniteLogger log;
+    private IgniteLogger log;
 
     /**
-     * Required by {@link Externalizable}.
+     * Empty constructor for serialization purposes.
      */
     public CacheContinuousQueryHandler() {
         // No-op.
@@ -411,11 +437,6 @@ public class CacheContinuousQueryHandler<K, V> implements 
GridContinuousHandler
         this.locOnly = locOnly;
     }
 
-    /** @return {@code True} if handler are local only, {@code false} 
otherwise. */
-    public boolean localOnly() {
-        return locOnly;
-    }
-
     /**
      * @param taskHash Task hash.
      */
@@ -942,8 +963,6 @@ public class CacheContinuousQueryHandler<K, V> implements 
GridContinuousHandler
      * @throws IgniteCheckedException In case of error.
      */
     void waitTopologyFuture(GridKernalContext ctx) throws 
IgniteCheckedException {
-        GridCacheContext<K, V> cctx = cacheContext(ctx);
-
         AffinityTopologyVersion topVer = initTopVer;
 
         
cacheContext(ctx).shared().exchange().affinityReadyFuture(topVer).get();
@@ -1392,73 +1411,89 @@ public class CacheContinuousQueryHandler<K, V> 
implements GridContinuousHandler
         assert ctx != null;
         assert ctx.config().isPeerClassLoadingEnabled();
 
-        if (requiresDeployment(rmtFilter))
+        // TODO : Remove this check after 
https://issues.apache.org/jira/browse/IGNITE-28945
+        if (rmtFilterDep != null || rmtFilterFactoryDep != null || 
rmtTransFactoryDep != null)
+            return;
+
+        /**
+         * Some filters, factories might be an Ignite-internals and do not 
require external marshaling. But there is no
+         * quarantine that a user-defuned class is not included in a wrap like 
{@link SecurityAwareFilter}. Hence, we always
+         * externally-marshall here.
+         */
+        if (rmtFilter != null)
             rmtFilterDep = new CacheContinuousQueryDeployableObject(rmtFilter, 
ctx);
 
-        if (requiresDeployment(rmtFilterFactory))
+        if (rmtFilterFactory != null)
             rmtFilterFactoryDep = new 
CacheContinuousQueryDeployableObject(rmtFilterFactory, ctx);
 
-        if (requiresDeployment(rmtTransFactory))
+        if (rmtTransFactory != null)
             rmtTransFactoryDep = new 
CacheContinuousQueryDeployableObject(rmtTransFactory, ctx);
     }
 
+    /** {@inheritDoc} */
+    @Override public void marshal(Marshaller marsh) throws 
IgniteCheckedException {
+        if (rmtFilter != null && rmtFilterDep == null)
+            rmtFilterBytes = marsh.marshal(rmtFilter);
+
+        if (rmtFilterFactory != null && rmtFilterFactoryDep == null)
+            rmtFilterFactoryBytes = marsh.marshal(rmtFilterFactory);
+
+        if (rmtTransFactory != null && rmtTransFactoryDep == null)
+            rmtTransFactoryBytes = marsh.marshal(rmtTransFactory);
+    }
+
+    /** {@inheritDoc} */
+    @Override public void unmarshal(Marshaller marsh, ClassLoader clsLdr) 
throws IgniteCheckedException {
+        if (rmtFilterBytes != null && rmtFilterDep == null)
+            rmtFilter = marsh.unmarshal(rmtFilterBytes, clsLdr);
+
+        if (rmtFilterFactoryBytes != null && rmtFilterFactoryDep == null)
+            rmtFilterFactory = marsh.unmarshal(rmtFilterFactoryBytes, clsLdr);
+
+        if (rmtTransFactoryBytes != null && rmtTransFactoryDep == null)
+            rmtTransFactory = marsh.unmarshal(rmtTransFactoryBytes, clsLdr);
+
+        if (rmtFilterDep != null || rmtFilterFactoryDep != null || 
rmtTransFactoryDep != null)
+            p2pUnmarshalFut = new GridFutureAdapter<>();
+
+        cacheId = CU.cacheId(cacheName);
+    }
+
     /** {@inheritDoc} */
     @Override public void p2pUnmarshal(UUID nodeId, GridKernalContext ctx) 
throws IgniteCheckedException {
         assert nodeId != null;
         assert ctx != null;
         assert ctx.config().isPeerClassLoadingEnabled();
 
-        if (rmtFilterDep != null)
-            rmtFilter = p2pUnmarshal(rmtFilterDep, nodeId, ctx);
+        // TODO : Remove this check after 
https://issues.apache.org/jira/browse/IGNITE-28945
+        if (rmtFilter != null || rmtFilterFactory != null || rmtTransFactory 
!= null)
+            return;
 
-        if (rmtFilterFactoryDep != null)
-            rmtFilterFactory = p2pUnmarshal(rmtFilterFactoryDep, nodeId, ctx);
+        try {
+            if (rmtFilterDep != null)
+                rmtFilter = rmtFilterDep.unmarshal(nodeId, ctx);
 
-        if (rmtTransFactoryDep != null)
-            rmtTransFactory = p2pUnmarshal(rmtTransFactoryDep, nodeId, ctx);
+            if (rmtFilterFactoryDep != null)
+                rmtFilterFactory = rmtFilterFactoryDep.unmarshal(nodeId, ctx);
 
-        if (!p2pUnmarshalFut.isDone())
-            ((GridFutureAdapter)p2pUnmarshalFut).onDone();
-    }
+            if (rmtTransFactoryDep != null)
+                rmtTransFactory = rmtTransFactoryDep.unmarshal(nodeId, ctx);
 
-    /**
-     * @return Whether the handler is marshalled for peer class loading.
-     */
-    public boolean isMarshalled() {
-        return (!requiresDeployment(rmtFilter) || rmtFilterDep != null)
-            && (!requiresDeployment(rmtFilterFactory) || rmtFilterFactoryDep 
!= null)
-            && (!requiresDeployment(rmtTransFactory) || rmtTransFactoryDep != 
null);
-    }
-
-    /**
-     * @param depObj Deployable object to unmarshal.
-     * @param nodeId Sender node Id.
-     * @param ctx Kernal context.
-     * @param <T> Result type.
-     * @return Unmarshalled object.
-     * @throws IgniteCheckedException In case of unmarshalling failures.
-     */
-    protected <T> T p2pUnmarshal(CacheContinuousQueryDeployableObject depObj,
-        UUID nodeId, GridKernalContext ctx) throws IgniteCheckedException {
-        if (depObj != null) {
-            try {
-                return depObj.unmarshal(nodeId, ctx);
-            }
-            catch (IgniteCheckedException e) {
-                ((GridFutureAdapter)p2pUnmarshalFut).onDone(e);
+            if (!p2pUnmarshalFut.isDone())
+                ((GridFutureAdapter)p2pUnmarshalFut).onDone();
+        }
+        catch (IgniteCheckedException e) {
+            ((GridFutureAdapter<?>)p2pUnmarshalFut).onDone(e);
 
-                throw e;
-            }
-            catch (ExceptionInInitializerError e) {
-                IgniteCheckedException err = new 
IgniteCheckedException("Failed to unmarshal deployable object.", e);
+            throw e;
+        }
+        catch (ExceptionInInitializerError e) {
+            IgniteCheckedException err = new IgniteCheckedException("Failed to 
unmarshal deployable object.", e);
 
-                ((GridFutureAdapter)p2pUnmarshalFut).onDone(err);
+            ((GridFutureAdapter<?>)p2pUnmarshalFut).onDone(err);
 
-                throw err;
-            }
+            throw err;
         }
-        else
-            return null;
     }
 
     /** {@inheritDoc} */
@@ -1554,87 +1589,6 @@ public class CacheContinuousQueryHandler<K, V> 
implements GridContinuousHandler
         return S.toString(CacheContinuousQueryHandler.class, this);
     }
 
-    /** {@inheritDoc} */
-    @Override public void writeExternal(ObjectOutput out) throws IOException {
-        U.writeString(out, cacheName);
-        out.writeObject(topic);
-
-        writeDeployable(out, rmtFilter, rmtFilterDep);
-
-        out.writeBoolean(internal);
-        out.writeBoolean(notifyExisting);
-        out.writeBoolean(oldValRequired);
-        out.writeBoolean(sync);
-        out.writeBoolean(ignoreExpired);
-        out.writeInt(taskHash);
-        out.writeBoolean(keepBinary);
-
-        writeDeployable(out, rmtFilterFactory, rmtFilterFactoryDep);
-
-        out.writeByte(types);
-
-        writeDeployable(out, rmtTransFactory, rmtTransFactoryDep);
-    }
-
-    /** {@inheritDoc} */
-    @Override public void readExternal(ObjectInput in) throws IOException, 
ClassNotFoundException {
-        cacheName = U.readString(in);
-        topic = in.readObject();
-
-        boolean b = in.readBoolean();
-
-        if (b) {
-            rmtFilterDep = 
(CacheContinuousQueryDeployableObject)in.readObject();
-
-            p2pUnmarshalFut = new GridFutureAdapter<>();
-        }
-        else
-            rmtFilter = (CacheEntryEventSerializableFilter<K, 
V>)in.readObject();
-
-        internal = in.readBoolean();
-        notifyExisting = in.readBoolean();
-        oldValRequired = in.readBoolean();
-        sync = in.readBoolean();
-        ignoreExpired = in.readBoolean();
-        taskHash = in.readInt();
-        keepBinary = in.readBoolean();
-
-        b = in.readBoolean();
-
-        if (b) {
-            rmtFilterFactoryDep = 
(CacheContinuousQueryDeployableObject)in.readObject();
-
-            if (p2pUnmarshalFut.isDone())
-                p2pUnmarshalFut = new GridFutureAdapter<>();
-        }
-        else
-            rmtFilterFactory = (Factory)in.readObject();
-
-        types = in.readByte();
-
-        b = in.readBoolean();
-
-        if (b) {
-            rmtTransFactoryDep = 
(CacheContinuousQueryDeployableObject)in.readObject();
-
-            if (p2pUnmarshalFut.isDone())
-                p2pUnmarshalFut = new GridFutureAdapter<>();
-        }
-        else
-            rmtTransFactory = (Factory<? extends 
IgniteClosure<CacheEntryEvent<? extends K, ? extends V>, ?>>)in.readObject();
-
-        cacheId = CU.cacheId(cacheName);
-    }
-
-    /** */
-    private static void writeDeployable(ObjectOutput out, Object obj, 
CacheContinuousQueryDeployableObject dep) throws IOException {
-        boolean b = dep != null;
-
-        out.writeBoolean(b);
-
-        out.writeObject(b ? dep : obj);
-    }
-
     /**
      * @param ctx Kernal context.
      * @return Cache context.
@@ -1812,9 +1766,4 @@ public class CacheContinuousQueryHandler<K, V> implements 
GridContinuousHandler
     Map<Integer, CacheContinuousQueryEventBuffer> 
partitionContinuesQueryEntryBuffers() {
         return Collections.unmodifiableMap(entryBufs);
     }
-
-    /** */
-    private static boolean requiresDeployment(@Nullable Object obj) {
-        return obj != null && !U.isGrid(obj.getClass());
-    }
 }
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryData.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryData.java
new file mode 100644
index 00000000000..44f0a2f231a
--- /dev/null
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryData.java
@@ -0,0 +1,72 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.ignite.internal.processors.continuous;
+
+import java.util.ArrayList;
+import java.util.Collection;
+import java.util.Map;
+import java.util.UUID;
+import org.apache.ignite.internal.Order;
+import org.apache.ignite.internal.util.tostring.GridToStringInclude;
+import org.apache.ignite.internal.util.typedef.internal.S;
+import org.apache.ignite.plugin.extensions.communication.Message;
+
+/** Continous routine Discovery data. */
+public final class ContinousRoutineDiscoveryData implements Message {
+    /** Node ID.  */
+    @Order(0)
+    UUID nodeId;
+
+    /** Items. */
+    @GridToStringInclude
+    @Order(1)
+    Collection<ContinousRoutineDiscoveryDataItem> items;
+
+    /** */
+    @Order(2)
+    Map<UUID, Map<UUID, ContinousRoutineLocalInfo>> clientInfos;
+
+    /** Empty constructor for serialization purposes. */
+    public ContinousRoutineDiscoveryData() {
+        // No-op.
+    }
+
+    /**
+     * @param nodeId Node ID.
+     * @param clientInfos Client information.
+     */
+    ContinousRoutineDiscoveryData(UUID nodeId, Map<UUID, Map<UUID, 
ContinousRoutineLocalInfo>> clientInfos) {
+        assert nodeId != null;
+
+        this.nodeId = nodeId;
+
+        this.clientInfos = clientInfos;
+
+        items = new ArrayList<>();
+    }
+
+    /** @param item Item. */
+    public void addItem(ContinousRoutineDiscoveryDataItem item) {
+        items.add(item);
+    }
+
+    /** {@inheritDoc} */
+    @Override public String toString() {
+        return S.toString(ContinousRoutineDiscoveryData.class, this);
+    }
+}
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java
new file mode 100644
index 00000000000..650229ae0ff
--- /dev/null
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineDiscoveryDataItem.java
@@ -0,0 +1,96 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.ignite.internal.processors.continuous;
+
+import java.util.UUID;
+import org.apache.ignite.cluster.ClusterNode;
+import org.apache.ignite.internal.Marshalled;
+import org.apache.ignite.internal.Order;
+import org.apache.ignite.internal.util.typedef.internal.S;
+import org.apache.ignite.lang.IgnitePredicate;
+import org.apache.ignite.plugin.extensions.communication.Message;
+import org.jetbrains.annotations.Nullable;
+
+/** Continous routine Discovery data item. */
+public final class ContinousRoutineDiscoveryDataItem implements Message {
+    /** */
+    @Order(0)
+    UUID routineId;
+
+    /** */
+    @Marshalled("prjPredBytes")
+    IgnitePredicate<ClusterNode> prjPred;
+
+    /** Marshalled {@link #prjPred}. */
+    @Order(1)
+    byte[] prjPredBytes;
+
+    /** Handler. */
+    @Order(2)
+    GridContinuousHandler hnd;
+
+    /** Buffer size. */
+    @Order(3)
+    int bufSize;
+
+    /** Time interval. */
+    @Order(4)
+    long interval;
+
+    /** Automatic unsubscribe flag. */
+    @Order(5)
+    boolean autoUnsubscribe;
+
+    /** Empty constructor for serialization porposes. */
+    public ContinousRoutineDiscoveryDataItem() {
+        // No-op.
+    }
+
+    /**
+     * @param routineId Consume ID.
+     * @param prjPred Projection predicate.
+     * @param hnd Handler.
+     * @param bufSize Buffer size.
+     * @param interval Time interval.
+     * @param autoUnsubscribe Automatic unsubscribe flag.
+     */
+    ContinousRoutineDiscoveryDataItem(UUID routineId,
+        @Nullable IgnitePredicate<ClusterNode> prjPred,
+        GridContinuousHandler hnd,
+        int bufSize,
+        long interval,
+        boolean autoUnsubscribe
+    ) {
+        assert routineId != null;
+        assert hnd != null;
+        assert bufSize > 0;
+        assert interval >= 0;
+
+        this.routineId = routineId;
+        this.prjPred = prjPred;
+        this.hnd = hnd;
+        this.bufSize = bufSize;
+        this.interval = interval;
+        this.autoUnsubscribe = autoUnsubscribe;
+    }
+
+    /** {@inheritDoc} */
+    @Override public String toString() {
+        return S.toString(ContinousRoutineDiscoveryDataItem.class, this);
+    }
+}
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java
new file mode 100644
index 00000000000..8bcbc4f8bd9
--- /dev/null
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinousRoutineLocalInfo.java
@@ -0,0 +1,133 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.ignite.internal.processors.continuous;
+
+import java.util.UUID;
+import org.apache.ignite.cluster.ClusterNode;
+import org.apache.ignite.internal.Marshalled;
+import org.apache.ignite.internal.Order;
+import org.apache.ignite.internal.UseBinaryMarshaller;
+import org.apache.ignite.internal.util.typedef.internal.S;
+import org.apache.ignite.lang.IgnitePredicate;
+import org.apache.ignite.plugin.extensions.communication.Message;
+import org.jetbrains.annotations.Nullable;
+
+/** Continous routine local info Discovery data.  */
+@UseBinaryMarshaller
+public final class ContinousRoutineLocalInfo implements Message, 
GridContinuousProcessor.RoutineInfo {
+    /** Source node id. */
+    @Order(0)
+    UUID nodeId;
+
+    /** Projection predicate. */
+    @Marshalled("prjPredBytes")
+    IgnitePredicate<ClusterNode> prjPred;
+
+    /** Marshalled {@link #prjPred}. */
+    @Order(1)
+    byte[] prjPredBytes;
+
+    /** Continuous routine handler. */
+    @Order(2)
+    GridContinuousHandler hnd;
+
+    /** Buffer size. */
+    @Order(3)
+    int bufSize;
+
+    /** Time interval. */
+    @Order(4)
+    long interval;
+
+    /** Automatic unsubscribe flag. */
+    @Order(5)
+    boolean autoUnsubscribe;
+
+    /** Empty constructor for serialization purposes. */
+    public ContinousRoutineLocalInfo() {
+        // No-op.
+    }
+
+    /**
+     * @param nodeId          Node id.
+     * @param prjPred         Projection predicate.
+     * @param hnd             Continuous routine handler.
+     * @param bufSize         Buffer size.
+     * @param interval        Interval.
+     * @param autoUnsubscribe Automatic unsubscribe flag.
+     */
+    ContinousRoutineLocalInfo(
+        UUID nodeId,
+        @Nullable IgnitePredicate<ClusterNode> prjPred,
+        GridContinuousHandler hnd,
+        int bufSize,
+        long interval,
+        boolean autoUnsubscribe
+    ) {
+        assert hnd != null;
+        assert bufSize > 0;
+        assert interval >= 0;
+
+        this.nodeId = nodeId;
+        this.prjPred = prjPred;
+        this.hnd = hnd;
+        this.bufSize = bufSize;
+        this.interval = interval;
+        this.autoUnsubscribe = autoUnsubscribe;
+    }
+
+    /** {@inheritDoc} */
+    @Override public GridContinuousHandler handler() {
+        return hnd;
+    }
+
+    /** {@inheritDoc} */
+    @Override public int bufferSize() {
+        return bufSize;
+    }
+
+    /** {@inheritDoc} */
+    @Override public long interval() {
+        return interval;
+    }
+
+    /** {@inheritDoc} */
+    @Override public boolean autoUnsubscribe() {
+        return autoUnsubscribe;
+    }
+
+    /** {@inheritDoc} */
+    @Override public long lastSendTime() {
+        return -1;
+    }
+
+    /** {@inheritDoc} */
+    @Override public boolean delayedRegister() {
+        return false;
+    }
+
+    /** {@inheritDoc} */
+    @Override public UUID nodeId() {
+        return nodeId;
+    }
+
+    /** {@inheritDoc} */
+    @Override public String toString() {
+        return S.toString(ContinousRoutineLocalInfo.class, this);
+    }
+}
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java
index 207a8f4fda9..8c5a920ffef 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutineInfo.java
@@ -17,45 +17,53 @@
 
 package org.apache.ignite.internal.processors.continuous;
 
-import java.io.Serializable;
 import java.util.UUID;
+import org.apache.ignite.internal.Order;
 import org.apache.ignite.internal.util.typedef.internal.S;
+import org.apache.ignite.plugin.extensions.communication.Message;
 
-/**
- *
- */
-class ContinuousRoutineInfo implements Serializable {
-    /** */
-    private static final long serialVersionUID = 0L;
-
+/** */
+public final class ContinuousRoutineInfo implements Message {
     /** */
+    @Order(0)
     UUID srcNodeId;
 
     /** */
-    final UUID routineId;
+    @Order(1)
+    UUID routineId;
 
     /** */
-    final byte[] hnd;
+    @Order(2)
+    GridContinuousHandler hnd;
 
     /** */
-    final byte[] nodeFilter;
+    @Order(3)
+    byte[] nodeFilter;
 
     /** */
-    final int bufSize;
+    @Order(4)
+    int bufSize;
 
     /** */
-    final long interval;
+    @Order(5)
+    long interval;
 
     /** */
-    final boolean autoUnsubscribe;
+    @Order(6)
+    boolean autoUnsubscribe;
 
-    /** */
-    transient boolean disconnected;
+    /** Transient. */
+    boolean disconnected;
+
+    /** Empty constructor for serialization purposes. */
+    public ContinuousRoutineInfo() {
+        // No-op.
+    }
 
     /**
      * @param srcNodeId Source node ID.
      * @param routineId Routine ID.
-     * @param hnd Marshalled handler.
+     * @param hnd Handler.
      * @param nodeFilter Marshalled node filter.
      * @param bufSize Handler buffer size.
      * @param interval Time interval.
@@ -64,7 +72,7 @@ class ContinuousRoutineInfo implements Serializable {
     ContinuousRoutineInfo(
         UUID srcNodeId,
         UUID routineId,
-        byte[] hnd,
+        GridContinuousHandler hnd,
         byte[] nodeFilter,
         int bufSize,
         long interval,
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesCommonDiscoveryData.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesCommonDiscoveryData.java
index d29de89b9b1..7b6ad6a3396 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesCommonDiscoveryData.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesCommonDiscoveryData.java
@@ -17,19 +17,23 @@
 
 package org.apache.ignite.internal.processors.continuous;
 
-import java.io.Serializable;
 import java.util.List;
+import org.apache.ignite.internal.Order;
 import org.apache.ignite.internal.util.typedef.internal.S;
+import org.apache.ignite.plugin.extensions.communication.Message;
 
 /**
  *
  */
-public class ContinuousRoutinesCommonDiscoveryData implements Serializable {
+public final class ContinuousRoutinesCommonDiscoveryData implements Message {
     /** */
-    private static final long serialVersionUID = 0L;
+    @Order(0)
+    List<ContinuousRoutineInfo> startedRoutines;
 
-    /** */
-    final List<ContinuousRoutineInfo> startedRoutines;
+    /** Empty constructor for serialization purposes. */
+    public ContinuousRoutinesCommonDiscoveryData() {
+        // No-op.
+    }
 
     /**
      * @param startedRoutines Routines started in cluster.
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java
index 9be6ef8e07e..ec40f56ca1f 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/ContinuousRoutinesJoiningNodeDiscoveryData.java
@@ -17,19 +17,23 @@
 
 package org.apache.ignite.internal.processors.continuous;
 
-import java.io.Serializable;
 import java.util.List;
+import org.apache.ignite.internal.Order;
 import org.apache.ignite.internal.util.typedef.internal.S;
+import org.apache.ignite.plugin.extensions.communication.Message;
 
 /**
  *
  */
-public class ContinuousRoutinesJoiningNodeDiscoveryData implements 
Serializable {
+public final class ContinuousRoutinesJoiningNodeDiscoveryData implements 
Message {
     /** */
-    private static final long serialVersionUID = 0L;
+    @Order(0)
+    List<ContinuousRoutineInfo> startedRoutines;
 
-    /** */
-    final List<ContinuousRoutineInfo> startedRoutines;
+    /** Empty constructor for serialization purposes. */
+    public ContinuousRoutinesJoiningNodeDiscoveryData() {
+        // No-op.
+    }
 
     /**
      * @param startedRoutines Routines registered on nodes, to be started in 
cluster.
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousHandler.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousHandler.java
index b1a3812f612..03ca0defa6d 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousHandler.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousHandler.java
@@ -17,20 +17,20 @@
 
 package org.apache.ignite.internal.processors.continuous;
 
-import java.io.Externalizable;
 import java.util.Collection;
 import java.util.Map;
 import java.util.UUID;
 import org.apache.ignite.IgniteCheckedException;
 import org.apache.ignite.internal.GridKernalContext;
 import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
+import org.apache.ignite.plugin.extensions.communication.Message;
 import org.jetbrains.annotations.Nullable;
 
 /**
  * Continuous routine handler.
  */
 @SuppressWarnings("PublicInnerClass")
-public interface GridContinuousHandler extends Externalizable, Cloneable {
+public interface GridContinuousHandler extends Cloneable, Message {
     /**
      * Listener registration status.
      */
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java
index 111d3b74049..7570b94f37e 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/GridContinuousProcessor.java
@@ -17,12 +17,6 @@
 
 package org.apache.ignite.internal.processors.continuous;
 
-import java.io.Externalizable;
-import java.io.IOException;
-import java.io.ObjectInput;
-import java.io.ObjectOutput;
-import java.io.Serializable;
-import java.util.ArrayList;
 import java.util.Collection;
 import java.util.Collections;
 import java.util.HashMap;
@@ -80,7 +74,6 @@ import 
org.apache.ignite.internal.util.future.GridFinishedFuture;
 import org.apache.ignite.internal.util.future.GridFutureAdapter;
 import org.apache.ignite.internal.util.lang.GridPlainRunnable;
 import org.apache.ignite.internal.util.lang.gridfunc.ReadOnlyCollectionView2X;
-import org.apache.ignite.internal.util.tostring.GridToStringInclude;
 import org.apache.ignite.internal.util.typedef.CI1;
 import org.apache.ignite.internal.util.typedef.F;
 import org.apache.ignite.internal.util.typedef.X;
@@ -125,10 +118,10 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
     public static final String CQ_SYS_VIEW_DESC = "Continuous queries";
 
     /** Local infos. */
-    private final ConcurrentMap<UUID, LocalRoutineInfo> locInfos = new 
ConcurrentHashMap<>();
+    private final ConcurrentMap<UUID, ContinousRoutineLocalInfo> locInfos = 
new ConcurrentHashMap<>();
 
     /** Local infos. */
-    private final ConcurrentMap<UUID, Map<UUID, LocalRoutineInfo>> clientInfos 
= new ConcurrentHashMap<>();
+    private final ConcurrentMap<UUID, Map<UUID, ContinousRoutineLocalInfo>> 
clientInfos = new ConcurrentHashMap<>();
 
     /** Remote infos. */
     private final ConcurrentMap<UUID, RemoteRoutineInfo> rmtInfos = new 
ConcurrentHashMap<>();
@@ -336,12 +329,12 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
     }
 
     /** */
-    Map<UUID, LocalRoutineInfo> localRoutineInfos() {
+    Map<UUID, ContinousRoutineLocalInfo> localRoutineInfos() {
         return Collections.unmodifiableMap(locInfos);
     }
 
     /** */
-    Map<UUID, Map<UUID, LocalRoutineInfo>> clientRoutineInfos() {
+    Map<UUID, Map<UUID, ContinousRoutineLocalInfo>> clientRoutineInfos() {
         return Collections.unmodifiableMap(clientInfos);
     }
 
@@ -408,7 +401,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
             return;
         }
 
-        Serializable data = getDiscoveryData(dataBag.joiningNodeId());
+        ContinousRoutineDiscoveryData data = 
getDiscoveryData(dataBag.joiningNodeId());
 
         if (data != null)
             dataBag.addJoiningNodeData(CONTINUOUS_PROC.ordinal(), data);
@@ -422,7 +415,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
             return;
         }
 
-        Serializable data = getDiscoveryData(dataBag.joiningNodeId());
+        ContinousRoutineDiscoveryData data = 
getDiscoveryData(dataBag.joiningNodeId());
 
         if (data != null)
             dataBag.addNodeSpecificData(CONTINUOUS_PROC.ordinal(), data);
@@ -431,7 +424,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
     /**
      * @param joiningNodeId Joining node id.
      */
-    private Serializable getDiscoveryData(UUID joiningNodeId) {
+    private @Nullable ContinousRoutineDiscoveryData getDiscoveryData(UUID 
joiningNodeId) {
         if (log.isDebugEnabled()) {
             log.debug("collectDiscoveryData [node=" + joiningNodeId +
                     ", loc=" + ctx.localNodeId() +
@@ -441,26 +434,22 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
         }
 
         if (!joiningNodeId.equals(ctx.localNodeId()) || !locInfos.isEmpty()) {
-            Map<UUID, Map<UUID, LocalRoutineInfo>> clientInfos0 = 
copyClientInfos(clientInfos);
+            Map<UUID, Map<UUID, ContinousRoutineLocalInfo>> clientInfos0 = 
copyClientInfos(clientInfos);
 
             if (joiningNodeId.equals(ctx.localNodeId()) && 
ctx.discovery().localNode().isClient()) {
-                Map<UUID, LocalRoutineInfo> infos = copyLocalInfos(locInfos);
+                Map<UUID, ContinousRoutineLocalInfo> infos = 
copyLocalInfos(locInfos);
 
                 clientInfos0.put(ctx.localNodeId(), infos);
             }
 
-            DiscoveryData data = new DiscoveryData(ctx.localNodeId(), 
clientInfos0);
+            ContinousRoutineDiscoveryData data = new 
ContinousRoutineDiscoveryData(ctx.localNodeId(), clientInfos0);
 
             // Collect listeners information (will be sent to joining node 
during discovery process).
-            for (Map.Entry<UUID, LocalRoutineInfo> e : locInfos.entrySet()) {
+            for (Map.Entry<UUID, ContinousRoutineLocalInfo> e : 
locInfos.entrySet()) {
                 UUID routineId = e.getKey();
-                LocalRoutineInfo info = e.getValue();
+                ContinousRoutineLocalInfo info = e.getValue();
 
-                assert !ctx.config().isPeerClassLoadingEnabled() ||
-                    !(info.hnd instanceof CacheContinuousQueryHandler) ||
-                    ((CacheContinuousQueryHandler)info.hnd).isMarshalled();
-
-                data.addItem(new DiscoveryDataItem(routineId,
+                data.addItem(new ContinousRoutineDiscoveryDataItem(routineId,
                     info.prjPred,
                     info.hnd,
                     info.bufSize,
@@ -477,13 +466,13 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
     /**
      * @param clientInfos Client infos.
      */
-    private Map<UUID, Map<UUID, LocalRoutineInfo>> copyClientInfos(Map<UUID, 
Map<UUID, LocalRoutineInfo>> clientInfos) {
-        Map<UUID, Map<UUID, LocalRoutineInfo>> res = 
U.newHashMap(clientInfos.size());
+    private Map<UUID, Map<UUID, ContinousRoutineLocalInfo>> 
copyClientInfos(Map<UUID, Map<UUID, ContinousRoutineLocalInfo>> clientInfos) {
+        Map<UUID, Map<UUID, ContinousRoutineLocalInfo>> res = 
U.newHashMap(clientInfos.size());
 
-        for (Map.Entry<UUID, Map<UUID, LocalRoutineInfo>> e : 
clientInfos.entrySet()) {
-            Map<UUID, LocalRoutineInfo> cp = U.newHashMap(e.getValue().size());
+        for (Map.Entry<UUID, Map<UUID, ContinousRoutineLocalInfo>> e : 
clientInfos.entrySet()) {
+            Map<UUID, ContinousRoutineLocalInfo> cp = 
U.newHashMap(e.getValue().size());
 
-            for (Map.Entry<UUID, LocalRoutineInfo> e0 : 
e.getValue().entrySet())
+            for (Map.Entry<UUID, ContinousRoutineLocalInfo> e0 : 
e.getValue().entrySet())
                 cp.put(e0.getKey(), e0.getValue());
 
             res.put(e.getKey(), cp);
@@ -495,10 +484,10 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
     /**
      * @param locInfos Locale infos.
      */
-    private Map<UUID, LocalRoutineInfo> copyLocalInfos(Map<UUID, 
LocalRoutineInfo> locInfos) {
-        Map<UUID, LocalRoutineInfo> res = U.newHashMap(locInfos.size());
+    private Map<UUID, ContinousRoutineLocalInfo> copyLocalInfos(Map<UUID, 
ContinousRoutineLocalInfo> locInfos) {
+        Map<UUID, ContinousRoutineLocalInfo> res = 
U.newHashMap(locInfos.size());
 
-        for (Map.Entry<UUID, LocalRoutineInfo> e : locInfos.entrySet())
+        for (Map.Entry<UUID, ContinousRoutineLocalInfo> e : 
locInfos.entrySet())
             res.put(e.getKey(), e.getValue());
 
         return res;
@@ -547,10 +536,10 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
             }
         }
         else {
-            Map<UUID, DiscoveryData> nodeSpecData = data.nodeSpecificData();
+            Map<UUID, ContinousRoutineDiscoveryData> nodeSpecData = 
data.nodeSpecificData();
 
             if (nodeSpecData != null) {
-                for (DiscoveryData val : nodeSpecData.values())
+                for (ContinousRoutineDiscoveryData val : nodeSpecData.values())
                     onDiscoveryDataReceivedMutable(val);
             }
         }
@@ -562,37 +551,37 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
      *
      * @param data received discovery data.
      */
-    private void onDiscoveryDataReceivedMutable(DiscoveryData data) {
+    private void onDiscoveryDataReceivedMutable(ContinousRoutineDiscoveryData 
data) {
         if (data != null) {
-            for (DiscoveryDataItem item : data.items) {
+            for (ContinousRoutineDiscoveryDataItem item : data.items) {
                 if (!locInfos.containsKey(item.routineId)) {
                     registerHandlerOnJoin(data.nodeId, item.routineId, 
item.prjPred,
                         item.hnd, item.bufSize, item.interval, 
item.autoUnsubscribe);
                 }
 
                 if (!item.autoUnsubscribe) {
-                    locInfos.putIfAbsent(item.routineId, new 
LocalRoutineInfo(data.nodeId,
+                    locInfos.putIfAbsent(item.routineId, new 
ContinousRoutineLocalInfo(data.nodeId,
                         item.prjPred, item.hnd, item.bufSize, item.interval, 
item.autoUnsubscribe));
                 }
             }
 
             // Process CQs started on clients.
-            for (Map.Entry<UUID, Map<UUID, LocalRoutineInfo>> entry : 
data.clientInfos.entrySet()) {
+            for (Map.Entry<UUID, Map<UUID, ContinousRoutineLocalInfo>> entry : 
data.clientInfos.entrySet()) {
                 UUID clientNodeId = entry.getKey();
 
                 if (!ctx.localNodeId().equals(clientNodeId)) {
-                    Map<UUID, LocalRoutineInfo> clientRoutineMap = 
entry.getValue();
+                    Map<UUID, ContinousRoutineLocalInfo> clientRoutineMap = 
entry.getValue();
 
-                    for (Map.Entry<UUID, LocalRoutineInfo> e : 
clientRoutineMap.entrySet()) {
+                    for (Map.Entry<UUID, ContinousRoutineLocalInfo> e : 
clientRoutineMap.entrySet()) {
                         UUID routineId = e.getKey();
-                        LocalRoutineInfo info = e.getValue();
+                        ContinousRoutineLocalInfo info = e.getValue();
 
                         registerHandlerOnJoin(clientNodeId, routineId, 
info.prjPred,
                             info.hnd, info.bufSize, info.interval, 
info.autoUnsubscribe);
                     }
                 }
 
-                Map<UUID, LocalRoutineInfo> map =
+                Map<UUID, ContinousRoutineLocalInfo> map =
                     clientInfos.computeIfAbsent(clientNodeId, k -> new 
HashMap<>());
 
                 map.putAll(entry.getValue());
@@ -628,23 +617,8 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
             return;
         }
 
-        GridContinuousHandler hnd;
-
-        try {
-            hnd = U.unmarshal(marsh, routineInfo.hnd, 
U.resolveClassLoader(ctx.config()));
-        }
-        catch (IgniteCheckedException e) {
-            U.error(log, "Failed to unmarshal continuous routine handler [" +
-                "routineId=" + routineInfo.routineId +
-                ", srcNodeId=" + routineInfo.srcNodeId + ']', e);
-
-            ctx.failure().process(new 
FailureContext(FailureType.CRITICAL_ERROR, e));
-
-            return;
-        }
-
         registerHandlerOnJoin(routineInfo.srcNodeId, routineInfo.routineId, 
nodeFilter,
-            hnd, routineInfo.bufSize, routineInfo.interval, 
routineInfo.autoUnsubscribe);
+            routineInfo.hnd, routineInfo.bufSize, routineInfo.interval, 
routineInfo.autoUnsubscribe);
     }
 
     /**
@@ -782,7 +756,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
 
         final UUID routineId = UUID.randomUUID();
 
-        LocalRoutineInfo routineInfo = new LocalRoutineInfo(ctx.localNodeId(), 
prjPred, hnd, 1, 0, true);
+        ContinousRoutineLocalInfo routineInfo = new 
ContinousRoutineLocalInfo(ctx.localNodeId(), prjPred, hnd, 1, 0, true);
 
         if (immutableDiscoCustomMsg) {
             routinesInfo.addRoutineInfo(createRoutineInfo(
@@ -822,14 +796,13 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
         long interval,
         boolean autoUnsubscribe)
         throws IgniteCheckedException {
-        byte[] hndBytes = marsh.marshal(hnd);
 
         byte[] filterBytes = nodeFilter != null ? marsh.marshal(nodeFilter) : 
null;
 
         return new ContinuousRoutineInfo(
             srcNodeId,
             routineId,
-            hndBytes,
+            hnd,
             filterBytes,
             bufSize,
             interval,
@@ -858,15 +831,12 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
         // Generate ID.
         final UUID routineId = UUID.randomUUID();
 
-        if (ctx.config().isPeerClassLoadingEnabled()) {
+        if (ctx.config().isPeerClassLoadingEnabled())
             hnd.p2pMarshal(ctx);
 
-            assert !(hnd instanceof CacheContinuousQueryHandler) || 
((CacheContinuousQueryHandler)hnd).isMarshalled();
-        }
-
         // Register routine locally.
         locInfos.put(routineId,
-            new LocalRoutineInfo(ctx.localNodeId(), prjPred, hnd, bufSize, 
interval, autoUnsubscribe));
+            new ContinousRoutineLocalInfo(ctx.localNodeId(), prjPred, hnd, 
bufSize, interval, autoUnsubscribe));
 
         if (locOnly) {
             try {
@@ -986,8 +956,6 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
             hnd.p2pMarshal(ctx);
         }
 
-        reqData.hndBytes = U.marshal(marsh, hnd);
-
         if (nodeFilter != null)
             reqData.nodeFilterBytes = U.marshal(marsh, nodeFilter);
 
@@ -1066,7 +1034,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
                 boolean stop = false;
 
                 // Unregister routine locally.
-                LocalRoutineInfo routine = locInfos.remove(routineId);
+                ContinousRoutineLocalInfo routine = locInfos.remove(routineId);
 
                 if (routine != null) {
                     stop = true;
@@ -1131,7 +1099,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
                 sendNotification(nodeId, routineId, null, toSnd, orderedTopic, 
true, null);
         }
         else {
-            LocalRoutineInfo locRoutineInfo = locInfos.get(routineId);
+            ContinousRoutineLocalInfo locRoutineInfo = locInfos.get(routineId);
 
             if (locRoutineInfo != null)
                 locRoutineInfo.handler().notifyCallback(nodeId, routineId, 
objs, ctx);
@@ -1256,7 +1224,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
                 unregisterRemote(e.getKey());
         }
 
-        for (LocalRoutineInfo routine : locInfos.values())
+        for (ContinousRoutineLocalInfo routine : locInfos.values())
             routine.hnd.onClientDisconnected();
 
         rmtInfos.clear();
@@ -1323,7 +1291,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
             unregisterRemote(routineId);
         }
 
-        for (Map<UUID, LocalRoutineInfo> clientInfo : clientInfos.values()) {
+        for (Map<UUID, ContinousRoutineLocalInfo> clientInfo : 
clientInfos.values()) {
             if (clientInfo.remove(msg.routineId()) != null)
                 break;
         }
@@ -1359,17 +1327,13 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
         data.nodeFilter = U.unmarshal(marsh, data.nodeFilterBytes,
             ctx.deploy().classLoader(data.depInfo, data.clsName, sndId));
 
-        if (data.hndBytes != null) {
-            data.hnd = U.unmarshal(marsh, data.hndBytes, 
U.resolveClassLoader(ctx.config()));
+        if (ctx.config().isPeerClassLoadingEnabled())
+            data.hnd.p2pUnmarshal(sndId, ctx);
 
-            if (ctx.config().isPeerClassLoadingEnabled())
-                data.hnd.p2pUnmarshal(sndId, ctx);
+        if (data.keepBinary) {
+            assert data.hnd instanceof CacheContinuousQueryHandler : data.hnd;
 
-            if (data.keepBinary) {
-                assert data.hnd instanceof CacheContinuousQueryHandler : 
data.hnd;
-
-                ((CacheContinuousQueryHandler<?, ?>)data.hnd).keepBinary(true);
-            }
+            ((CacheContinuousQueryHandler<?, ?>)data.hnd).keepBinary(true);
         }
     }
 
@@ -1409,17 +1373,17 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
         GridContinuousHandler hnd = data.handler();
 
         if (node.isClient()) {
-            Map<UUID, LocalRoutineInfo> clientRoutineMap = 
clientInfos.get(node.id());
+            Map<UUID, ContinousRoutineLocalInfo> clientRoutineMap = 
clientInfos.get(node.id());
 
             if (clientRoutineMap == null) {
                 clientRoutineMap = new HashMap<>();
 
-                Map<UUID, LocalRoutineInfo> old = clientInfos.put(node.id(), 
clientRoutineMap);
+                Map<UUID, ContinousRoutineLocalInfo> old = 
clientInfos.put(node.id(), clientRoutineMap);
 
                 assert old == null;
             }
 
-            clientRoutineMap.put(routineId, new LocalRoutineInfo(node.id(),
+            clientRoutineMap.put(routineId, new 
ContinousRoutineLocalInfo(node.id(),
                 data.nodeFilter(),
                 hnd,
                 data.bufferSize(),
@@ -1457,7 +1421,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
 
                 if (!data.autoUnsubscribe())
                     // Register routine locally.
-                    locInfos.putIfAbsent(routineId, new LocalRoutineInfo(
+                    locInfos.putIfAbsent(routineId, new 
ContinousRoutineLocalInfo(
                         node.id(), prjPred, hnd, data.bufferSize(), 
data.interval(), data.autoUnsubscribe()));
             }
             catch (IgniteCheckedException e) {
@@ -1503,7 +1467,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
 
         ContinuousRoutineInfo routineInfo = new ContinuousRoutineInfo(snd.id(),
             msg.routineId(),
-            reqData.hndBytes,
+            reqData.hnd,
             reqData.nodeFilterBytes,
             reqData.bufferSize(),
             reqData.interval(),
@@ -1642,7 +1606,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
         UUID routineId = msg.routineId();
 
         try {
-            LocalRoutineInfo routine = locInfos.get(routineId);
+            ContinousRoutineLocalInfo routine = locInfos.get(routineId);
 
             if (routine != null)
                 routine.hnd.notifyCallback(nodeId, routineId, 
(Collection<?>)msg.data(), ctx);
@@ -1804,7 +1768,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
     @SuppressWarnings("TooBroadScope")
     private void unregisterRemote(UUID routineId) {
         RemoteRoutineInfo remote;
-        LocalRoutineInfo loc;
+        ContinousRoutineLocalInfo loc;
 
         stopLock.lock();
 
@@ -2006,100 +1970,6 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
         boolean delayedRegister();
     }
 
-    /**
-     * Local routine info.
-     */
-    public static class LocalRoutineInfo implements Serializable, RoutineInfo {
-        /** */
-        private static final long serialVersionUID = 0L;
-
-        /** Source node id. */
-        private final UUID nodeId;
-
-        /** Projection predicate. */
-        private final IgnitePredicate<ClusterNode> prjPred;
-
-        /** Continuous routine handler. */
-        private final GridContinuousHandler hnd;
-
-        /** Buffer size. */
-        private final int bufSize;
-
-        /** Time interval. */
-        private final long interval;
-
-        /** Automatic unsubscribe flag. */
-        private boolean autoUnsubscribe;
-
-        /**
-         * @param nodeId Node id.
-         * @param prjPred Projection predicate.
-         * @param hnd Continuous routine handler.
-         * @param bufSize Buffer size.
-         * @param interval Interval.
-         * @param autoUnsubscribe Automatic unsubscribe flag.
-         */
-        LocalRoutineInfo(
-            UUID nodeId,
-            @Nullable IgnitePredicate<ClusterNode> prjPred,
-            GridContinuousHandler hnd,
-            int bufSize,
-            long interval,
-            boolean autoUnsubscribe
-        ) {
-            assert hnd != null;
-            assert bufSize > 0;
-            assert interval >= 0;
-
-            this.nodeId = nodeId;
-            this.prjPred = prjPred;
-            this.hnd = hnd;
-            this.bufSize = bufSize;
-            this.interval = interval;
-            this.autoUnsubscribe = autoUnsubscribe;
-        }
-
-        /** {@inheritDoc} */
-        @Override public GridContinuousHandler handler() {
-            return hnd;
-        }
-
-        /** {@inheritDoc} */
-        @Override public int bufferSize() {
-            return bufSize;
-        }
-
-        /** {@inheritDoc} */
-        @Override public long interval() {
-            return interval;
-        }
-
-        /** {@inheritDoc} */
-        @Override public boolean autoUnsubscribe() {
-            return autoUnsubscribe;
-        }
-
-        /** {@inheritDoc} */
-        @Override public long lastSendTime() {
-            return -1;
-        }
-
-        /** {@inheritDoc} */
-        @Override public boolean delayedRegister() {
-            return false;
-        }
-
-        /** {@inheritDoc} */
-        @Override public UUID nodeId() {
-            return nodeId;
-        }
-
-        /** {@inheritDoc} */
-        @Override public String toString() {
-            return S.toString(LocalRoutineInfo.class, this);
-        }
-    }
-
     /**
      * Remote routine info.
      */
@@ -2321,157 +2191,6 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
         }
     }
 
-    /**
-     * Discovery data.
-     */
-    private static class DiscoveryData implements Externalizable {
-        /** */
-        private static final long serialVersionUID = 0L;
-
-        /** Node ID. */
-        private UUID nodeId;
-
-        /** Items. */
-        @GridToStringInclude
-        private Collection<DiscoveryDataItem> items;
-
-        /** */
-        private Map<UUID, Map<UUID, LocalRoutineInfo>> clientInfos;
-
-        /**
-         * Required by {@link Externalizable}.
-         */
-        public DiscoveryData() {
-            // No-op.
-        }
-
-        /**
-         * @param nodeId Node ID.
-         * @param clientInfos Client information.
-         */
-        DiscoveryData(UUID nodeId, Map<UUID, Map<UUID, LocalRoutineInfo>> 
clientInfos) {
-            assert nodeId != null;
-
-            this.nodeId = nodeId;
-
-            this.clientInfos = clientInfos;
-
-            items = new ArrayList<>();
-        }
-
-        /**
-         * @param item Item.
-         */
-        public void addItem(DiscoveryDataItem item) {
-            items.add(item);
-        }
-
-        /** {@inheritDoc} */
-        @Override public void writeExternal(ObjectOutput out) throws 
IOException {
-            U.writeUuid(out, nodeId);
-            U.writeCollection(out, items);
-            U.writeMap(out, clientInfos);
-        }
-
-        /** {@inheritDoc} */
-        @Override public void readExternal(ObjectInput in) throws IOException, 
ClassNotFoundException {
-            nodeId = U.readUuid(in);
-            items = U.readCollection(in);
-            clientInfos = U.readMap(in);
-        }
-
-        /** {@inheritDoc} */
-        @Override public String toString() {
-            return S.toString(DiscoveryData.class, this);
-        }
-    }
-
-    /**
-     * Discovery data item.
-     */
-    private static class DiscoveryDataItem implements Externalizable {
-        /** */
-        private static final long serialVersionUID = 0L;
-
-        /** Consume ID. */
-        private UUID routineId;
-
-        /** Projection predicate. */
-        private IgnitePredicate<ClusterNode> prjPred;
-
-        /** Handler. */
-        private GridContinuousHandler hnd;
-
-        /** Buffer size. */
-        private int bufSize;
-
-        /** Time interval. */
-        private long interval;
-
-        /** Automatic unsubscribe flag. */
-        private boolean autoUnsubscribe;
-
-        /**
-         * Required by {@link Externalizable}.
-         */
-        public DiscoveryDataItem() {
-            // No-op.
-        }
-
-        /**
-         * @param routineId Consume ID.
-         * @param prjPred Projection predicate.
-         * @param hnd Handler.
-         * @param bufSize Buffer size.
-         * @param interval Time interval.
-         * @param autoUnsubscribe Automatic unsubscribe flag.
-         */
-        DiscoveryDataItem(UUID routineId,
-            @Nullable IgnitePredicate<ClusterNode> prjPred,
-            GridContinuousHandler hnd,
-            int bufSize,
-            long interval,
-            boolean autoUnsubscribe
-        ) {
-            assert routineId != null;
-            assert hnd != null;
-            assert bufSize > 0;
-            assert interval >= 0;
-
-            this.routineId = routineId;
-            this.prjPred = prjPred;
-            this.hnd = hnd;
-            this.bufSize = bufSize;
-            this.interval = interval;
-            this.autoUnsubscribe = autoUnsubscribe;
-        }
-
-        /** {@inheritDoc} */
-        @Override public void writeExternal(ObjectOutput out) throws 
IOException {
-            U.writeUuid(out, routineId);
-            out.writeObject(prjPred);
-            out.writeObject(hnd);
-            out.writeInt(bufSize);
-            out.writeLong(interval);
-            out.writeBoolean(autoUnsubscribe);
-        }
-
-        /** {@inheritDoc} */
-        @Override public void readExternal(ObjectInput in) throws IOException, 
ClassNotFoundException {
-            routineId = U.readUuid(in);
-            prjPred = (IgnitePredicate<ClusterNode>)in.readObject();
-            hnd = (GridContinuousHandler)in.readObject();
-            bufSize = in.readInt();
-            interval = in.readLong();
-            autoUnsubscribe = in.readBoolean();
-        }
-
-        /** {@inheritDoc} */
-        @Override public String toString() {
-            return S.toString(DiscoveryDataItem.class, this);
-        }
-    }
-
     /**
      * Future for start routine.
      */
@@ -2556,7 +2275,7 @@ public class GridContinuousProcessor extends 
GridProcessorAdapter {
             Map<Integer, Long> cntrs) {
             try {
                 if (errs == null || errs.isEmpty()) {
-                    LocalRoutineInfo routine = locInfos.get(routineId);
+                    ContinousRoutineLocalInfo routine = 
locInfos.get(routineId);
 
                     // Update partition counters.
                     if (routine != null && routine.handler().isQuery()) {
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/StartRequestData.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/StartRequestData.java
index 595aaa56c58..2d52f2dfe0b 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/StartRequestData.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/continuous/StartRequestData.java
@@ -44,11 +44,8 @@ public class StartRequestData implements Message {
     GridDeploymentInfoMessage depInfo;
 
     /** Handler, restored by the processor reading this request. */
-    GridContinuousHandler hnd;
-
-    /** Serialized handler. */
     @Order(3)
-    byte[] hndBytes;
+    GridContinuousHandler hnd;
 
     /** Buffer size. */
     @Order(4)
diff --git 
a/modules/core/src/main/java/org/apache/ignite/spi/IgniteSpiAdapter.java 
b/modules/core/src/main/java/org/apache/ignite/spi/IgniteSpiAdapter.java
index 850b196a747..79abc47ac7f 100644
--- a/modules/core/src/main/java/org/apache/ignite/spi/IgniteSpiAdapter.java
+++ b/modules/core/src/main/java/org/apache/ignite/spi/IgniteSpiAdapter.java
@@ -81,7 +81,7 @@ public abstract class IgniteSpiAdapter implements IgniteSpi {
     protected IgniteLogger log;
 
     /** Ignite instance. */
-    protected Ignite ignite;
+    protected IgniteEx ignite;
 
     /** Ignite instance name. */
     protected String igniteInstanceName;
@@ -268,7 +268,7 @@ public abstract class IgniteSpiAdapter implements IgniteSpi 
{
      */
     @IgniteInstanceResource
     protected void injectResources(Ignite ignite) {
-        this.ignite = ignite;
+        this.ignite = (IgniteEx)ignite;
 
         if (ignite != null && igniteInstanceName == null)
             igniteInstanceName = ignite.name();
diff --git 
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/DiscoveryDataBag.java
 
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/DiscoveryDataBag.java
index a2f3af751af..df44689dff8 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/DiscoveryDataBag.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/DiscoveryDataBag.java
@@ -280,14 +280,6 @@ public class DiscoveryDataBag {
         commonData.put(cmpId, data);
     }
 
-    /**
-     * @param cmpId Component ID.
-     * @param data Serializable data.
-     */
-    public void addNodeSpecificData(Integer cmpId, Serializable data) {
-        addNodeSpecificData(cmpId, new SerializableDataBagItemWrapper(data));
-    }
-
     /**
      * @param cmpId Component ID.
      * @param data Message data.
diff --git 
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
 
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
index b894f0b8772..e6d27b3bb48 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
@@ -864,7 +864,7 @@ class ClientImpl extends TcpDiscoveryImpl {
     }
 
     /**
-     * Marshalls credentials with discovery SPI marshaller (will replace 
attribute value).
+     * Marshals credentials with discovery SPI marshaller (will replace 
attribute value).
      *
      * @param node Node to marshall credentials for.
      * @throws IgniteSpiException If marshalling failed.
diff --git 
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
 
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
index e84c716c8f2..89cc89e3359 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
@@ -1672,7 +1672,7 @@ class ServerImpl extends TcpDiscoveryImpl {
     }
 
     /**
-     * Marshalls credentials with discovery SPI marshaller (will replace 
attribute value).
+     * Marshals credentials with discovery SPI marshaller (will replace 
attribute value).
      *
      * @param node Node to marshall credentials for.
      * @param cred Credentials for marshall.
@@ -1693,7 +1693,7 @@ class ServerImpl extends TcpDiscoveryImpl {
     }
 
     /**
-     * Unmarshalls credentials with discovery SPI marshaller (will not replace 
attribute value).
+     * Unmarshals credentials with discovery SPI marshaller (will not replace 
attribute value).
      *
      * @param node Node to unmarshall credentials for.
      * @return Security credentials.
diff --git 
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpi.java
 
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpi.java
index 9c52f9504be..b92ca657ad7 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpi.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpi.java
@@ -54,6 +54,7 @@ import org.apache.ignite.cluster.ClusterNode;
 import org.apache.ignite.configuration.AddressResolver;
 import org.apache.ignite.configuration.IgniteConfiguration;
 import org.apache.ignite.failure.FailureContext;
+import org.apache.ignite.failure.FailureType;
 import org.apache.ignite.internal.IgniteEx;
 import org.apache.ignite.internal.IgniteInterruptedCheckedException;
 import 
org.apache.ignite.internal.managers.communication.UnknownMessageException;
@@ -1888,6 +1889,18 @@ public class TcpDiscoverySpi extends IgniteSpiAdapter 
implements IgniteDiscovery
                 if (msg != null && sslMsgPattern.matcher(msg).matches())
                     streamCorruptedCause.initCause(new SSLException("Detected 
SSL alert in StreamCorruptedException"));
             }
+
+            if (X.hasCause(e, ClassNotFoundException.class)) {
+                LT.error(log, e, "Failed to read message due to an unknown 
class to unmarshal received. Unable to " +
+                    "process the Discovery protocol. Stopping the Discovery 
SPI and invoking the failure handler. " +
+                    "RmtAddr=" + sock.getRemoteSocketAddress() + ", rmtPort=" 
+ sock.getPort() + ']');
+
+                ignite.context().failure().process(new 
FailureContext(FailureType.CRITICAL_ERROR, e));
+
+                // Prevents following cycling attempts to reconnect and logs 
flooding.
+                spiStop();
+            }
+
             throw e;
         }
         finally {
diff --git 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryEntriesExpireTest.java
 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryEntriesExpireTest.java
index 736b1742bef..317076bc92a 100644
--- 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryEntriesExpireTest.java
+++ 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/CacheContinuousQueryEntriesExpireTest.java
@@ -28,7 +28,7 @@ import org.apache.ignite.cache.query.ContinuousQuery;
 import org.apache.ignite.configuration.CacheConfiguration;
 import org.apache.ignite.internal.IgniteEx;
 import org.apache.ignite.internal.processors.cache.GridCacheContext;
-import 
org.apache.ignite.internal.processors.continuous.GridContinuousProcessor;
+import 
org.apache.ignite.internal.processors.continuous.ContinousRoutineLocalInfo;
 import org.apache.ignite.internal.util.typedef.F;
 import org.apache.ignite.internal.util.typedef.internal.CU;
 import org.apache.ignite.testframework.GridTestUtils;
@@ -57,7 +57,7 @@ public class CacheContinuousQueryEntriesExpireTest extends 
GridCommonAbstractTes
         for (int i = 0; i < 1_000; i++)
             cache.put(i, i);
 
-        ConcurrentMap<UUID, GridContinuousProcessor.LocalRoutineInfo> locInfos 
=
+        ConcurrentMap<UUID, ContinousRoutineLocalInfo> locInfos =
             getFieldValue(srv1.context().continuous(), "locInfos");
 
         assertEquals(1, locInfos.size());
diff --git 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/ContinuousQueryRemoteFilterMissingInClassPathSelfTest.java
 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/ContinuousQueryRemoteFilterMissingInClassPathSelfTest.java
index 62873ce51fd..b2b9a27d605 100644
--- 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/ContinuousQueryRemoteFilterMissingInClassPathSelfTest.java
+++ 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/ContinuousQueryRemoteFilterMissingInClassPathSelfTest.java
@@ -28,7 +28,6 @@ import javax.cache.event.CacheEntryListenerException;
 import javax.cache.event.CacheEntryUpdatedListener;
 import org.apache.ignite.Ignite;
 import org.apache.ignite.IgniteCache;
-import org.apache.ignite.IgniteCheckedException;
 import org.apache.ignite.IgniteLogger;
 import org.apache.ignite.cache.CacheEntryEventSerializableFilter;
 import org.apache.ignite.cache.CacheMode;
@@ -37,13 +36,14 @@ import org.apache.ignite.configuration.CacheConfiguration;
 import org.apache.ignite.configuration.IgniteConfiguration;
 import org.apache.ignite.internal.util.typedef.X;
 import org.apache.ignite.testframework.GridStringLogger;
-import org.apache.ignite.testframework.GridTestUtils;
 import org.apache.ignite.testframework.ListeningTestLogger;
 import org.apache.ignite.testframework.LogListener;
 import org.apache.ignite.testframework.config.GridTestProperties;
 import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest;
 import org.junit.Test;
 
+import static 
org.apache.ignite.testframework.GridTestUtils.assertThrowsWithCause;
+
 /**
  *
  */
@@ -110,21 +110,8 @@ public class 
ContinuousQueryRemoteFilterMissingInClassPathSelfTest extends GridC
      * @throws Exception If fail.
      */
     @Test
-    public void testClientJoinsMissingClassWarning() throws Exception {
-        setExternalLoader = true;
-        Ignite ignite0 = startGrid(1);
-
-        executeContinuousQuery(ignite0.cache(DEFAULT_CACHE_NAME));
-
-        log = new GridStringLogger();
-        setExternalLoader = false;
-
-        startClientGrid(2);
-
-        String logStr = log.toString();
-
-        assertTrue(logStr.contains("Failed to unmarshal continuous query 
remote filter on client node. " +
-            "Can be ignored.") || logStr.contains("Failed to unmarshal 
continuous routine handler"));
+    public void testClientJoinsMissingClass() throws Exception {
+        doTestNodeJoinsWithNoClassLoaderForContinuousQuery(false);
     }
 
     /**
@@ -151,6 +138,11 @@ public class 
ContinuousQueryRemoteFilterMissingInClassPathSelfTest extends GridC
      */
     @Test
     public void testServerJoinsMissingClassException() throws Exception {
+        doTestNodeJoinsWithNoClassLoaderForContinuousQuery(true);
+    }
+
+    /** */
+    private void doTestNodeJoinsWithNoClassLoaderForContinuousQuery(boolean 
server) throws Exception {
         setExternalLoader = true;
         Ignite ignite0 = startGrid(1);
 
@@ -160,19 +152,30 @@ public class 
ContinuousQueryRemoteFilterMissingInClassPathSelfTest extends GridC
 
         log = listeningLog;
 
-        LogListener lsnr = LogListener.matches(logStr ->
-            logStr.contains("class org.apache.ignite.IgniteCheckedException: " 
+
-                "Failed to find class with given class loader for 
unmarshalling")
-                || logStr.contains("Failed to unmarshal continuous routine 
handler"
-            )).build();
+        LogListener lsnr1 = LogListener.matches("Failed to initialize a 
continuous query").build();
+        LogListener lsnr2 = LogListener.matches("ClassNotFoundException: " + 
EXT_FILTER_CLASS).build();
 
-        listeningLog.registerListener(lsnr);
+        listeningLog.registerListener(lsnr1);
+        listeningLog.registerListener(lsnr2);
 
         setExternalLoader = false;
 
-        GridTestUtils.assertThrows(log, () -> startGrid(2), 
IgniteCheckedException.class, "Failed to start");
+        if (server) {
+            assertThrowsWithCause(() -> startGrid(2), 
ClassNotFoundException.class);
 
-        assertTrue(lsnr.check());
+            assertTrue(lsnr1.check());
+            assertTrue(lsnr2.check());
+        }
+        else {
+            startClientGrid(2);
+
+            /**
+             * Client successfuly starts because continous query isn't 
actually deployed on a client. It is filtered out
+             * by {@link AttributeNodeFilter} with {@link 
IgniteNodeAttributes#ATTR_CLIENT_MODE}.
+             **/
+            assertFalse(lsnr1.check());
+            assertFalse(lsnr2.check());
+        }
     }
 
     /**
diff --git 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/DiscoveryDataDeserializationFailureHanderTest.java
 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/DiscoveryDataDeserializationFailureHanderTest.java
index e52d06fa0d3..f0a0c281782 100644
--- 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/DiscoveryDataDeserializationFailureHanderTest.java
+++ 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/DiscoveryDataDeserializationFailureHanderTest.java
@@ -57,6 +57,7 @@ public class DiscoveryDataDeserializationFailureHanderTest 
extends GridCommonAbs
     @Test
     public void testFailureHander() throws Exception {
         CountDownLatch latch = new CountDownLatch(1);
+
         failureHnd = new TestFailureHandler(latch);
 
         Ignite node1 = startGrid(1);
diff --git 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/GridCacheContinuousQueryNodesFilteringTest.java
 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/GridCacheContinuousQueryNodesFilteringTest.java
index 7688a7b12ab..7efdad76be4 100644
--- 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/GridCacheContinuousQueryNodesFilteringTest.java
+++ 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/query/continuous/GridCacheContinuousQueryNodesFilteringTest.java
@@ -34,6 +34,7 @@ import org.apache.ignite.configuration.CacheConfiguration;
 import org.apache.ignite.configuration.IgniteConfiguration;
 import org.apache.ignite.failure.FailureHandler;
 import org.apache.ignite.failure.TestFailureHandler;
+import org.apache.ignite.internal.util.typedef.X;
 import org.apache.ignite.lang.IgnitePredicate;
 import org.apache.ignite.testframework.GridStringLogger;
 import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest;
@@ -74,9 +75,17 @@ public class GridCacheContinuousQueryNodesFilteringTest 
extends GridCommonAbstra
             IgniteConfiguration node2Cfg = getConfiguration("node2", true, 
null)
                 .setFailureHandler(failHnd);
 
-            try (Ignite node2 = startGrid(node2Cfg)) {
-                assertTrue("Failure handler hasn't been invoked on the joined 
node.",
-                    latch.await(5, TimeUnit.SECONDS));
+            try {
+                startGrid(node2Cfg);
+            }
+            catch (Throwable t) {
+                // Node start failure with unknow class of continous query 
filter is ok.
+                if (!X.hasCause(t, "Class not found for continuous query 
remote filter [name=" + ENTRY_FILTER_CLS_NAME,
+                    IgniteException.class))
+                    throw t;
+            }
+            finally {
+                assertTrue("Failure handler hasn't been invoked on the joined 
node.", latch.await(5, TimeUnit.SECONDS));
             }
         }
     }
diff --git 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java
 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java
index 47b107edeee..c16bf257775 100644
--- 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java
+++ 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/continuous/GridEventConsumeSelfTest.java
@@ -63,7 +63,6 @@ import static 
org.apache.ignite.events.EventType.EVT_JOB_FINISHED;
 import static org.apache.ignite.events.EventType.EVT_JOB_STARTED;
 import static org.apache.ignite.events.EventType.EVT_NODE_FAILED;
 import static org.apache.ignite.events.EventType.EVT_NODE_LEFT;
-import static 
org.apache.ignite.internal.processors.continuous.GridContinuousProcessor.LocalRoutineInfo;
 import static org.apache.ignite.testframework.GridTestUtils.noop;
 
 /**
@@ -157,10 +156,10 @@ public class GridEventConsumeSelfTest extends 
GridCommonAbstractTest {
      * @param proc Continuous processor.
      * @return Local event routines.
      */
-    private Collection<LocalRoutineInfo> localRoutines(GridContinuousProcessor 
proc) {
-        return F.view(U.<Map<UUID, LocalRoutineInfo>>field(proc, 
"locInfos").values(),
-            new IgnitePredicate<LocalRoutineInfo>() {
-                @Override public boolean apply(LocalRoutineInfo info) {
+    private Collection<ContinousRoutineLocalInfo> 
localRoutines(GridContinuousProcessor proc) {
+        return F.view(U.<Map<UUID, ContinousRoutineLocalInfo>>field(proc, 
"locInfos").values(),
+            new IgnitePredicate<>() {
+                @Override public boolean apply(ContinousRoutineLocalInfo info) 
{
                     return info.handler().isEvents();
                 }
             });
diff --git 
a/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java
 
b/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java
index eabb741a605..87799c52e3c 100644
--- 
a/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java
+++ 
b/modules/zookeeper/src/main/java/org/apache/ignite/spi/discovery/zk/internal/ZookeeperDiscoveryImpl.java
@@ -1156,7 +1156,7 @@ public class ZookeeperDiscoveryImpl {
     }
 
     /**
-     * Marshalls credentials with discovery SPI marshaller (will replace 
attribute value).
+     * Marshals credentials with discovery SPI marshaller (will replace 
attribute value).
      *
      * @param node Node to marshall credentials for.
      * @throws IgniteSpiException If marshalling failed.

Reply via email to