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

anton-vinogradov pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ignite.git


The following commit(s) were added to refs/heads/master by this push:
     new 54159731e86 IGNITE-27977 Refactor bytes serialization for 
DataStreamerRequest (#13454)
54159731e86 is described below

commit 54159731e86c904397e507d39e6849f42247625b
Author: Anton Vinogradov <[email protected]>
AuthorDate: Mon Aug 10 23:30:05 2026 +0300

    IGNITE-27977 Refactor bytes serialization for DataStreamerRequest (#13454)
---
 .../ignite/internal/CoreMessagesProvider.java      |  2 +
 .../datastreamer/DataStreamProcessor.java          | 13 ++-
 .../datastreamer/DataStreamerBuiltInUpdater.java   | 72 +++++++++++++++++
 .../processors/datastreamer/DataStreamerImpl.java  | 53 +++++++-----
 .../datastreamer/DataStreamerReceiverMessage.java  | 68 ++++++++++++++++
 .../datastreamer/DataStreamerRequest.java          | 25 +++---
 .../datastreamer/DataStreamerImplSelfTest.java     | 94 +++++++++++++++++++++-
 7 files changed, 289 insertions(+), 38 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 3b066879eea..6513e174962 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
@@ -208,6 +208,7 @@ import 
org.apache.ignite.internal.processors.continuous.StartRoutineDiscoveryMes
 import 
org.apache.ignite.internal.processors.continuous.StopRoutineAckDiscoveryMessage;
 import 
org.apache.ignite.internal.processors.continuous.StopRoutineDiscoveryMessage;
 import org.apache.ignite.internal.processors.datastreamer.DataStreamerEntry;
+import 
org.apache.ignite.internal.processors.datastreamer.DataStreamerReceiverMessage;
 import org.apache.ignite.internal.processors.datastreamer.DataStreamerRequest;
 import org.apache.ignite.internal.processors.datastreamer.DataStreamerResponse;
 import org.apache.ignite.internal.processors.marshaller.MappedName;
@@ -662,6 +663,7 @@ public class CoreMessagesProvider extends 
AbstractMessageFactoryProvider {
         register(DataStreamerEntry.class);
         register(DataStreamerRequest.class);
         register(DataStreamerResponse.class);
+        register(DataStreamerReceiverMessage.class);
 
         // [11900 - 12000]: Metrics, monitoring messages.
         msgIdx = 11900;
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamProcessor.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamProcessor.java
index 92ec36cc210..42fa9112a2f 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamProcessor.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamProcessor.java
@@ -28,6 +28,7 @@ import 
org.apache.ignite.internal.IgniteInterruptedCheckedException;
 import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException;
 import org.apache.ignite.internal.managers.communication.GridIoManager;
 import org.apache.ignite.internal.managers.communication.GridMessageListener;
+import org.apache.ignite.internal.managers.communication.MessageMarshalling;
 import org.apache.ignite.internal.managers.deployment.GridDeployment;
 import org.apache.ignite.internal.processors.GridProcessorAdapter;
 import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
@@ -49,7 +50,6 @@ import 
org.apache.ignite.internal.util.worker.queue.IgniteDelayedObjectHandler;
 import org.apache.ignite.lang.IgniteClosure;
 import org.apache.ignite.lang.IgniteFuture;
 import org.apache.ignite.lang.IgniteInClosure;
-import org.apache.ignite.marshaller.Marshaller;
 import org.apache.ignite.stream.StreamReceiver;
 import org.jetbrains.annotations.Nullable;
 
@@ -69,9 +69,6 @@ public class DataStreamProcessor extends GridProcessorAdapter 
{
     /** Data Streamer flusher. */
     private final DataStreamerFlusher flusher = new DataStreamerFlusher();
 
-    /** Marshaller. */
-    private final Marshaller marsh;
-
     /**
      * @param ctx Kernal context.
      */
@@ -87,8 +84,6 @@ public class DataStreamProcessor extends GridProcessorAdapter 
{
                 }
             });
         }
-
-        marsh = ctx.marshaller();
     }
 
     /** {@inheritDoc} */
