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(),