Repository: ignite Updated Branches: refs/heads/ignite-5658 19a2147cd -> 647be5693
wip on data streamer Project: http://git-wip-us.apache.org/repos/asf/ignite/repo Commit: http://git-wip-us.apache.org/repos/asf/ignite/commit/647be569 Tree: http://git-wip-us.apache.org/repos/asf/ignite/tree/647be569 Diff: http://git-wip-us.apache.org/repos/asf/ignite/diff/647be569 Branch: refs/heads/ignite-5658 Commit: 647be569357411069509cfa6b4d5c5dbe662de3a Parents: 19a2147 Author: Yakov Zhdanov <[email protected]> Authored: Fri Jul 14 14:57:21 2017 +0300 Committer: Yakov Zhdanov <[email protected]> Committed: Fri Jul 14 14:57:21 2017 +0300 ---------------------------------------------------------------------- .../ignite/codegen/MessageCodeGenerator.java | 2 - .../managers/communication/GridIoMessage.java | 3 + .../datastreamer/DataStreamerImpl.java | 26 +++++---- .../datastreamer/DataStreamerRequest.java | 59 +++++++++++++++----- .../datastreamer/DataStreamerImplSelfTest.java | 5 +- 5 files changed, 66 insertions(+), 29 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/ignite/blob/647be569/modules/codegen/src/main/java/org/apache/ignite/codegen/MessageCodeGenerator.java ---------------------------------------------------------------------- diff --git a/modules/codegen/src/main/java/org/apache/ignite/codegen/MessageCodeGenerator.java b/modules/codegen/src/main/java/org/apache/ignite/codegen/MessageCodeGenerator.java index 99ec08a..7ff9f27 100644 --- a/modules/codegen/src/main/java/org/apache/ignite/codegen/MessageCodeGenerator.java +++ b/modules/codegen/src/main/java/org/apache/ignite/codegen/MessageCodeGenerator.java @@ -44,8 +44,6 @@ import org.apache.ignite.internal.GridDirectCollection; import org.apache.ignite.internal.GridDirectMap; import org.apache.ignite.internal.GridDirectTransient; import org.apache.ignite.internal.IgniteCodeGeneratingFail; -import org.apache.ignite.internal.processors.cache.version.GridCacheVersion; -import org.apache.ignite.internal.processors.cache.version.GridCacheVersionEx; import org.apache.ignite.internal.util.typedef.internal.SB; import org.apache.ignite.internal.util.typedef.internal.U; import org.apache.ignite.lang.IgniteUuid; http://git-wip-us.apache.org/repos/asf/ignite/blob/647be569/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java ---------------------------------------------------------------------- diff --git a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java index dccd336..c17717a 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/managers/communication/GridIoMessage.java @@ -23,6 +23,7 @@ import java.nio.ByteBuffer; import org.apache.ignite.internal.ExecutorAwareMessage; import org.apache.ignite.internal.GridDirectTransient; import org.apache.ignite.internal.processors.cache.GridCacheMessage; +import org.apache.ignite.internal.processors.datastreamer.DataStreamerRequest; 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; @@ -336,6 +337,8 @@ public class GridIoMessage implements Message { public int partition() { if (msg instanceof GridCacheMessage) return ((GridCacheMessage)msg).partition(); + if (msg instanceof DataStreamerRequest) + return ((DataStreamerRequest)msg).partition(); // TODO introduce partitionAware interface. else return STRIPE_DISABLED_PART; } http://git-wip-us.apache.org/repos/asf/ignite/blob/647be569/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImpl.java ---------------------------------------------------------------------- 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 b7c359c..5d1b0a3 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 @@ -62,6 +62,7 @@ import org.apache.ignite.internal.IgniteInternalFuture; import org.apache.ignite.internal.IgniteInterruptedCheckedException; import org.apache.ignite.internal.cluster.ClusterTopologyCheckedException; import org.apache.ignite.internal.cluster.ClusterTopologyServerNotFoundException; +import org.apache.ignite.internal.managers.communication.GridIoPolicy; import org.apache.ignite.internal.managers.communication.GridMessageListener; import org.apache.ignite.internal.managers.deployment.GridDeployment; import org.apache.ignite.internal.managers.eventstorage.GridLocalEventListener; @@ -1439,7 +1440,7 @@ public class DataStreamerImpl<K, V> implements IgniteDataStreamer<K, V>, Delayed "[batchTopVer=" + curBatchTopVer + ", topVer=" + topVer + "]")); } else if (entries0 != null) { - submit(entries0, curBatchTopVer, curFut0, remap); + submit(entries0, curBatchTopVer, curFut0, remap, b.partId); if (cancelled) curFut0.onDone(new IgniteCheckedException("Data streamer has been cancelled: " + @@ -1483,7 +1484,7 @@ public class DataStreamerImpl<K, V> implements IgniteDataStreamer<K, V>, Delayed } if (entries0 != null) - submit(entries0, batchTopVer, curFut0, false); + submit(entries0, batchTopVer, curFut0, false, b.partId); } // Create compound future for this flush. @@ -1624,15 +1625,16 @@ public class DataStreamerImpl<K, V> implements IgniteDataStreamer<K, V>, Delayed * @param topVer Topology version. * @param curFut Current future. * @param remap Remapping flag. + * @param partId Partition ID. * @throws IgniteInterruptedCheckedException If interrupted. */ private void submit( final Collection<DataStreamerEntry> entries, @Nullable AffinityTopologyVersion topVer, final GridFutureAdapter<Object> curFut, - boolean remap - ) - throws IgniteInterruptedCheckedException { + boolean remap, + int partId + ) throws IgniteInterruptedCheckedException { assert entries != null; assert !entries.isEmpty(); assert curFut != null; @@ -1715,7 +1717,7 @@ public class DataStreamerImpl<K, V> implements IgniteDataStreamer<K, V>, Delayed reqId, topicBytes, cacheName, - updaterBytes, + updaterBytes, // TODO why do we always send updater bytes? entries, true, skipStore, @@ -1726,10 +1728,12 @@ public class DataStreamerImpl<K, V> implements IgniteDataStreamer<K, V>, Delayed dep != null ? dep.participants() : null, dep != null ? dep.classLoaderId() : null, dep == null, - topVer); + topVer, + rcvr == ISOLATED_UPDATER ? partId : -1); try { - ctx.io().sendToGridTopic(node, TOPIC_DATASTREAM, req, plc); + ctx.io().sendToGridTopic(node, TOPIC_DATASTREAM, req, + partId == -1 ? plc : GridIoPolicy.SYSTEM_POOL); if (log.isDebugEnabled()) log.debug("Sent request to node [nodeId=" + node.id() + ", req=" + req + ']'); @@ -1948,8 +1952,10 @@ public class DataStreamerImpl<K, V> implements IgniteDataStreamer<K, V>, Delayed private static final long serialVersionUID = 0L; /** {@inheritDoc} */ - @Override public void receive(IgniteCache<KeyCacheObject, CacheObject> cache, - Collection<Map.Entry<KeyCacheObject, CacheObject>> entries) { + @Override public void receive( + IgniteCache<KeyCacheObject, CacheObject> cache, + Collection<Map.Entry<KeyCacheObject, CacheObject>> entries + ) { IgniteCacheProxy<KeyCacheObject, CacheObject> proxy = (IgniteCacheProxy<KeyCacheObject, CacheObject>)cache; GridCacheAdapter<KeyCacheObject, CacheObject> internalCache = proxy.context().cache(); http://git-wip-us.apache.org/repos/asf/ignite/blob/647be569/modules/core/src/main/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerRequest.java ---------------------------------------------------------------------- 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 b4cbf66..f70ee9c 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 @@ -90,6 +90,9 @@ public class DataStreamerRequest implements Message { /** Topology version. */ private AffinityTopologyVersion topVer; + /** */ + private int partId; + /** * {@code Externalizable} support. */ @@ -113,8 +116,10 @@ public class DataStreamerRequest implements Message { * @param clsLdrId Class loader ID. * @param forceLocDep Force local deployment. * @param topVer Topology version. + * @param partId Partition ID. */ - public DataStreamerRequest(long reqId, + public DataStreamerRequest( + long reqId, byte[] resTopicBytes, @Nullable String cacheName, byte[] updaterBytes, @@ -128,7 +133,9 @@ public class DataStreamerRequest implements Message { Map<UUID, IgniteUuid> ldrParticipants, IgniteUuid clsLdrId, boolean forceLocDep, - @NotNull AffinityTopologyVersion topVer) { + @NotNull AffinityTopologyVersion topVer, + int partId + ) { assert topVer != null; this.reqId = reqId; @@ -146,6 +153,7 @@ public class DataStreamerRequest implements Message { this.clsLdrId = clsLdrId; this.forceLocDep = forceLocDep; this.topVer = topVer; + this.partId = partId; } /** @@ -253,6 +261,13 @@ public class DataStreamerRequest implements Message { return topVer; } + /** + * @return Partition ID. + */ + public int partition() { + return partId; + } + /** {@inheritDoc} */ @Override public void onAckReceived() { // No-op. @@ -324,42 +339,48 @@ public class DataStreamerRequest implements Message { writer.incrementState(); case 8: - if (!writer.writeLong("reqId", reqId)) + if (!writer.writeInt("partId", partId)) return false; writer.incrementState(); case 9: - if (!writer.writeByteArray("resTopicBytes", resTopicBytes)) + if (!writer.writeLong("reqId", reqId)) return false; writer.incrementState(); case 10: - if (!writer.writeString("sampleClsName", sampleClsName)) + if (!writer.writeByteArray("resTopicBytes", resTopicBytes)) return false; writer.incrementState(); case 11: - if (!writer.writeBoolean("skipStore", skipStore)) + if (!writer.writeString("sampleClsName", sampleClsName)) return false; writer.incrementState(); case 12: - if (!writer.writeMessage("topVer", topVer)) + if (!writer.writeBoolean("skipStore", skipStore)) return false; writer.incrementState(); case 13: - if (!writer.writeByteArray("updaterBytes", updaterBytes)) + if (!writer.writeMessage("topVer", topVer)) return false; writer.incrementState(); case 14: + if (!writer.writeByteArray("updaterBytes", updaterBytes)) + return false; + + writer.incrementState(); + + case 15: if (!writer.writeString("userVer", userVer)) return false; @@ -447,7 +468,7 @@ public class DataStreamerRequest implements Message { reader.incrementState(); case 8: - reqId = reader.readLong("reqId"); + partId = reader.readInt("partId"); if (!reader.isLastRead()) return false; @@ -455,7 +476,7 @@ public class DataStreamerRequest implements Message { reader.incrementState(); case 9: - resTopicBytes = reader.readByteArray("resTopicBytes"); + reqId = reader.readLong("reqId"); if (!reader.isLastRead()) return false; @@ -463,7 +484,7 @@ public class DataStreamerRequest implements Message { reader.incrementState(); case 10: - sampleClsName = reader.readString("sampleClsName"); + resTopicBytes = reader.readByteArray("resTopicBytes"); if (!reader.isLastRead()) return false; @@ -471,7 +492,7 @@ public class DataStreamerRequest implements Message { reader.incrementState(); case 11: - skipStore = reader.readBoolean("skipStore"); + sampleClsName = reader.readString("sampleClsName"); if (!reader.isLastRead()) return false; @@ -479,7 +500,7 @@ public class DataStreamerRequest implements Message { reader.incrementState(); case 12: - topVer = reader.readMessage("topVer"); + skipStore = reader.readBoolean("skipStore"); if (!reader.isLastRead()) return false; @@ -487,7 +508,7 @@ public class DataStreamerRequest implements Message { reader.incrementState(); case 13: - updaterBytes = reader.readByteArray("updaterBytes"); + topVer = reader.readMessage("topVer"); if (!reader.isLastRead()) return false; @@ -495,6 +516,14 @@ public class DataStreamerRequest implements Message { reader.incrementState(); case 14: + updaterBytes = reader.readByteArray("updaterBytes"); + + if (!reader.isLastRead()) + return false; + + reader.incrementState(); + + case 15: userVer = reader.readString("userVer"); if (!reader.isLastRead()) @@ -514,6 +543,6 @@ public class DataStreamerRequest implements Message { /** {@inheritDoc} */ @Override public byte fieldsCount() { - return 15; + return 16; } } http://git-wip-us.apache.org/repos/asf/ignite/blob/647be569/modules/core/src/test/java/org/apache/ignite/internal/processors/datastreamer/DataStreamerImplSelfTest.java ---------------------------------------------------------------------- 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 e72a9b4..ca0e150 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 @@ -346,7 +346,8 @@ public class DataStreamerImplSelfTest extends GridCommonAbstractTest { req.participants(), req.classLoaderId(), req.forceLocalDeployment(), - staleTop); + staleTop, + -1); msg = new GridIoMessage( GridTestUtils.<Byte>getFieldValue(ioMsg, "plc"), @@ -365,4 +366,4 @@ public class DataStreamerImplSelfTest extends GridCommonAbstractTest { super.sendMessage(node, msg, ackC); } } -} \ No newline at end of file +}