@@ -235,9 +230,11 @@ public class DataStreamProcessor extends 
GridProcessorAdapter {
             StreamReceiver<?, ?> updater;
 
             try {
-                updater = U.unmarshal(marsh, req.updaterBytes(), 
U.resolveClassLoader(clsLdr, ctx.config()));
+                MessageMarshalling.unmarshal(req, ctx, null, 
U.resolveClassLoader(clsLdr, ctx.config()));
+
+                updater = req.updater();
 
-                if (updater != null)
+                if (req.customUpdater())
                     ctx.resource().injectGeneric(updater);
             }
             catch (IgniteCheckedException e) {
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerBuiltInUpdater.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerBuiltInUpdater.java
new file mode 100644
index 00000000000..a98122e8c48
--- /dev/null
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerBuiltInUpdater.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.datastreamer;
+
+import org.apache.ignite.stream.StreamReceiver;
+import org.jetbrains.annotations.Nullable;
+
+/**
+ * The built-in updaters. Every node has them, so a request names the one it 
needs instead of carrying a serialized
+ * copy.
+ */
+enum DataStreamerBuiltInUpdater {
+    /** {@link DataStreamerImpl#ISOLATED_UPDATER}. */
+    ISOLATED(DataStreamerImpl.ISOLATED_UPDATER),
+
+    /** {@link DataStreamerCacheUpdaters#individual()}. */
+    INDIVIDUAL(DataStreamerCacheUpdaters.individual()),
+
+    /** {@link DataStreamerCacheUpdaters#batched()}. */
+    BATCHED(DataStreamerCacheUpdaters.batched()),
+
+    /** {@link DataStreamerCacheUpdaters#batchedSorted()}. */
+    BATCHED_SORTED(DataStreamerCacheUpdaters.batchedSorted());
+
+    /** */
+    private final StreamReceiver<?, ?> updater;
+
+    /** @param updater Updater this constant stands for. */
+    DataStreamerBuiltInUpdater(StreamReceiver<?, ?> updater) {
+        this.updater = updater;
+    }
+
+    /** @return Updater of this node. */
+    StreamReceiver<?, ?> updater() {
+        return updater;
+    }
+
+    /** @return New message naming this updater. */
+    DataStreamerReceiverMessage message() {
+        return new DataStreamerReceiverMessage(this);
+    }
+
+    /**
+     * Matches by class, so an updater built anew is still recognized as 
built-in: these updaters hold no state.
+     *
+     * @param updater Updater to look up.
+     * @return Constant standing for {@code updater}, or {@code null} when it 
is a custom one.
+     */
+    static @Nullable DataStreamerBuiltInUpdater of(StreamReceiver<?, ?> 
updater) {
+        for (DataStreamerBuiltInUpdater builtIn : values()) {
+            if (builtIn.updater.getClass() == updater.getClass())
+                return builtIn;
+        }
+
+        return null;
+    }
+}
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java
index ec42ac560ae..c9115f309bc 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java
@@ -146,17 +146,14 @@ public class DataStreamerImpl<K, V> implements 
IgniteDataStreamer<K, V>, Delayed
      */
     private final Map<Long, ThreadBuffer> threadBufMap = new 
ConcurrentHashMap<>();
 
-    /** Isolated receiver. */
-    private static final StreamReceiver ISOLATED_UPDATER = new 
IsolatedUpdater();
+    /** Default, Isolated receiver. */
+    static final StreamReceiver ISOLATED_UPDATER = new IsolatedUpdater();
 
     /** Amount of permissions should be available to continue new data 
processing. */
     private static final int REMAP_SEMAPHORE_PERMISSIONS_COUNT = 
Integer.MAX_VALUE;
 
-    /** Cache receiver. */
-    private StreamReceiver<K, V> rcvr = ISOLATED_UPDATER;
-
-    /** */
-    private byte[] updaterBytes;
+    /** Message of the cache receiver; {@code null} until a receiver is set, 
the Isolated updater being used so far. */
+    private volatile @Nullable DataStreamerReceiverMessage rcvrMsg;
 
     /** IO policy resovler for data load request. */
     private IgniteClosure<ClusterNode, Byte> ioPlcRslvr;
@@ -489,12 +486,29 @@ public class DataStreamerImpl<K, V> implements 
IgniteDataStreamer<K, V>, Delayed
     @Override public void receiver(StreamReceiver<K, V> rcvr) {
         A.notNull(rcvr, "rcvr");
 
-        this.rcvr = rcvr;
+        DataStreamerBuiltInUpdater builtIn = 
DataStreamerBuiltInUpdater.of(rcvr);
+
+        rcvrMsg = builtIn != null ? builtIn.message() : new 
DataStreamerReceiverMessage(rcvr);
+    }
+
+    /** @return Message of the receiver in use, one of the Isolated updater 
until a receiver is set. */
+    private DataStreamerReceiverMessage receiverMessage() {
+        DataStreamerReceiverMessage rcvrMsg0 = rcvrMsg;
+
+        return rcvrMsg0 != null ? rcvrMsg0 : 
DataStreamerBuiltInUpdater.ISOLATED.message();
+    }
+
+    /** @return Cache receiver, custom or built-in, the Isolated updater until 
one is set. */
+    @SuppressWarnings("unchecked")
+    private StreamReceiver<K, V> receiver() {
+        DataStreamerReceiverMessage rcvrMsg0 = rcvrMsg;
+
+        return (StreamReceiver<K, V>)(rcvrMsg0 == null ? ISOLATED_UPDATER : 
rcvrMsg0.receiver());
     }
 
     /** {@inheritDoc} */
     @Override public boolean allowOverwrite() {
-        return rcvr != ISOLATED_UPDATER;
+        return receiver() != ISOLATED_UPDATER;
     }
 
     /** {@inheritDoc} */
@@ -507,7 +521,7 @@ public class DataStreamerImpl<K, V> implements 
IgniteDataStreamer<K, V>, Delayed
         if (node == null)
             throw new CacheException("Failed to get node for cache: " + 
cacheName);
 
-        rcvr = allow ? DataStreamerCacheUpdaters.<K, V>individual() : 
ISOLATED_UPDATER;
+        receiver(allow ? DataStreamerCacheUpdaters.individual() : 
ISOLATED_UPDATER);
     }
 
     /** {@inheritDoc} */
@@ -655,7 +669,8 @@ public class DataStreamerImpl<K, V> implements 
IgniteDataStreamer<K, V>, Delayed
 
         lock(false);
 
-        if (rcvr instanceof IsolatedUpdater && 
inconsistencyWarned.compareAndSet(false, true))
+        // Without overwrite the Isolated updater is in use, and it is the one 
the warning is about.
+        if (!allowOverwrite() && inconsistencyWarned.compareAndSet(false, 
true))
             log.warning(WRN_INCONSISTENT_UPDATES);
 
         try {
@@ -886,6 +901,8 @@ public class DataStreamerImpl<K, V> implements 
IgniteDataStreamer<K, V>, Delayed
                         assert key != null;
 
                         if (initPda) {
+                            StreamReceiver<K, V> rcvr = receiver();
+
                             if (cacheObjCtx.addDeploymentInfo())
                                 jobPda = new 
DataStreamerPda(key.value(cacheObjCtx, false),
                                     entry.getValue() != null ? 
entry.getValue().value(cacheObjCtx, false) : null,
@@ -1850,7 +1867,7 @@ public class DataStreamerImpl<K, V> implements 
IgniteDataStreamer<K, V>, Delayed
                                 false,
                                 skipStore,
                                 keepBinary,
-                                rcvr),
+                                receiver()),
                             plc);
 
                         locFuts.add(callFut);
@@ -1943,12 +1960,6 @@ public class DataStreamerImpl<K, V> implements 
IgniteDataStreamer<K, V>, Delayed
                         if (val != null)
                             val.marshal(cacheObjCtx);
                     }
