5578

Project: http://git-wip-us.apache.org/repos/asf/ignite/repo
Commit: http://git-wip-us.apache.org/repos/asf/ignite/commit/c6221718
Tree: http://git-wip-us.apache.org/repos/asf/ignite/tree/c6221718
Diff: http://git-wip-us.apache.org/repos/asf/ignite/diff/c6221718

Branch: refs/heads/ignite-5578
Commit: c622171857ee34d5b680b8eb0b2e6fcf40a0f6bf
Parents: d72ff59
Author: sboikov <[email protected]>
Authored: Thu Jul 13 09:00:14 2017 +0300
Committer: sboikov <[email protected]>
Committed: Thu Jul 13 09:00:14 2017 +0300

----------------------------------------------------------------------
 .../GridDhtPartitionsExchangeFuture.java        |  4 ++--
 .../CacheLateAffinityAssignmentTest.java        | 24 ++++++++++++++++++++
 2 files changed, 26 insertions(+), 2 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/ignite/blob/c6221718/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsExchangeFuture.java
----------------------------------------------------------------------
diff --git 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsExchangeFuture.java
 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsExchangeFuture.java
index ab66df3..a8d1589 100644
--- 
a/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsExchangeFuture.java
+++ 
b/modules/core/src/main/java/org/apache/ignite/internal/processors/cache/distributed/dht/preloader/GridDhtPartitionsExchangeFuture.java
@@ -1123,9 +1123,9 @@ public class GridDhtPartitionsExchangeFuture extends 
GridDhtTopologyFutureAdapte
                 resetLostPartitions(caches);
         }
 
-        if (cctx.kernalContext().clientNode() || (localJoinExchange() && 
!cctx.database().persistenceEnabled())) {
+        if (cctx.kernalContext().clientNode()) {
             msg = new GridDhtPartitionsSingleMessage(exchangeId(),
-                cctx.kernalContext().clientNode(),
+                true,
                 cctx.versions().last(),
                 true);
         }

http://git-wip-us.apache.org/repos/asf/ignite/blob/c6221718/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/CacheLateAffinityAssignmentTest.java
----------------------------------------------------------------------
diff --git 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/CacheLateAffinityAssignmentTest.java
 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/CacheLateAffinityAssignmentTest.java
index 23043d1..840dda1 100644
--- 
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/CacheLateAffinityAssignmentTest.java
+++ 
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/CacheLateAffinityAssignmentTest.java
@@ -37,6 +37,7 @@ import javax.cache.processor.EntryProcessor;
 import javax.cache.processor.MutableEntry;
 import org.apache.ignite.Ignite;
 import org.apache.ignite.IgniteCache;
+import org.apache.ignite.IgniteDataStreamer;
 import org.apache.ignite.IgniteException;
 import org.apache.ignite.IgniteServices;
 import org.apache.ignite.IgniteTransactions;
@@ -1891,6 +1892,29 @@ public class CacheLateAffinityAssignmentTest extends 
GridCommonAbstractTest {
     }
 
     /**
+     * @throws Exception If failed.
+     */
+    public void testStreamer1() throws Exception {
+        cacheC = new IgniteClosure<String, CacheConfiguration[]>() {
+            @Override public CacheConfiguration[] apply(String s) {
+                return null;
+            }
+        };
+
+        startServer(0, 1);
+
+        cacheC = null;
+        cacheNodeFilter = new 
CacheNodeFilter(Collections.singletonList(getTestIgniteInstanceName(0)));
+
+        startServer(1, 2);
+
+        IgniteDataStreamer<Object, Object> streamer = 
ignite(0).dataStreamer(CACHE_NAME1);
+
+        streamer.addData(1, 1);
+        streamer.flush();
+    }
+
+    /**
      * @param cache Cache
      */
     private void cacheOperations(IgniteCache<Object, Object> cache) {

Reply via email to