-
-                    if (updaterBytes == null) {
-                        assert rcvr != null;
-
-                        updaterBytes = U.marshal(ctx, rcvr);
-                    }
                 }
                 catch (IgniteCheckedException e) {
                     U.error(log, "Failed to marshal.", e);
@@ -1991,11 +2002,13 @@ public class DataStreamerImpl<K, V> implements 
IgniteDataStreamer<K, V>, Delayed
                 if (topVer == null)
                     topVer = 
ctx.cache().context().exchange().readyAffinityVersion();
 
+                DataStreamerReceiverMessage rcvrMsg0 = receiverMessage();
+
                 DataStreamerRequest req = new DataStreamerRequest(
                     reqId,
                     topicId,
                     cacheName,
-                    updaterBytes,
+                    rcvrMsg0,
                     entries,
                     true,
                     skipStore,
@@ -2004,7 +2017,7 @@ public class DataStreamerImpl<K, V> implements 
IgniteDataStreamer<K, V>, Delayed
                     dep != null ? jobPda0.deployClass().getName() : null,
                     dep == null,
                     topVer,
-                    (rcvr == ISOLATED_UPDATER) ? partId : NO_STRIPE);
+                    rcvrMsg0.receiver() == ISOLATED_UPDATER ? partId : 
NO_STRIPE);
 
                 try {
                     ctx.io().sendToGridTopic(node, TOPIC_DATASTREAM, req, plc);
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerReceiverMessage.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerReceiverMessage.java
new file mode 100644
index 00000000000..6eafdb8b5b0
--- /dev/null
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerReceiverMessage.java
@@ -0,0 +1,68 @@
+/*
+ * 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.datastreamer;
+
+import org.apache.ignite.internal.Marshalled;
+import org.apache.ignite.internal.Order;
+import org.apache.ignite.internal.UseBinaryMarshaller;
+import org.apache.ignite.plugin.extensions.communication.Message;
+import org.apache.ignite.stream.StreamReceiver;
+import org.jetbrains.annotations.Nullable;
+
+/** DataStreamer cache receiver/updater message. */
+@UseBinaryMarshaller
+public class DataStreamerReceiverMessage implements Message {
+    /** Custom cache receiver/updater; {@code null} when {@link #builtIn} is 
effective. */
+    @Marshalled("rcvrBytes")
+    @Nullable StreamReceiver<?, ?> rcvr;
+
+    /** Serialized {@link #rcvr}; {@code null} when {@link #builtIn} is 
effective. */
+    @Order(0)
+    volatile @Nullable byte[] rcvrBytes;
+
+    /** A built-in updater every node has; {@code null} when {@link #rcvr} is 
effective. */
+    @Order(1)
+    @Nullable DataStreamerBuiltInUpdater builtIn;
+
+    /** Empty constructor for serialization purposes. */
+    public DataStreamerReceiverMessage() {
+        // No-op.
+    }
+
+    /** @param rcvr Custom receiver. */
+    DataStreamerReceiverMessage(StreamReceiver<?, ?> rcvr) {
+        assert DataStreamerBuiltInUpdater.of(rcvr) == null : "A built-in 
updater travels by name: " + rcvr;
+
+        this.rcvr = rcvr;
+    }
+
+    /** @param builtIn Built-in updater, sent as a name rather than as a copy. 
*/
+    DataStreamerReceiverMessage(DataStreamerBuiltInUpdater builtIn) {
+        this.builtIn = builtIn;
+    }
+
+    /** @return {@code True} if this is a custom receiver, {@code false} if a 
built-in one. */
+    boolean customUpdater() {
+        return builtIn == null;
+    }
+
+    /** @return Receiver: the custom one sent here, or the built-in one of 
this node. */
+    StreamReceiver<?, ?> receiver() {
+        return builtIn == null ? rcvr : builtIn.updater();
+    }
+}
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java
index 6d941ce969a..f690eac5f38 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java
@@ -25,15 +25,17 @@ import 
org.apache.ignite.internal.managers.deployment.GridDeploymentInfo;
 import 
org.apache.ignite.internal.managers.deployment.GridDeploymentInfoMessage;
 import org.apache.ignite.internal.processors.affinity.AffinityTopologyVersion;
 import org.apache.ignite.internal.processors.cache.GridCacheUtils;
+import org.apache.ignite.internal.util.tostring.GridToStringExclude;
 import org.apache.ignite.internal.util.typedef.internal.S;
 import org.apache.ignite.lang.IgniteUuid;
 import org.apache.ignite.plugin.extensions.communication.CacheIdAware;
+import org.apache.ignite.stream.StreamReceiver;
 import org.jetbrains.annotations.NotNull;
 import org.jetbrains.annotations.Nullable;
 
 import static org.apache.ignite.internal.GridTopic.TOPIC_DATASTREAM;
 
-/** */
+/** Batch of streamed entries. */
 public class DataStreamerRequest implements DeferredUnmarshalMessage, 
CacheIdAware, StripedMessage {
     /** */
     @Order(0)
@@ -47,10 +49,10 @@ public class DataStreamerRequest implements 
DeferredUnmarshalMessage, CacheIdAwa
     @Order(2)
     String cacheName;
 
-    /** */
-    // TODO: Refactor bytes serialization - IGNITE-27977
+    /** Cache updater. */
+    @GridToStringExclude
     @Order(3)
-    byte[] updaterBytes;
+    DataStreamerReceiverMessage updaterMsg;
 
     /** Entries to update. */
     @Order(4)
@@ -97,7 +99,7 @@ public class DataStreamerRequest implements 
DeferredUnmarshalMessage, CacheIdAwa
      * @param reqId Request ID.
      * @param resTopicId Response topic ID.
      * @param cacheName Cache name.
-     * @param updaterBytes Cache receiver.
+     * @param updaterMsg Cache updater.
      * @param entries Entries to put.
      * @param ignoreDepOwnership Ignore ownership.
      * @param skipStore Skip store flag.
@@ -112,7 +114,7 @@ public class DataStreamerRequest implements 
DeferredUnmarshalMessage, CacheIdAwa
         long reqId,
         IgniteUuid resTopicId,
         @Nullable String cacheName,
-        byte[] updaterBytes,
+        DataStreamerReceiverMessage updaterMsg,
         Collection<DataStreamerEntry> entries,
         boolean ignoreDepOwnership,
         boolean skipStore,
@@ -128,7 +130,7 @@ public class DataStreamerRequest implements 
DeferredUnmarshalMessage, CacheIdAwa
         this.reqId = reqId;
         this.resTopicId = resTopicId;
         this.cacheName = cacheName;
-        this.updaterBytes = updaterBytes;
+        this.updaterMsg = updaterMsg;
         this.entries = entries;
         this.ignoreDepOwnership = ignoreDepOwnership;
         this.skipStore = skipStore;
@@ -155,9 +157,14 @@ public class DataStreamerRequest implements 
DeferredUnmarshalMessage, CacheIdAwa
         return cacheName;
     }
 
+    /** @return {@code True} if the request sends a custom updater rather than 
naming a built-in one. */
+    boolean customUpdater() {
+        return updaterMsg.customUpdater();
+    }
+
     /** @return Updater. */
-    byte[] updaterBytes() {
-        return updaterBytes;
+    StreamReceiver<?, ?> updater() {
+        return updaterMsg.receiver();
     }
 
     /** @return Entries to update. */
diff --git 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java
 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java
index b6e8e1efdd1..f38ccc03b31 100644
--- 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java
+++ 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java
@@ -20,9 +20,12 @@ package org.apache.ignite.internal.processors.datastreamer;
 import java.io.StringWriter;
 import java.util.ArrayList;
 import java.util.Collection;
+import java.util.Collections;
+import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Random;
+import java.util.Set;
 import java.util.concurrent.Callable;
 import java.util.concurrent.CountDownLatch;
 import java.util.concurrent.CyclicBarrier;
@@ -88,6 +91,9 @@ public class DataStreamerImplSelfTest extends 
GridCommonAbstractTest {
     /** Indicates whether we need to make the topology stale */
     private static boolean needStaleTop = false;
 
+    /** Distinct updaters sent since the current test started: the serialized 
bytes, or the built-in constant. */
+    private static final Set<Object> SENT_UPDATERS = 
Collections.synchronizedSet(new HashSet<>());
+
     /** {@inheritDoc} */
     @Override protected void afterTest() throws Exception {
         super.afterTest();
@@ -142,6 +148,71 @@ public class DataStreamerImplSelfTest extends 
GridCommonAbstractTest {
             assertTrue(fut.isDone());
     }
 
+    /**
+     * The receiver does not change between batches, so it is marshalled once: 
every request carries the very bytes
+     * produced for the first one.
+     *
+     * @throws Exception If failed.
+     */
+    @Test
+    public void testReceiverMarshalledOncePerStreamer() throws Exception {
+        startGridsAndStream(new TestReceiver());
+
+        assertEquals("The receiver was marshalled more than once", 1, 
SENT_UPDATERS.size());
+
+        assertTrue("A custom receiver must be sent, not named", 
F.first(SENT_UPDATERS) instanceof byte[]);
+    }
+
+    /**
+     * Every built-in updater is sent as a name rather than as a copy, and the 
data still lands.
+     *
+     * @throws Exception If failed.
+     */
+    @Test
+    public void testBuiltInUpdaterIsNotSent() throws Exception {
+        for (DataStreamerBuiltInUpdater builtIn : 
DataStreamerBuiltInUpdater.values()) {
+            startGridsAndStream(builtIn.updater());
+
+            assertEquals("Expected " + builtIn + " to be sent as a name", 
Collections.singleton(builtIn),
+                SENT_UPDATERS);
+
+            IgniteCache<Object, Object> cache = 
grid(1).cache(DEFAULT_CACHE_NAME);
+
+            for (int i = 0; i < KEYS_COUNT; i++)
+                assertEquals(i, cache.get(i));
+
+            stopAllGrids();
+        }
+    }
+
+    /**
+     * Starts two nodes and streams {@link #KEYS_COUNT} entries from the first 
one, a request per entry, collecting
+     * the updaters they carry, custom or built-in. Waits for the partition 
map first: until it is ready every
+     * partition is primary here, and a streamer that overwrites sends nothing 
to the remote node.
+     *
+     * @param rcvr Receiver to stream with.
+     * @throws Exception If failed.
+     */
+    @SuppressWarnings("unchecked")
+    private void startGridsAndStream(StreamReceiver<?, ?> rcvr) throws 
Exception {
+        cnt = 0;
+
+        startGrids(2);
+
+        awaitPartitionMapExchange();
+
+        SENT_UPDATERS.clear();
+
+        try (IgniteDataStreamer<Object, Object> ldr = 
grid(0).dataStreamer(DEFAULT_CACHE_NAME)) {
+            ldr.receiver((StreamReceiver<Object, Object>)rcvr);
+
+            ldr.perNodeBufferSize(1);
+
+            for (int i = 0; i < KEYS_COUNT; i++)
+                ldr.addData(i, i);
+        }
+    }
+
     /**
      * Test inconsistency log warning of the streamer. Default receiver goes 
first and is set again after a consistent
      * receiver. The warning must appear only once.
@@ -664,12 +735,33 @@ public class DataStreamerImplSelfTest extends 
GridCommonAbstractTest {
         return cacheCfg;
     }
 
+    /** A custom receiver: unlike the built-in ones, it is sent with the 
requests. */
+    private static class TestReceiver implements StreamReceiver<Object, 
Object> {
+        /** */
+        private static final long serialVersionUID = 0L;
+
+        /** {@inheritDoc} */
+        @Override public void receive(IgniteCache<Object, Object> cache, 
Collection<Map.Entry<Object, Object>> entries) {
+            for (Map.Entry<Object, Object> e : entries)
+                cache.put(e.getKey(), e.getValue());
+        }
+    }
+
     /**
      * Simulate stale (not up-to-date) topology
      */
     private static class StaleTopologyCommunicationSpi extends 
TcpCommunicationSpi {
         /** {@inheritDoc} */
         @Override public void sendMessage(ClusterNode node, Message msg, 
IgniteInClosure<IgniteException> ackC) {
+            Message sentMsg = msg instanceof GridIoMessage ? 
((GridIoMessage)msg).message() : null;
+
+            // Already marshalled at this point.
+            if (sentMsg instanceof DataStreamerRequest) {
+                DataStreamerReceiverMessage updaterMsg = 
((DataStreamerRequest)sentMsg).updaterMsg;
+
+                SENT_UPDATERS.add(updaterMsg.customUpdater() ? 
updaterMsg.rcvrBytes : updaterMsg.builtIn);
+            }
+
             // Send stale topology only in the first request to avoid 
indefinitely getting failures.
             if (needStaleTop) {
                 if (msg instanceof GridIoMessage) {
@@ -692,7 +784,7 @@ public class DataStreamerImplSelfTest extends 
GridCommonAbstractTest {
                             req.requestId(),
                             req.resTopicId,
                             req.cacheName(),
-                            req.updaterBytes(),
+                            req.updaterMsg,
                             req.entries(),
                             req.ignoreDeploymentOwnership(),
                             req.skipStore(),

Reply via email to