From 566821ebbf4c5e6b34d463ba4491b6be113b58c5 Mon Sep 17 00:00:00 2001 From: Vladimir Steshin Date: Tue, 28 Jul 2026 16:07:29 +0300 Subject: [PATCH 01/25] raw --- .../main/java/org/apache/ignite/internal/GridComponent.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridComponent.java b/modules/core/src/main/java/org/apache/ignite/internal/GridComponent.java index 95913aec95f1c..d346bbcd6152a 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridComponent.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridComponent.java @@ -130,7 +130,9 @@ enum DiscoveryDataExchangeType { * * @param dataBag container object to store discovery data in. */ - public void collectJoiningNodeData(DiscoveryDataBag dataBag); + default void collectJoiningNodeData(DiscoveryDataBag dataBag) { + // No-op. + } /** * Collects discovery data on nodes already in grid on receiving From 3b10bea30156ab0e852adbbd4376394223325ae5 Mon Sep 17 00:00:00 2001 From: Vladimir Steshin Date: Wed, 29 Jul 2026 01:00:14 +0300 Subject: [PATCH 02/25] impl --- .../ignite/internal/CoreMessagesProvider.java | 6 ++ ...DistributedMetaStorageClusterNodeData.java | 86 ++++++++++++----- .../DistributedMetaStorageHistoryItem.java | 10 ++ ...tributedMetaStorageHistoryItemMessage.java | 67 +++++++++++++ .../DistributedMetaStorageImpl.java | 96 +++++++++---------- ...DistributedMetaStorageJoiningNodeData.java | 59 ++++++++---- .../DistributedMetaStorageVersion.java | 2 +- .../persistence/DmsDataWriter.java | 25 ++--- ...oryCachedDistributedMetaStorageBridge.java | 8 +- .../persistence/DmsDataWriterTest.java | 8 +- 10 files changed, 253 insertions(+), 114 deletions(-) create mode 100644 modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java 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 bba4c1a6ca6b4..0d32590a48913 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 @@ -211,6 +211,9 @@ import org.apache.ignite.internal.processors.marshaller.MissingMappingResponseMessage; import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageCasAckMessage; import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageCasMessage; +import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageClusterNodeData; +import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageHistoryItemMessage; +import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageJoiningNodeData; import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageUpdateAckMessage; import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageUpdateMessage; import org.apache.ignite.internal.processors.plugin.PluginsDataBagItem; @@ -709,6 +712,9 @@ public CoreMessagesProvider(Marshaller dfltMarsh, Marshaller schemaAwareMarsh, C withNoSchema(ClusterUpdateNotifierDataBagItem.class); withNoSchemaResolvedClassLoader(PluginsDataBagItem.class); withSchema(EventsDataBagItem.class); + withSchema(DistributedMetaStorageHistoryItemMessage.class); + withSchema(DistributedMetaStorageJoiningNodeData.class); + withSchema(DistributedMetaStorageClusterNodeData.class); // [13400 - 13500]: Operation context messages. msgIdx = 13400; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java index 01375e055af15..56df3f88d8ec8 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java @@ -17,51 +17,85 @@ package org.apache.ignite.internal.processors.metastorage.persistence; -import java.io.Serializable; +import java.io.Externalizable; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.processors.cache.persistence.metastorage.MetaStorage; +import org.apache.ignite.internal.util.tostring.GridToStringInclude; +import org.apache.ignite.plugin.extensions.communication.Message; +import org.jetbrains.annotations.Nullable; /** - * Distributed metastorage data that cluster sends to joining node. + * Distributed metastorage data that cluster sends to joining node. Contains unwrapped {@link DistributedMetaStorageVersion}, + * arrays of unwrapped {@link DistributedMetaStorageKeyValuePair} (to reduce messages number) and + * {@link DistributedMetaStorageHistoryItemMessage} wraps. The original data holders are {@link Externalizable}s and + * are persistent by {@link MetaStorage with the dedicated code-generated serializers. Thus, we do not make them directly a {@link Message}. + * + * @see DmsDataWriter#write(String, byte[]) + * @see MetaStorage#write(String, Serializable) */ -@SuppressWarnings({"PublicField", "AssignmentOrReturnOfFieldWithMutableType"}) -class DistributedMetaStorageClusterNodeData implements Serializable { - /** */ - private static final long serialVersionUID = 0L; +public class DistributedMetaStorageClusterNodeData implements Message { + /** @see DistributedMetaStorageVersion#id */ + @Order(0) + @GridToStringInclude + long dVerId; - /** - * Distributed metastorage version of cluster. If {@link #fullData} is not null then this version corresponds to - * its content. - */ - public final DistributedMetaStorageVersion ver; + /** @see DistributedMetaStorageVersion#hash */ + @Order(1) + @GridToStringInclude + long dVerHash; - /** - * Full data is sent if there's not enough history items on local node. - */ - public final DistributedMetaStorageKeyValuePair[] fullData; + /** @see DistributedMetaStorageKeyValuePair#key */ + @GridToStringInclude + @Order(2) + String[] fullDataKeys; + + /** @see DistributedMetaStorageKeyValuePair#valBytes */ + @GridToStringInclude + @Order(3) + byte[][] fullDataValsBytes; /** - * Required updates for joining nodes or full available history of local node if {@link #fullData} is + * Required updates for joining nodes or full available history of local node if the full data is * not {@code null}. */ - public final DistributedMetaStorageHistoryItem[] hist; + @Order(4) + @Nullable DistributedMetaStorageHistoryItemMessage[] hist; /** - * Additional updates. Makes sence only if {@link #fullData} is not {@code null}. + * Additional updates. Makes sence only if the full data is not {@code null}. */ - public DistributedMetaStorageHistoryItem[] updates; + @Order(5) + @Nullable DistributedMetaStorageHistoryItemMessage[] updates; + + /** Empty constructor for serialization purposes. */ + public DistributedMetaStorageClusterNodeData() { + // No-op. + } /** */ public DistributedMetaStorageClusterNodeData( DistributedMetaStorageVersion ver, - DistributedMetaStorageKeyValuePair[] fullData, - DistributedMetaStorageHistoryItem[] hist, - DistributedMetaStorageHistoryItem[] updates + @Nullable DistributedMetaStorageKeyValuePair[] fullData, + @Nullable DistributedMetaStorageHistoryItem[] hist, + @Nullable DistributedMetaStorageHistoryItem[] updates ) { assert ver != null; assert fullData == null || hist != null; - this.fullData = fullData; - this.ver = ver; - this.hist = hist; - this.updates = updates; + dVerId = ver.id; + dVerHash = ver.hash; + + if (fullData != null) { + fullDataKeys = new String[fullData.length]; + fullDataValsBytes = new byte[fullData.length][]; + + for (int i = 0; i < fullData.length; ++i) { + fullDataKeys[i] = fullData[i].key; + fullDataValsBytes[i] = fullData[i].valBytes; + } + } + + this.hist = DistributedMetaStorageHistoryItemMessage.of(hist); + this.updates = DistributedMetaStorageHistoryItemMessage.of(updates); } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java index b24a275e64b53..3795cbaf0b282 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java @@ -62,6 +62,16 @@ public DistributedMetaStorageHistoryItem(String[] keys, byte[][] valBytesArr) { this.valBytesArr = valBytesArr; } + /** */ + static DistributedMetaStorageHistoryItem[] of(DistributedMetaStorageHistoryItemMessage[] histMsgs) { + DistributedMetaStorageHistoryItem[] res = new DistributedMetaStorageHistoryItem[histMsgs.length]; + + for (int i = 0; i < histMsgs.length; ++i) + res[i] = new DistributedMetaStorageHistoryItem(histMsgs[i].keys, histMsgs[i].valBytes); + + return res; + } + /** */ public long estimateSize() { int len = keys.length; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java new file mode 100644 index 0000000000000..cc8c07a839313 --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java @@ -0,0 +1,67 @@ +/* + * 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.metastorage.persistence; + +import java.io.Serializable; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.dto.IgniteDataTransferObject; +import org.apache.ignite.internal.processors.cache.persistence.metastorage.MetaStorage; +import org.apache.ignite.plugin.extensions.communication.Message; +import org.jetbrains.annotations.Nullable; + +/** + * Message for {@link DistributedMetaStorageHistoryItem} which is a persistent {@link IgniteDataTransferObject} stored + * by {@link MetaStorage} with dedicated DTO-serializer. + * + * @see MetaStorage#write(String, Serializable) + */ +public class DistributedMetaStorageHistoryItemMessage implements Message { + /** */ + @Order(0) + String[] keys; + + /** */ + @Order(1) + byte[][] valBytes; + + /** Empty constructor for serialization purposes. */ + public DistributedMetaStorageHistoryItemMessage() { + // No-op. + } + + /** */ + DistributedMetaStorageHistoryItemMessage(String[] keys, byte[][] valBytes) { + this.keys = keys; + this.valBytes = valBytes; + } + + /** @return {@link Message} wrap array for {@link DistributedMetaStorageHistoryItem} array. The history item is a persistent + * {@link IgniteDataTransferObject} stored by {@link MetaStorage} with dedicated code-generated serializer. + * Thus, we keep a {@link Message} wrap for them to transfer across the cluster. */ + static @Nullable DistributedMetaStorageHistoryItemMessage[] of(@Nullable DistributedMetaStorageHistoryItem[] hist) { + if (hist == null) + return null; + + var res = new DistributedMetaStorageHistoryItemMessage[hist.length]; + + for (int i = 0; i < hist.length; ++i) + res[i] = new DistributedMetaStorageHistoryItemMessage(hist[i].keys, hist[i].valBytesArr); + + return res; + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index ad9022f2e9747..1318ea0b046a4 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -564,22 +564,10 @@ private void onMetaStorageReadyForWrite(ReadWriteMetastorage metastorage) { lock.readLock().lock(); try { - if (isClient) { - Serializable data = new DistributedMetaStorageJoiningNodeData( - getBaselineTopologyId(), - ver, - EMPTY_ARRAY - ); - - dataBag.addJoiningNodeData(COMPONENT_ID, data); - - return; - } - - Serializable data = new DistributedMetaStorageJoiningNodeData( + var data = new DistributedMetaStorageJoiningNodeData( getBaselineTopologyId(), ver, - histCache.toArray() + isClient ? EMPTY_ARRAY : histCache.toArray() ); dataBag.addJoiningNodeData(COMPONENT_ID, data); @@ -641,9 +629,9 @@ private int getBaselineTopologyId() { if (!isPersistenceEnabled) return null; - DistributedMetaStorageVersion remoteVer = joiningData.ver; + DistributedMetaStorageVersion remoteVer = new DistributedMetaStorageVersion(joiningData.dVerId, joiningData.dVerHash); - DistributedMetaStorageHistoryItem[] remoteHist = joiningData.hist; + DistributedMetaStorageHistoryItem[] remoteHist = DistributedMetaStorageHistoryItem.of(joiningData.hist); int remoteHistSize = remoteHist.length; @@ -725,7 +713,7 @@ else if (remoteBltId < locBltId) } if (errorMsg == null) - errorMsg = validatePayload(joiningData); + errorMsg = validatePayload(remoteHist); return (errorMsg == null) ? null : new IgniteNodeValidationResult(node.id(), errorMsg); } @@ -735,11 +723,11 @@ else if (remoteBltId < locBltId) } /** - * @param joiningData Joining data to validate. + * @param remoteHist Joining data to validate. * @return {@code null} if contained data is valid otherwise error message. */ - private String validatePayload(DistributedMetaStorageJoiningNodeData joiningData) { - for (DistributedMetaStorageHistoryItem item : joiningData.hist) { + private String validatePayload(DistributedMetaStorageHistoryItem[] remoteHist) { + for (DistributedMetaStorageHistoryItem item : remoteHist) { for (int i = 0; i < item.keys().length; i++) { try { unmarshal(marshaller, item.valuesBytesArray()[i]); @@ -766,19 +754,17 @@ private String validatePayload(DistributedMetaStorageJoiningNodeData joiningData DistributedMetaStorageJoiningNodeData joiningData = discoData.joiningNodeData(); - DistributedMetaStorageVersion remoteVer = joiningData.ver; - lock.writeLock().lock(); try { DistributedMetaStorageVersion locVer = ver; - if (remoteVer.id() > locVer.id()) { - DistributedMetaStorageHistoryItem[] hist = joiningData.hist; + if (joiningData.dVerId > locVer.id()) { + DistributedMetaStorageHistoryItem[] hist = DistributedMetaStorageHistoryItem.of(joiningData.hist); - if (remoteVer.id() - locVer.id() <= hist.length) { - for (long v = locVer.id() + 1; v <= remoteVer.id(); v++) { - int hv = (int)(v - remoteVer.id() + hist.length - 1); + if (joiningData.dVerId - locVer.id() <= hist.length) { + for (long v = locVer.id() + 1; v <= joiningData.dVerId; v++) { + int hv = (int)(v - joiningData.dVerId + hist.length - 1); try { completeWrite(hist[hv]); @@ -788,8 +774,10 @@ private String validatePayload(DistributedMetaStorageJoiningNodeData joiningData } } } - else - assert false : "Joining node is too far ahead [remoteVer=" + remoteVer + "]"; + else { + assert false : "Joining node is too far ahead [remoteVerId=" + joiningData.dVerId + ", remoteVerHash=" + + joiningData.dVerHash + "]"; + } } } finally { @@ -821,23 +809,21 @@ private String validatePayload(DistributedMetaStorageJoiningNodeData joiningData DistributedMetaStorageJoiningNodeData joiningData = discoData.joiningNodeData(); - DistributedMetaStorageVersion remoteVer = joiningData.ver; - lock.readLock().lock(); try { DistributedMetaStorageVersion locVer = ver; - if (remoteVer.id() >= locVer.id()) { - Serializable nodeData = new DistributedMetaStorageClusterNodeData(remoteVer, null, null, null); + if (joiningData.dVerId >= locVer.id()) { + var rmtVer = new DistributedMetaStorageVersion(joiningData.dVerId, joiningData.dVerHash); - dataBag.addGridCommonData(COMPONENT_ID, nodeData); + dataBag.addGridCommonData(COMPONENT_ID, new DistributedMetaStorageClusterNodeData(rmtVer, null, null, null)); } else { - if (locVer.id() - remoteVer.id() <= histCache.size() && !dataBag.isJoiningNodeClient()) { - DistributedMetaStorageHistoryItem[] updates = history(remoteVer.id() + 1, locVer.id()); + if (locVer.id() - joiningData.dVerId <= histCache.size() && !dataBag.isJoiningNodeClient()) { + DistributedMetaStorageHistoryItem[] updates = history(joiningData.dVerId + 1, locVer.id()); - Serializable nodeData = new DistributedMetaStorageClusterNodeData(ver, null, null, updates); + var nodeData = new DistributedMetaStorageClusterNodeData(ver, null, null, updates); dataBag.addGridCommonData(COMPONENT_ID, nodeData); } @@ -853,7 +839,7 @@ private String validatePayload(DistributedMetaStorageJoiningNodeData joiningData else hist = history(ver.id() - histCache.size() + 1, locVer.id()); - Serializable nodeData = new DistributedMetaStorageClusterNodeData(ver0, fullData, hist, null); + var nodeData = new DistributedMetaStorageClusterNodeData(ver0, fullData, hist, null); dataBag.addGridCommonData(COMPONENT_ID, nodeData); } @@ -961,30 +947,44 @@ private DistributedMetaStorageKeyValuePair[] localFullData() { DistributedMetaStorageClusterNodeData nodeData = data.commonData(); if (nodeData != null) { - if (nodeData.fullData != null) { - ver = nodeData.ver; + DistributedMetaStorageKeyValuePair[] newfullData = null; + DistributedMetaStorageHistoryItem[] newHist = null; + + if (nodeData.fullDataKeys != null) { + assert nodeData.fullDataKeys != null ^ nodeData.fullDataValsBytes != null; + + ver = new DistributedMetaStorageVersion(nodeData.dVerId, nodeData.dVerHash); + + newfullData = new DistributedMetaStorageKeyValuePair[nodeData.fullDataKeys.length]; + + for (int i = 0; i < newfullData.length; ++i) + newfullData[i] = new DistributedMetaStorageKeyValuePair(nodeData.fullDataKeys[i], nodeData.fullDataValsBytes[i]); - notifyListenersBeforeReadyForWrite(nodeData.fullData); + notifyListenersBeforeReadyForWrite(newfullData); bridge.writeFullNodeData(nodeData); } if (nodeData.hist != null) { + newHist = new DistributedMetaStorageHistoryItem[newfullData.length]; + clearHistoryCache(); for (int i = 0, len = nodeData.hist.length; i < len; i++) { - DistributedMetaStorageHistoryItem histItem = nodeData.hist[i]; + var histItem = new DistributedMetaStorageHistoryItem(nodeData.hist[i].keys, nodeData.hist[i].valBytes); + + newHist[i] = histItem; addToHistoryCache(ver.id() + i - (len - 1), histItem); } } - if (isPersistenceEnabled && nodeData.fullData != null) - dataWriter.addUpdateTask(nodeData); + if (isPersistenceEnabled && newfullData != null) + dataWriter.addUpdateTask(ver, newHist, newfullData); if (nodeData.updates != null) { - for (DistributedMetaStorageHistoryItem update : nodeData.updates) - completeWrite(update); + for (DistributedMetaStorageHistoryItemMessage updateMsg : nodeData.updates) + completeWrite(new DistributedMetaStorageHistoryItem(updateMsg.keys, updateMsg.valBytes)); } } else if (!isClient && ver.id() > 0) { @@ -1178,9 +1178,7 @@ private RuntimeException criticalError(Throwable e) { * @param histItem {@code } pair to process. * @throws IgniteCheckedException In case of IO/unmarshalling errors. */ - private void completeWrite( - DistributedMetaStorageHistoryItem histItem - ) throws IgniteCheckedException { + private void completeWrite(DistributedMetaStorageHistoryItem histItem) throws IgniteCheckedException { assert lock.writeLock().isHeldByCurrentThread(); histItem = optimizeHistoryItem(histItem); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java index c9be2671d3105..bcb7084c38e42 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java @@ -17,30 +17,47 @@ package org.apache.ignite.internal.processors.metastorage.persistence; -import java.io.Serializable; +import java.io.Externalizable; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.processors.cache.persistence.metastorage.MetaStorage; +import org.apache.ignite.internal.util.tostring.GridToStringInclude; +import org.apache.ignite.plugin.extensions.communication.Message; /** - * Distributed metastorage data that joining node sends to cluster. + * Distributed metastorage data message that a joining node sends to a cluster. Contains unwrapped + * {@link DistributedMetaStorageVersion} (to reduce messages number) and {@link DistributedMetaStorageHistoryItemMessage}s. + * They are a {@link Externalizable} and persistent by {@link MetaStorage with the dedicated code-generated serializers. + * Thus, we do not make them directly a {@link Message}. + * + * @see DmsDataWriter#write(String, byte[]) + * @see MetaStorage#write(String, Serializable) */ -@SuppressWarnings("PublicField") -class DistributedMetaStorageJoiningNodeData implements Serializable { - /** */ - private static final long serialVersionUID = 0L; - +public class DistributedMetaStorageJoiningNodeData implements Message { /** * Baseline topology id of node, {@code -1} if baseline topology is null. */ - public final int bltId; + @Order(0) + int bltId; - /** - * Distributed metastorage version of joining node. - */ - public final DistributedMetaStorageVersion ver; + /** @see DistributedMetaStorageVersion#id */ + @Order(1) + @GridToStringInclude + long dVerId; - /** - * Available history of joining node. - */ - public final DistributedMetaStorageHistoryItem[] hist; + /** @see DistributedMetaStorageVersion#hash */ + @Order(2) + @GridToStringInclude + long dVerHash; + + /** Available history of joining node. */ + @Order(3) + @GridToStringInclude + DistributedMetaStorageHistoryItemMessage[] hist; + + /** For serialization purposes. */ + public DistributedMetaStorageJoiningNodeData() { + // No-op. + } /** */ public DistributedMetaStorageJoiningNodeData( @@ -48,8 +65,14 @@ public DistributedMetaStorageJoiningNodeData( DistributedMetaStorageVersion ver, DistributedMetaStorageHistoryItem[] hist ) { + assert ver != null; + assert hist != null; + this.bltId = bltId; - this.ver = ver; - this.hist = hist; + + dVerId = ver.id; + dVerHash = ver.hash; + + this.hist = DistributedMetaStorageHistoryItemMessage.of(hist); } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageVersion.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageVersion.java index 88b4a37084626..31102ba9f6482 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageVersion.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageVersion.java @@ -67,7 +67,7 @@ public DistributedMetaStorageVersion() { * @param id Id. * @param hash Hash. */ - private DistributedMetaStorageVersion(long id, long hash) { + DistributedMetaStorageVersion(long id, long hash) { this.id = id; this.hash = hash; } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java index 06f7337e884ae..fe223ea4517a2 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java @@ -140,29 +140,32 @@ public void addUpdateTask(DistributedMetaStorageHistoryItem histItem) { } /** */ - public void addUpdateTask(DistributedMetaStorageClusterNodeData fullNodeData) { - assert fullNodeData.fullData != null; - assert fullNodeData.hist != null; + public void addUpdateTask( + DistributedMetaStorageVersion ver, + DistributedMetaStorageHistoryItem[] hist, + DistributedMetaStorageKeyValuePair[] fullNodeData + ) { + assert ver != null; + assert hist != null; + assert fullNodeData != null; addToQueue(newDmsTask(() -> { metastorage.writeRaw(cleanupGuardKey(), DUMMY_VALUE); doCleanup(); - for (DistributedMetaStorageKeyValuePair item : fullNodeData.fullData) + for (DistributedMetaStorageKeyValuePair item : fullNodeData) metastorage.writeRaw(localKey(item.key), item.valBytes); - for (int i = 0, len = fullNodeData.hist.length; i < len; i++) { - DistributedMetaStorageHistoryItem histItem = fullNodeData.hist[i]; - - long histItemVer = fullNodeData.ver.id() + i - (len - 1); + for (int i = 0, len = hist.length; i < len; i++) { + long histItemVer = ver.id() + i - (len - 1); - metastorage.write(historyItemKey(histItemVer), histItem); + metastorage.write(historyItemKey(histItemVer), hist[i]); } - metastorage.write(versionKey(), fullNodeData.ver); + metastorage.write(versionKey(), ver); - workerDmsVer = fullNodeData.ver; + workerDmsVer = ver; metastorage.remove(cleanupGuardKey()); })); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridge.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridge.java index 298d0b5c50900..4f35148b438e6 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridge.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridge.java @@ -115,12 +115,14 @@ public DistributedMetaStorageKeyValuePair[] localFullData() { /** */ public void writeFullNodeData(DistributedMetaStorageClusterNodeData fullNodeData) { - assert fullNodeData.fullData != null; + assert fullNodeData.fullDataKeys != null; + assert fullNodeData.fullDataValsBytes != null; + assert fullNodeData.fullDataKeys.length == fullNodeData.fullDataValsBytes.length; cache.clear(); - for (DistributedMetaStorageKeyValuePair item : fullNodeData.fullData) - cache.put(item.key, item.valBytes); + for (int i = 0; i < fullNodeData.fullDataKeys.length; ++i) + cache.put(fullNodeData.fullDataKeys[i], fullNodeData.fullDataValsBytes[i]); } /** */ diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriterTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriterTest.java index 29f36e40589c0..700f46dcc750c 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriterTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriterTest.java @@ -203,12 +203,8 @@ public void testUpdateFullNodeData() throws Exception { DistributedMetaStorageVersion ver = INITIAL_VERSION.nextVersion(update); - dmsDataWriter.addUpdateTask(new DistributedMetaStorageClusterNodeData( - ver, - new DistributedMetaStorageKeyValuePair[] {toKeyValuePair(update)}, - new DistributedMetaStorageHistoryItem[] {update}, - new DistributedMetaStorageHistoryItem[] {histItem("key4", "val4")} // Has to be ignored. - )); + dmsDataWriter.addUpdateTask(ver, new DistributedMetaStorageHistoryItem[] {update}, + new DistributedMetaStorageKeyValuePair[] {toKeyValuePair(update)}); stopWorker(); From f086a96f30b558e1cf1ff5b260eb11efb5441586 Mon Sep 17 00:00:00 2001 From: Vladimir Steshin Date: Wed, 29 Jul 2026 01:00:14 +0300 Subject: [PATCH 03/25] manual review --- .../ignite/internal/CoreMessagesProvider.java | 6 ++ .../apache/ignite/internal/GridComponent.java | 4 +- ...DistributedMetaStorageClusterNodeData.java | 91 ++++++++++++----- .../DistributedMetaStorageHistoryItem.java | 17 +++- ...tributedMetaStorageHistoryItemMessage.java | 66 +++++++++++++ .../DistributedMetaStorageImpl.java | 98 +++++++++---------- ...DistributedMetaStorageJoiningNodeData.java | 61 ++++++++---- .../DistributedMetaStorageVersion.java | 2 +- .../persistence/DmsDataWriter.java | 24 ++--- ...oryCachedDistributedMetaStorageBridge.java | 8 +- .../persistence/DmsDataWriterTest.java | 8 +- 11 files changed, 264 insertions(+), 121 deletions(-) create mode 100644 modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java 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 bba4c1a6ca6b4..0d32590a48913 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 @@ -211,6 +211,9 @@ import org.apache.ignite.internal.processors.marshaller.MissingMappingResponseMessage; import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageCasAckMessage; import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageCasMessage; +import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageClusterNodeData; +import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageHistoryItemMessage; +import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageJoiningNodeData; import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageUpdateAckMessage; import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageUpdateMessage; import org.apache.ignite.internal.processors.plugin.PluginsDataBagItem; @@ -709,6 +712,9 @@ public CoreMessagesProvider(Marshaller dfltMarsh, Marshaller schemaAwareMarsh, C withNoSchema(ClusterUpdateNotifierDataBagItem.class); withNoSchemaResolvedClassLoader(PluginsDataBagItem.class); withSchema(EventsDataBagItem.class); + withSchema(DistributedMetaStorageHistoryItemMessage.class); + withSchema(DistributedMetaStorageJoiningNodeData.class); + withSchema(DistributedMetaStorageClusterNodeData.class); // [13400 - 13500]: Operation context messages. msgIdx = 13400; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/GridComponent.java b/modules/core/src/main/java/org/apache/ignite/internal/GridComponent.java index d346bbcd6152a..95913aec95f1c 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/GridComponent.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/GridComponent.java @@ -130,9 +130,7 @@ enum DiscoveryDataExchangeType { * * @param dataBag container object to store discovery data in. */ - default void collectJoiningNodeData(DiscoveryDataBag dataBag) { - // No-op. - } + public void collectJoiningNodeData(DiscoveryDataBag dataBag); /** * Collects discovery data on nodes already in grid on receiving diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java index 01375e055af15..4af7cecbcc721 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java @@ -17,51 +17,88 @@ package org.apache.ignite.internal.processors.metastorage.persistence; -import java.io.Serializable; +import java.io.Externalizable; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.processors.cache.persistence.metastorage.MetaStorage; +import org.apache.ignite.internal.util.tostring.GridToStringInclude; +import org.apache.ignite.plugin.extensions.communication.Message; +import org.jetbrains.annotations.Nullable; /** - * Distributed metastorage data that cluster sends to joining node. + * Distributed metastorage data that cluster sends to joining node. Contains unwrapped {@link DistributedMetaStorageVersion}, + * arrays of unwrapped {@link DistributedMetaStorageKeyValuePair} (to reduce messages number) and wrapped {@link DistributedMetaStorageHistoryItem}s. + * The original data holders are {@link Externalizable}s and are persistent by {@link MetaStorage} with the dedicated + * code-generated serializers. Thus, we do not make them directly a {@link Message}. + * + * @see DmsDataWriter#write(String, byte[]) + * @see MetaStorage#write(String, Serializable) */ -@SuppressWarnings({"PublicField", "AssignmentOrReturnOfFieldWithMutableType"}) -class DistributedMetaStorageClusterNodeData implements Serializable { - /** */ - private static final long serialVersionUID = 0L; +public class DistributedMetaStorageClusterNodeData implements Message { + /** @see DistributedMetaStorageVersion#id */ + @Order(0) + @GridToStringInclude + long dVerId; - /** - * Distributed metastorage version of cluster. If {@link #fullData} is not null then this version corresponds to - * its content. - */ - public final DistributedMetaStorageVersion ver; + /** @see DistributedMetaStorageVersion#hash */ + @Order(1) + @GridToStringInclude + long dVerHash; /** - * Full data is sent if there's not enough history items on local node. + * Array of the full data keys. + * + * @see DistributedMetaStorageKeyValuePair#key */ - public final DistributedMetaStorageKeyValuePair[] fullData; + @GridToStringInclude + @Order(2) + @Nullable String[] fullDataKeys; /** - * Required updates for joining nodes or full available history of local node if {@link #fullData} is - * not {@code null}. + * Arrays of the full data bytes. + * + * @see DistributedMetaStorageKeyValuePair#valBytes */ - public final DistributedMetaStorageHistoryItem[] hist; + @GridToStringInclude + @Order(3) + @Nullable byte[][] fullDataValsBytes; - /** - * Additional updates. Makes sence only if {@link #fullData} is not {@code null}. - */ - public DistributedMetaStorageHistoryItem[] updates; + /** Required updates for joining nodes or full available history of local node if the full data is not {@code null}. */ + @Order(4) + @Nullable DistributedMetaStorageHistoryItemMessage[] hist; + + /** Additional updates. Makes sence only if the full data is not {@code null}. */ + @Order(5) + @Nullable DistributedMetaStorageHistoryItemMessage[] updates; + + /** Empty constructor for serialization purposes. */ + public DistributedMetaStorageClusterNodeData() { + // No-op. + } /** */ public DistributedMetaStorageClusterNodeData( DistributedMetaStorageVersion ver, - DistributedMetaStorageKeyValuePair[] fullData, - DistributedMetaStorageHistoryItem[] hist, - DistributedMetaStorageHistoryItem[] updates + @Nullable DistributedMetaStorageKeyValuePair[] fullData, + @Nullable DistributedMetaStorageHistoryItem[] hist, + @Nullable DistributedMetaStorageHistoryItem[] updates ) { assert ver != null; assert fullData == null || hist != null; - this.fullData = fullData; - this.ver = ver; - this.hist = hist; - this.updates = updates; + dVerId = ver.id; + dVerHash = ver.hash; + + if (fullData != null) { + fullDataKeys = new String[fullData.length]; + fullDataValsBytes = new byte[fullData.length][]; + + for (int i = 0; i < fullData.length; ++i) { + fullDataKeys[i] = fullData[i].key; + fullDataValsBytes[i] = fullData[i].valBytes; + } + } + + this.hist = DistributedMetaStorageHistoryItemMessage.of(hist); + this.updates = DistributedMetaStorageHistoryItemMessage.of(updates); } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java index b24a275e64b53..ee943aeff79c1 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java @@ -21,8 +21,13 @@ import org.apache.ignite.internal.dto.IgniteDataTransferObject; 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; -/** */ +/** + * History item holder of Distributed Metastorage. + * + * @see DistributedMetaStorageHistoryItemMessage + */ final class DistributedMetaStorageHistoryItem extends IgniteDataTransferObject { /** */ private static final long serialVersionUID = 0L; @@ -62,6 +67,16 @@ public DistributedMetaStorageHistoryItem(String[] keys, byte[][] valBytesArr) { this.valBytesArr = valBytesArr; } + /** @return Array of {@link DistributedMetaStorageHistoryItem} created of the related {@link Message} transfer wraps. */ + static DistributedMetaStorageHistoryItem[] of(DistributedMetaStorageHistoryItemMessage[] histMsgs) { + DistributedMetaStorageHistoryItem[] res = new DistributedMetaStorageHistoryItem[histMsgs.length]; + + for (int i = 0; i < histMsgs.length; ++i) + res[i] = new DistributedMetaStorageHistoryItem(histMsgs[i].keys, histMsgs[i].valBytes); + + return res; + } + /** */ public long estimateSize() { int len = keys.length; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java new file mode 100644 index 0000000000000..c09466ea6ee7e --- /dev/null +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java @@ -0,0 +1,66 @@ +/* + * 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.metastorage.persistence; + +import java.io.Serializable; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.dto.IgniteDataTransferObject; +import org.apache.ignite.internal.processors.cache.persistence.metastorage.MetaStorage; +import org.apache.ignite.plugin.extensions.communication.Message; +import org.jetbrains.annotations.Nullable; + +/** + * Message wrap for {@link DistributedMetaStorageHistoryItem} which is a persistent {@link IgniteDataTransferObject} stored + * by {@link MetaStorage} using the dedicated code-generated DTO-serializer. + * + * @see DmsDataWriter#write(String, byte[]) + * @see MetaStorage#write(String, Serializable) + */ +public class DistributedMetaStorageHistoryItemMessage implements Message { + /** */ + @Order(0) + String[] keys; + + /** */ + @Order(1) + byte[][] valBytes; + + /** Empty constructor for serialization purposes. */ + public DistributedMetaStorageHistoryItemMessage() { + // No-op. + } + + /** */ + DistributedMetaStorageHistoryItemMessage(String[] keys, byte[][] valBytes) { + this.keys = keys; + this.valBytes = valBytes; + } + + /** @return {@link Message} wraps array for {@link DistributedMetaStorageHistoryItem} array. */ + static @Nullable DistributedMetaStorageHistoryItemMessage[] of(@Nullable DistributedMetaStorageHistoryItem[] hist) { + if (hist == null) + return null; + + var res = new DistributedMetaStorageHistoryItemMessage[hist.length]; + + for (int i = 0; i < hist.length; ++i) + res[i] = new DistributedMetaStorageHistoryItemMessage(hist[i].keys, hist[i].valBytesArr); + + return res; + } +} diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index ad9022f2e9747..9e410638e9e83 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -564,22 +564,10 @@ private void onMetaStorageReadyForWrite(ReadWriteMetastorage metastorage) { lock.readLock().lock(); try { - if (isClient) { - Serializable data = new DistributedMetaStorageJoiningNodeData( - getBaselineTopologyId(), - ver, - EMPTY_ARRAY - ); - - dataBag.addJoiningNodeData(COMPONENT_ID, data); - - return; - } - - Serializable data = new DistributedMetaStorageJoiningNodeData( + var data = new DistributedMetaStorageJoiningNodeData( getBaselineTopologyId(), ver, - histCache.toArray() + isClient ? EMPTY_ARRAY : histCache.toArray() ); dataBag.addJoiningNodeData(COMPONENT_ID, data); @@ -641,9 +629,9 @@ private int getBaselineTopologyId() { if (!isPersistenceEnabled) return null; - DistributedMetaStorageVersion remoteVer = joiningData.ver; + DistributedMetaStorageVersion remoteVer = new DistributedMetaStorageVersion(joiningData.dVerId, joiningData.dVerHash); - DistributedMetaStorageHistoryItem[] remoteHist = joiningData.hist; + DistributedMetaStorageHistoryItem[] remoteHist = DistributedMetaStorageHistoryItem.of(joiningData.hist); int remoteHistSize = remoteHist.length; @@ -725,7 +713,7 @@ else if (remoteBltId < locBltId) } if (errorMsg == null) - errorMsg = validatePayload(joiningData); + errorMsg = validatePayload(remoteHist); return (errorMsg == null) ? null : new IgniteNodeValidationResult(node.id(), errorMsg); } @@ -735,11 +723,11 @@ else if (remoteBltId < locBltId) } /** - * @param joiningData Joining data to validate. + * @param remoteHist Joining history data to validate. * @return {@code null} if contained data is valid otherwise error message. */ - private String validatePayload(DistributedMetaStorageJoiningNodeData joiningData) { - for (DistributedMetaStorageHistoryItem item : joiningData.hist) { + private String validatePayload(DistributedMetaStorageHistoryItem[] remoteHist) { + for (DistributedMetaStorageHistoryItem item : remoteHist) { for (int i = 0; i < item.keys().length; i++) { try { unmarshal(marshaller, item.valuesBytesArray()[i]); @@ -766,19 +754,17 @@ private String validatePayload(DistributedMetaStorageJoiningNodeData joiningData DistributedMetaStorageJoiningNodeData joiningData = discoData.joiningNodeData(); - DistributedMetaStorageVersion remoteVer = joiningData.ver; - lock.writeLock().lock(); try { DistributedMetaStorageVersion locVer = ver; - if (remoteVer.id() > locVer.id()) { - DistributedMetaStorageHistoryItem[] hist = joiningData.hist; + if (joiningData.dVerId > locVer.id()) { + DistributedMetaStorageHistoryItem[] hist = DistributedMetaStorageHistoryItem.of(joiningData.hist); - if (remoteVer.id() - locVer.id() <= hist.length) { - for (long v = locVer.id() + 1; v <= remoteVer.id(); v++) { - int hv = (int)(v - remoteVer.id() + hist.length - 1); + if (joiningData.dVerId - locVer.id() <= hist.length) { + for (long v = locVer.id() + 1; v <= joiningData.dVerId; v++) { + int hv = (int)(v - joiningData.dVerId + hist.length - 1); try { completeWrite(hist[hv]); @@ -788,8 +774,10 @@ private String validatePayload(DistributedMetaStorageJoiningNodeData joiningData } } } - else - assert false : "Joining node is too far ahead [remoteVer=" + remoteVer + "]"; + else { + assert false : "Joining node is too far ahead [remoteVerId=" + joiningData.dVerId + ", remoteVerHash=" + + joiningData.dVerHash + "]"; + } } } finally { @@ -821,23 +809,21 @@ private String validatePayload(DistributedMetaStorageJoiningNodeData joiningData DistributedMetaStorageJoiningNodeData joiningData = discoData.joiningNodeData(); - DistributedMetaStorageVersion remoteVer = joiningData.ver; - lock.readLock().lock(); try { DistributedMetaStorageVersion locVer = ver; - if (remoteVer.id() >= locVer.id()) { - Serializable nodeData = new DistributedMetaStorageClusterNodeData(remoteVer, null, null, null); + if (joiningData.dVerId >= locVer.id()) { + var rmtVer = new DistributedMetaStorageVersion(joiningData.dVerId, joiningData.dVerHash); - dataBag.addGridCommonData(COMPONENT_ID, nodeData); + dataBag.addGridCommonData(COMPONENT_ID, new DistributedMetaStorageClusterNodeData(rmtVer, null, null, null)); } else { - if (locVer.id() - remoteVer.id() <= histCache.size() && !dataBag.isJoiningNodeClient()) { - DistributedMetaStorageHistoryItem[] updates = history(remoteVer.id() + 1, locVer.id()); + if (locVer.id() - joiningData.dVerId <= histCache.size() && !dataBag.isJoiningNodeClient()) { + DistributedMetaStorageHistoryItem[] updates = history(joiningData.dVerId + 1, locVer.id()); - Serializable nodeData = new DistributedMetaStorageClusterNodeData(ver, null, null, updates); + var nodeData = new DistributedMetaStorageClusterNodeData(ver, null, null, updates); dataBag.addGridCommonData(COMPONENT_ID, nodeData); } @@ -853,7 +839,7 @@ private String validatePayload(DistributedMetaStorageJoiningNodeData joiningData else hist = history(ver.id() - histCache.size() + 1, locVer.id()); - Serializable nodeData = new DistributedMetaStorageClusterNodeData(ver0, fullData, hist, null); + var nodeData = new DistributedMetaStorageClusterNodeData(ver0, fullData, hist, null); dataBag.addGridCommonData(COMPONENT_ID, nodeData); } @@ -961,30 +947,46 @@ private DistributedMetaStorageKeyValuePair[] localFullData() { DistributedMetaStorageClusterNodeData nodeData = data.commonData(); if (nodeData != null) { - if (nodeData.fullData != null) { - ver = nodeData.ver; + // Cached unwrapped full data. + DistributedMetaStorageKeyValuePair[] newfullData = null; + // Cached unwrapped history. + DistributedMetaStorageHistoryItem[] newHist = null; + + if (nodeData.fullDataKeys != null) { + assert nodeData.fullDataValsBytes != null && nodeData.fullDataValsBytes.length == nodeData.fullDataKeys.length; + + ver = new DistributedMetaStorageVersion(nodeData.dVerId, nodeData.dVerHash); + + newfullData = new DistributedMetaStorageKeyValuePair[nodeData.fullDataKeys.length]; + + for (int i = 0; i < newfullData.length; ++i) + newfullData[i] = new DistributedMetaStorageKeyValuePair(nodeData.fullDataKeys[i], nodeData.fullDataValsBytes[i]); - notifyListenersBeforeReadyForWrite(nodeData.fullData); + notifyListenersBeforeReadyForWrite(newfullData); bridge.writeFullNodeData(nodeData); } if (nodeData.hist != null) { + newHist = new DistributedMetaStorageHistoryItem[newfullData.length]; + clearHistoryCache(); for (int i = 0, len = nodeData.hist.length; i < len; i++) { - DistributedMetaStorageHistoryItem histItem = nodeData.hist[i]; + var histItem = new DistributedMetaStorageHistoryItem(nodeData.hist[i].keys, nodeData.hist[i].valBytes); + + newHist[i] = histItem; addToHistoryCache(ver.id() + i - (len - 1), histItem); } } - if (isPersistenceEnabled && nodeData.fullData != null) - dataWriter.addUpdateTask(nodeData); + if (isPersistenceEnabled && newfullData != null) + dataWriter.addUpdateTask(ver, newHist, newfullData); if (nodeData.updates != null) { - for (DistributedMetaStorageHistoryItem update : nodeData.updates) - completeWrite(update); + for (DistributedMetaStorageHistoryItemMessage updateMsg : nodeData.updates) + completeWrite(new DistributedMetaStorageHistoryItem(updateMsg.keys, updateMsg.valBytes)); } } else if (!isClient && ver.id() > 0) { @@ -1178,9 +1180,7 @@ private RuntimeException criticalError(Throwable e) { * @param histItem {@code } pair to process. * @throws IgniteCheckedException In case of IO/unmarshalling errors. */ - private void completeWrite( - DistributedMetaStorageHistoryItem histItem - ) throws IgniteCheckedException { + private void completeWrite(DistributedMetaStorageHistoryItem histItem) throws IgniteCheckedException { assert lock.writeLock().isHeldByCurrentThread(); histItem = optimizeHistoryItem(histItem); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java index c9be2671d3105..11baf3b4934b2 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java @@ -17,30 +17,45 @@ package org.apache.ignite.internal.processors.metastorage.persistence; -import java.io.Serializable; +import java.io.Externalizable; +import org.apache.ignite.internal.Order; +import org.apache.ignite.internal.processors.cache.persistence.metastorage.MetaStorage; +import org.apache.ignite.internal.util.tostring.GridToStringInclude; +import org.apache.ignite.plugin.extensions.communication.Message; /** - * Distributed metastorage data that joining node sends to cluster. + * Distributed metastorage data message that a joining node sends to a cluster. Contains unwrapped + * {@link DistributedMetaStorageVersion} (to reduce the messages number) and {@link DistributedMetaStorageHistoryItemMessage}s. + * They are a {@link Externalizable} and persistent by {@link MetaStorage} with the dedicated code-generated serializers. + * Thus, we do not make them directly a {@link Message}. + * + * @see DmsDataWriter#write(String, byte[]) + * @see MetaStorage#write(String, Serializable) */ -@SuppressWarnings("PublicField") -class DistributedMetaStorageJoiningNodeData implements Serializable { - /** */ - private static final long serialVersionUID = 0L; +public class DistributedMetaStorageJoiningNodeData implements Message { + /** Baseline topology id of node, {@code -1} if baseline topology is null. */ + @Order(0) + int bltId; - /** - * Baseline topology id of node, {@code -1} if baseline topology is null. - */ - public final int bltId; + /** @see DistributedMetaStorageVersion#id */ + @Order(1) + @GridToStringInclude + long dVerId; - /** - * Distributed metastorage version of joining node. - */ - public final DistributedMetaStorageVersion ver; + /** @see DistributedMetaStorageVersion#hash */ + @Order(2) + @GridToStringInclude + long dVerHash; - /** - * Available history of joining node. - */ - public final DistributedMetaStorageHistoryItem[] hist; + /** Available history of joining node. */ + @Order(3) + @GridToStringInclude + DistributedMetaStorageHistoryItemMessage[] hist; + + /** For serialization purposes. */ + public DistributedMetaStorageJoiningNodeData() { + // No-op. + } /** */ public DistributedMetaStorageJoiningNodeData( @@ -48,8 +63,14 @@ public DistributedMetaStorageJoiningNodeData( DistributedMetaStorageVersion ver, DistributedMetaStorageHistoryItem[] hist ) { + assert ver != null; + assert hist != null; + this.bltId = bltId; - this.ver = ver; - this.hist = hist; + + dVerId = ver.id; + dVerHash = ver.hash; + + this.hist = DistributedMetaStorageHistoryItemMessage.of(hist); } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageVersion.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageVersion.java index 88b4a37084626..31102ba9f6482 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageVersion.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageVersion.java @@ -67,7 +67,7 @@ public DistributedMetaStorageVersion() { * @param id Id. * @param hash Hash. */ - private DistributedMetaStorageVersion(long id, long hash) { + DistributedMetaStorageVersion(long id, long hash) { this.id = id; this.hash = hash; } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java index 06f7337e884ae..b1c6bc540b360 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java @@ -140,29 +140,31 @@ public void addUpdateTask(DistributedMetaStorageHistoryItem histItem) { } /** */ - public void addUpdateTask(DistributedMetaStorageClusterNodeData fullNodeData) { - assert fullNodeData.fullData != null; - assert fullNodeData.hist != null; + public void addUpdateTask( + DistributedMetaStorageVersion ver, + DistributedMetaStorageHistoryItem[] hist, + DistributedMetaStorageKeyValuePair[] fullNodeData + ) { + assert fullNodeData != null; + assert hist != null; addToQueue(newDmsTask(() -> { metastorage.writeRaw(cleanupGuardKey(), DUMMY_VALUE); doCleanup(); - for (DistributedMetaStorageKeyValuePair item : fullNodeData.fullData) + for (DistributedMetaStorageKeyValuePair item : fullNodeData) metastorage.writeRaw(localKey(item.key), item.valBytes); - for (int i = 0, len = fullNodeData.hist.length; i < len; i++) { - DistributedMetaStorageHistoryItem histItem = fullNodeData.hist[i]; - - long histItemVer = fullNodeData.ver.id() + i - (len - 1); + for (int i = 0, len = hist.length; i < len; i++) { + long histItemVer = ver.id() + i - (len - 1); - metastorage.write(historyItemKey(histItemVer), histItem); + metastorage.write(historyItemKey(histItemVer), hist[i]); } - metastorage.write(versionKey(), fullNodeData.ver); + metastorage.write(versionKey(), ver); - workerDmsVer = fullNodeData.ver; + workerDmsVer = ver; metastorage.remove(cleanupGuardKey()); })); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridge.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridge.java index 298d0b5c50900..4f35148b438e6 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridge.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridge.java @@ -115,12 +115,14 @@ public DistributedMetaStorageKeyValuePair[] localFullData() { /** */ public void writeFullNodeData(DistributedMetaStorageClusterNodeData fullNodeData) { - assert fullNodeData.fullData != null; + assert fullNodeData.fullDataKeys != null; + assert fullNodeData.fullDataValsBytes != null; + assert fullNodeData.fullDataKeys.length == fullNodeData.fullDataValsBytes.length; cache.clear(); - for (DistributedMetaStorageKeyValuePair item : fullNodeData.fullData) - cache.put(item.key, item.valBytes); + for (int i = 0; i < fullNodeData.fullDataKeys.length; ++i) + cache.put(fullNodeData.fullDataKeys[i], fullNodeData.fullDataValsBytes[i]); } /** */ diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriterTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriterTest.java index 29f36e40589c0..700f46dcc750c 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriterTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriterTest.java @@ -203,12 +203,8 @@ public void testUpdateFullNodeData() throws Exception { DistributedMetaStorageVersion ver = INITIAL_VERSION.nextVersion(update); - dmsDataWriter.addUpdateTask(new DistributedMetaStorageClusterNodeData( - ver, - new DistributedMetaStorageKeyValuePair[] {toKeyValuePair(update)}, - new DistributedMetaStorageHistoryItem[] {update}, - new DistributedMetaStorageHistoryItem[] {histItem("key4", "val4")} // Has to be ignored. - )); + dmsDataWriter.addUpdateTask(ver, new DistributedMetaStorageHistoryItem[] {update}, + new DistributedMetaStorageKeyValuePair[] {toKeyValuePair(update)}); stopWorker(); From 15240e9ac7bb20792d39c3e9b7a2bbdd82b185c8 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Wed, 29 Jul 2026 12:36:00 +0300 Subject: [PATCH 04/25] cjeckstyle --- .../persistence/DistributedMetaStorageClusterNodeData.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java index 4af7cecbcc721..54ce954345294 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java @@ -26,9 +26,9 @@ /** * Distributed metastorage data that cluster sends to joining node. Contains unwrapped {@link DistributedMetaStorageVersion}, - * arrays of unwrapped {@link DistributedMetaStorageKeyValuePair} (to reduce messages number) and wrapped {@link DistributedMetaStorageHistoryItem}s. - * The original data holders are {@link Externalizable}s and are persistent by {@link MetaStorage} with the dedicated - * code-generated serializers. Thus, we do not make them directly a {@link Message}. + * arrays of unwrapped {@link DistributedMetaStorageKeyValuePair} (to reduce messages number) and wrapped + * {@link DistributedMetaStorageHistoryItem}s. The original data holders are {@link Externalizable}s and are persistent by + * {@link MetaStorage} with the dedicated code-generated serializers. Thus, we do not make them directly a {@link Message}. * * @see DmsDataWriter#write(String, byte[]) * @see MetaStorage#write(String, Serializable) From 75ce04732ccdfc97b1aaf5f36f4131ee2d4ec1ca Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Wed, 29 Jul 2026 12:48:22 +0300 Subject: [PATCH 05/25] + master, checkstyle, manual review --- .../DistributedMetaStorageHistoryItem.java | 12 +----------- .../DistributedMetaStorageHistoryItemMessage.java | 1 + 2 files changed, 2 insertions(+), 11 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java index 08f3bbb685baf..8ca5a84845a9f 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java @@ -24,7 +24,7 @@ import org.apache.ignite.plugin.extensions.communication.Message; /** - * History item holder of Distributed Metastorage. + * Distributed Metastorage рistory items holder. * * @see DistributedMetaStorageHistoryItemMessage */ @@ -77,16 +77,6 @@ static DistributedMetaStorageHistoryItem[] of(DistributedMetaStorageHistoryItemM return res; } - /** */ - static DistributedMetaStorageHistoryItem[] of(DistributedMetaStorageHistoryItemMessage[] histMsgs) { - DistributedMetaStorageHistoryItem[] res = new DistributedMetaStorageHistoryItem[histMsgs.length]; - - for (int i = 0; i < histMsgs.length; ++i) - res[i] = new DistributedMetaStorageHistoryItem(histMsgs[i].keys, histMsgs[i].valBytes); - - return res; - } - /** */ public long estimateSize() { int len = keys.length; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java index c09466ea6ee7e..6511de64dafd8 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java @@ -28,6 +28,7 @@ * Message wrap for {@link DistributedMetaStorageHistoryItem} which is a persistent {@link IgniteDataTransferObject} stored * by {@link MetaStorage} using the dedicated code-generated DTO-serializer. * + * @see DistributedMetaStorageHistoryItem * @see DmsDataWriter#write(String, byte[]) * @see MetaStorage#write(String, Serializable) */ From 8aad319b302bd26a3625ba4b057da1ffbcc12abe Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Wed, 29 Jul 2026 13:03:43 +0300 Subject: [PATCH 06/25] manual review, better javadocs --- .../DistributedMetaStorageClusterNodeData.java | 9 +++++---- .../persistence/DistributedMetaStorageHistoryItem.java | 6 +++++- .../DistributedMetaStorageHistoryItemMessage.java | 2 +- .../DistributedMetaStorageJoiningNodeData.java | 8 ++++---- 4 files changed, 15 insertions(+), 10 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java index 54ce954345294..0291720f12a7b 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java @@ -25,10 +25,11 @@ import org.jetbrains.annotations.Nullable; /** - * Distributed metastorage data that cluster sends to joining node. Contains unwrapped {@link DistributedMetaStorageVersion}, - * arrays of unwrapped {@link DistributedMetaStorageKeyValuePair} (to reduce messages number) and wrapped - * {@link DistributedMetaStorageHistoryItem}s. The original data holders are {@link Externalizable}s and are persistent by - * {@link MetaStorage} with the dedicated code-generated serializers. Thus, we do not make them directly a {@link Message}. + * Distributed metastorage data that cluster sends to joining node. To reduce messages number, contains plain representation + * of {@link DistributedMetaStorageVersion}, arrays of plain representations of {@link DistributedMetaStorageKeyValuePair}. + * And wrapped {@link DistributedMetaStorageHistoryItem}s. The original data holders are {@link Externalizable}s and + * are persistent by {@link MetaStorage} with the dedicated code-generated serializers. Thus, we do not make them + * directly a {@link Message}. * * @see DmsDataWriter#write(String, byte[]) * @see MetaStorage#write(String, Serializable) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java index 8ca5a84845a9f..3935775ab4185 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java @@ -19,14 +19,18 @@ import java.util.Arrays; import org.apache.ignite.internal.dto.IgniteDataTransferObject; +import org.apache.ignite.internal.processors.cache.persistence.metastorage.MetaStorage; 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; /** - * Distributed Metastorage рistory items holder. + * Distributed Metastorage history items holder. Is a persistent {@link IgniteDataTransferObject} stored by {@link MetaStorage} + * using the dedicated code-generated DTO-serializer. Then, has a transfer wrap {@link DistributedMetaStorageHistoryItemMessage}. * * @see DistributedMetaStorageHistoryItemMessage + * @see DmsDataWriter#write(String, byte[]) + * @see MetaStorage#write(String, Serializable) */ final class DistributedMetaStorageHistoryItem extends IgniteDataTransferObject { /** */ diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java index 6511de64dafd8..773464c862936 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java @@ -52,7 +52,7 @@ public DistributedMetaStorageHistoryItemMessage() { this.valBytes = valBytes; } - /** @return {@link Message} wraps array for {@link DistributedMetaStorageHistoryItem} array. */ + /** @return {@link Message} wraps array for {@code hist}. */ static @Nullable DistributedMetaStorageHistoryItemMessage[] of(@Nullable DistributedMetaStorageHistoryItem[] hist) { if (hist == null) return null; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java index 11baf3b4934b2..b93d3d098fdab 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java @@ -24,10 +24,10 @@ import org.apache.ignite.plugin.extensions.communication.Message; /** - * Distributed metastorage data message that a joining node sends to a cluster. Contains unwrapped - * {@link DistributedMetaStorageVersion} (to reduce the messages number) and {@link DistributedMetaStorageHistoryItemMessage}s. - * They are a {@link Externalizable} and persistent by {@link MetaStorage} with the dedicated code-generated serializers. - * Thus, we do not make them directly a {@link Message}. + * Distributed metastorage data message that a joining node sends to a cluster. To reduce the messages number, contains + * plain representation of {@link DistributedMetaStorageVersion}. And {@link Message}s wraps of {@link DistributedMetaStorageHistoryItem}s. + * The original data holders are a {@link Externalizable} and persistent by {@link MetaStorage} with the dedicated + * code-generated serializers. Thus, we do not make them directly a {@link Message}. * * @see DmsDataWriter#write(String, byte[]) * @see MetaStorage#write(String, Serializable) From d48f781b9cb5de29b48439ccf4c2161e1cc441a1 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Wed, 29 Jul 2026 13:07:55 +0300 Subject: [PATCH 07/25] manual review, better javadocs --- .../persistence/DistributedMetaStorageClusterNodeData.java | 2 +- .../persistence/DistributedMetaStorageHistoryItemMessage.java | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java index 0291720f12a7b..b6ae3026f6144 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java @@ -27,7 +27,7 @@ /** * Distributed metastorage data that cluster sends to joining node. To reduce messages number, contains plain representation * of {@link DistributedMetaStorageVersion}, arrays of plain representations of {@link DistributedMetaStorageKeyValuePair}. - * And wrapped {@link DistributedMetaStorageHistoryItem}s. The original data holders are {@link Externalizable}s and + * And wrapped {@link DistributedMetaStorageHistoryItem}s. The version and the full data holders are {@link Externalizable}s and * are persistent by {@link MetaStorage} with the dedicated code-generated serializers. Thus, we do not make them * directly a {@link Message}. * diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java index 773464c862936..4315821d8f74f 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java @@ -25,7 +25,7 @@ import org.jetbrains.annotations.Nullable; /** - * Message wrap for {@link DistributedMetaStorageHistoryItem} which is a persistent {@link IgniteDataTransferObject} stored + * Transfer wrap for {@link DistributedMetaStorageHistoryItem} which is a persistent {@link IgniteDataTransferObject} stored * by {@link MetaStorage} using the dedicated code-generated DTO-serializer. * * @see DistributedMetaStorageHistoryItem From 512725cca2a5e80b9fbfb3b25cac66c5057027f0 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Wed, 29 Jul 2026 13:10:43 +0300 Subject: [PATCH 08/25] manual review, better javadocs --- .../persistence/DistributedMetaStorageClusterNodeData.java | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java index b6ae3026f6144..fa17001ffb603 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java @@ -27,9 +27,8 @@ /** * Distributed metastorage data that cluster sends to joining node. To reduce messages number, contains plain representation * of {@link DistributedMetaStorageVersion}, arrays of plain representations of {@link DistributedMetaStorageKeyValuePair}. - * And wrapped {@link DistributedMetaStorageHistoryItem}s. The version and the full data holders are {@link Externalizable}s and - * are persistent by {@link MetaStorage} with the dedicated code-generated serializers. Thus, we do not make them - * directly a {@link Message}. + * And wrapped {@link DistributedMetaStorageHistoryItem}s. The version and the full data holders are {@link Externalizable}s + * persistent by {@link MetaStorage} with the dedicated code-generated serializers. Thus, we do not make them directly a {@link Message}. * * @see DmsDataWriter#write(String, byte[]) * @see MetaStorage#write(String, Serializable) From 63de1149a5ed406156e9ca84ea021ff668144a7f Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Wed, 29 Jul 2026 17:36:38 +0300 Subject: [PATCH 09/25] fix --- .../metastorage/persistence/DistributedMetaStorageImpl.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index a556b410ac3ac..7d7391b56c66d 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -968,7 +968,7 @@ private DistributedMetaStorageKeyValuePair[] localFullData() { } if (nodeData.hist != null) { - newHist = new DistributedMetaStorageHistoryItem[newfullData.length]; + newHist = new DistributedMetaStorageHistoryItem[nodeData.hist.length]; clearHistoryCache(); From 4fb509998d23d6af51e90036a036fe53cc0ecdac Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Fri, 31 Jul 2026 14:53:56 +0300 Subject: [PATCH 10/25] review fixes --- ...DistributedMetaStorageClusterNodeData.java | 4 +- .../DistributedMetaStorageHistoryItem.java | 2 +- ...tributedMetaStorageHistoryItemMessage.java | 2 +- .../DistributedMetaStorageImpl.java | 46 +++++++------------ ...DistributedMetaStorageJoiningNodeData.java | 2 +- .../persistence/DmsDataWriter.java | 8 ++-- .../persistence/DmsDataWriterTest.java | 10 +--- 7 files changed, 27 insertions(+), 47 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java index fa17001ffb603..f9119780b02da 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java @@ -98,7 +98,7 @@ public DistributedMetaStorageClusterNodeData( } } - this.hist = DistributedMetaStorageHistoryItemMessage.of(hist); - this.updates = DistributedMetaStorageHistoryItemMessage.of(updates); + this.hist = DistributedMetaStorageHistoryItemMessage.toMessages(hist); + this.updates = DistributedMetaStorageHistoryItemMessage.toMessages(updates); } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java index 3935775ab4185..761d1057da8bf 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java @@ -72,7 +72,7 @@ public DistributedMetaStorageHistoryItem(String[] keys, byte[][] valBytesArr) { } /** @return Array of {@link DistributedMetaStorageHistoryItem} created of the related {@link Message} transfer wraps. */ - static DistributedMetaStorageHistoryItem[] of(DistributedMetaStorageHistoryItemMessage[] histMsgs) { + static DistributedMetaStorageHistoryItem[] fromMessage(DistributedMetaStorageHistoryItemMessage[] histMsgs) { DistributedMetaStorageHistoryItem[] res = new DistributedMetaStorageHistoryItem[histMsgs.length]; for (int i = 0; i < histMsgs.length; ++i) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java index 4315821d8f74f..838dc2505463c 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItemMessage.java @@ -53,7 +53,7 @@ public DistributedMetaStorageHistoryItemMessage() { } /** @return {@link Message} wraps array for {@code hist}. */ - static @Nullable DistributedMetaStorageHistoryItemMessage[] of(@Nullable DistributedMetaStorageHistoryItem[] hist) { + static @Nullable DistributedMetaStorageHistoryItemMessage[] toMessages(@Nullable DistributedMetaStorageHistoryItem[] hist) { if (hist == null) return null; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index 7d7391b56c66d..61fed697c887f 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -564,13 +564,8 @@ private void onMetaStorageReadyForWrite(ReadWriteMetastorage metastorage) { lock.readLock().lock(); try { - var data = new DistributedMetaStorageJoiningNodeData( - getBaselineTopologyId(), - ver, - isClient ? EMPTY_ARRAY : histCache.toArray() - ); - - dataBag.addJoiningNodeData(COMPONENT_ID, data); + dataBag.addJoiningNodeData(COMPONENT_ID, new DistributedMetaStorageJoiningNodeData( + getBaselineTopologyId(), ver, isClient ? EMPTY_ARRAY : histCache.toArray())); } finally { lock.readLock().unlock(); @@ -631,7 +626,7 @@ private int getBaselineTopologyId() { DistributedMetaStorageVersion remoteVer = new DistributedMetaStorageVersion(joiningData.dVerId, joiningData.dVerHash); - DistributedMetaStorageHistoryItem[] remoteHist = DistributedMetaStorageHistoryItem.of(joiningData.hist); + DistributedMetaStorageHistoryItem[] remoteHist = DistributedMetaStorageHistoryItem.fromMessage(joiningData.hist); int remoteHistSize = remoteHist.length; @@ -760,7 +755,7 @@ private String validatePayload(DistributedMetaStorageHistoryItem[] remoteHist) { DistributedMetaStorageVersion locVer = ver; if (joiningData.dVerId > locVer.id()) { - DistributedMetaStorageHistoryItem[] hist = DistributedMetaStorageHistoryItem.of(joiningData.hist); + DistributedMetaStorageHistoryItem[] hist = DistributedMetaStorageHistoryItem.fromMessage(joiningData.hist); if (joiningData.dVerId - locVer.id() <= hist.length) { for (long v = locVer.id() + 1; v <= joiningData.dVerId; v++) { @@ -947,8 +942,6 @@ private DistributedMetaStorageKeyValuePair[] localFullData() { DistributedMetaStorageClusterNodeData nodeData = data.commonData(); if (nodeData != null) { - // Cached unwrapped full data. - DistributedMetaStorageKeyValuePair[] newfullData = null; // Cached unwrapped history. DistributedMetaStorageHistoryItem[] newHist = null; @@ -957,12 +950,7 @@ private DistributedMetaStorageKeyValuePair[] localFullData() { ver = new DistributedMetaStorageVersion(nodeData.dVerId, nodeData.dVerHash); - newfullData = new DistributedMetaStorageKeyValuePair[nodeData.fullDataKeys.length]; - - for (int i = 0; i < newfullData.length; ++i) - newfullData[i] = new DistributedMetaStorageKeyValuePair(nodeData.fullDataKeys[i], nodeData.fullDataValsBytes[i]); - - notifyListenersBeforeReadyForWrite(newfullData); + notifyListenersBeforeReadyForWrite(nodeData.fullDataKeys, nodeData.fullDataValsBytes); bridge.writeFullNodeData(nodeData); } @@ -981,8 +969,8 @@ private DistributedMetaStorageKeyValuePair[] localFullData() { } } - if (isPersistenceEnabled && newfullData != null) - dataWriter.addUpdateTask(ver, newHist, newfullData); + if (isPersistenceEnabled && nodeData.fullDataKeys != null) + dataWriter.addUpdateTask(ver, newHist, nodeData.fullDataKeys, nodeData.fullDataValsBytes); if (nodeData.updates != null) { for (DistributedMetaStorageHistoryItemMessage updateMsg : nodeData.updates) @@ -1318,23 +1306,22 @@ void clearHistoryCache() { /** * Notify listeners on node start. Even if there was no data restoring. * - * @param newData Data about which listeners should be notified. + * @param newDataKeys Data keys about which listeners should be notified. + * @param newDataVals Data values about which listeners should be notified. */ - private void notifyListenersBeforeReadyForWrite( - DistributedMetaStorageKeyValuePair[] newData - ) throws IgniteCheckedException { + private void notifyListenersBeforeReadyForWrite(String[] newDataKeys, byte[][] newDataVals) throws IgniteCheckedException { assert lock.isWriteLockedByCurrentThread(); DistributedMetaStorageKeyValuePair[] oldData = bridge.localFullData(); int oldIdx = 0, newIdx = 0; - while (oldIdx < oldData.length && newIdx < newData.length) { + while (oldIdx < oldData.length && newIdx < newDataKeys.length) { String oldKey = oldData[oldIdx].key; byte[] oldValBytes = oldData[oldIdx].valBytes; - String newKey = newData[newIdx].key; - byte[] newValBytes = newData[newIdx].valBytes; + String newKey = newDataKeys[newIdx]; + byte[] newValBytes = newDataVals[newIdx]; int c = oldKey.compareTo(newKey); @@ -1365,9 +1352,10 @@ else if (c > 0) { notifyListeners(oldData[oldIdx].key, () -> unmarshal(marshaller, oldValBytes), () -> null); } - for (; newIdx < newData.length; ++newIdx) { - byte[] newValBytes = newData[newIdx].valBytes; - notifyListeners(newData[newIdx].key, () -> null, () -> unmarshal(marshaller, newValBytes)); + for (; newIdx < newDataKeys.length; ++newIdx) { + byte[] newValBytes = newDataVals[newIdx]; + + notifyListeners(newDataKeys[newIdx], () -> null, () -> unmarshal(marshaller, newValBytes)); } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java index b93d3d098fdab..11748fb9c99ae 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageJoiningNodeData.java @@ -71,6 +71,6 @@ public DistributedMetaStorageJoiningNodeData( dVerId = ver.id; dVerHash = ver.hash; - this.hist = DistributedMetaStorageHistoryItemMessage.of(hist); + this.hist = DistributedMetaStorageHistoryItemMessage.toMessages(hist); } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java index b1c6bc540b360..1d8afc49ff10d 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java @@ -143,9 +143,9 @@ public void addUpdateTask(DistributedMetaStorageHistoryItem histItem) { public void addUpdateTask( DistributedMetaStorageVersion ver, DistributedMetaStorageHistoryItem[] hist, - DistributedMetaStorageKeyValuePair[] fullNodeData + String[] newDataKeys, + byte[][] newDataVals ) { - assert fullNodeData != null; assert hist != null; addToQueue(newDmsTask(() -> { @@ -153,8 +153,8 @@ public void addUpdateTask( doCleanup(); - for (DistributedMetaStorageKeyValuePair item : fullNodeData) - metastorage.writeRaw(localKey(item.key), item.valBytes); + for (int i = 0; i < newDataKeys.length; ++i) + metastorage.writeRaw(localKey(newDataKeys[i]), newDataVals[i]); for (int i = 0, len = hist.length; i < len; i++) { long histItemVer = ver.id() + i - (len - 1); diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriterTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriterTest.java index 700f46dcc750c..ed48a4e8c2348 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriterTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriterTest.java @@ -203,8 +203,7 @@ public void testUpdateFullNodeData() throws Exception { DistributedMetaStorageVersion ver = INITIAL_VERSION.nextVersion(update); - dmsDataWriter.addUpdateTask(ver, new DistributedMetaStorageHistoryItem[] {update}, - new DistributedMetaStorageKeyValuePair[] {toKeyValuePair(update)}); + dmsDataWriter.addUpdateTask(ver, new DistributedMetaStorageHistoryItem[] {update}, update.keys, update.valBytesArr); stopWorker(); @@ -325,13 +324,6 @@ public void testHalt() throws Exception { assertEquals("val1", metastorage.read(localKey("key1"))); } - /** */ - private DistributedMetaStorageKeyValuePair toKeyValuePair(DistributedMetaStorageHistoryItem histItem) { - assertEquals(1, histItem.keys().length); - - return new DistributedMetaStorageKeyValuePair(histItem.keys()[0], histItem.valuesBytesArray()[0]); - } - /** */ private void write(String key, String val) throws IgniteCheckedException { dmsDataWriter.addUpdateTask(histItem(key, val)); From eebe0d611a7c301f375242cf222a5b598f650771 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Fri, 31 Jul 2026 15:01:45 +0300 Subject: [PATCH 11/25] fix --- .../org/apache/ignite/internal/CoreMessagesProvider.java | 6 +++--- 1 file changed, 3 insertions(+), 3 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 d9531b01596c3..b025fe2cbca81 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 @@ -714,9 +714,9 @@ public CoreMessagesProvider(Marshaller dfltMarsh, Marshaller schemaAwareMarsh) { withNoSchema(ClusterUpdateNotifierDataBagItem.class); withNoSchema(PluginsDataBagItem.class); withSchema(EventsDataBagItem.class); - withSchema(DistributedMetaStorageHistoryItemMessage.class); - withSchema(DistributedMetaStorageJoiningNodeData.class); - withSchema(DistributedMetaStorageClusterNodeData.class); + withNoSchema(DistributedMetaStorageHistoryItemMessage.class); + withNoSchema(DistributedMetaStorageJoiningNodeData.class); + withNoSchema(DistributedMetaStorageClusterNodeData.class); // [13400 - 13500]: Operation context messages. msgIdx = 13400; From cf87198d70bbea0684a1985ab655c0235e9e2cbd Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Fri, 31 Jul 2026 15:53:49 +0300 Subject: [PATCH 12/25] review fixes --- .../persistence/DistributedMetaStorageHistoryItem.java | 2 +- .../persistence/DistributedMetaStorageImpl.java | 8 +++----- 2 files changed, 4 insertions(+), 6 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java index 761d1057da8bf..7f05c90eb0603 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageHistoryItem.java @@ -72,7 +72,7 @@ public DistributedMetaStorageHistoryItem(String[] keys, byte[][] valBytesArr) { } /** @return Array of {@link DistributedMetaStorageHistoryItem} created of the related {@link Message} transfer wraps. */ - static DistributedMetaStorageHistoryItem[] fromMessage(DistributedMetaStorageHistoryItemMessage[] histMsgs) { + static DistributedMetaStorageHistoryItem[] fromMessages(DistributedMetaStorageHistoryItemMessage[] histMsgs) { DistributedMetaStorageHistoryItem[] res = new DistributedMetaStorageHistoryItem[histMsgs.length]; for (int i = 0; i < histMsgs.length; ++i) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index 61fed697c887f..53e212ff66527 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -626,7 +626,7 @@ private int getBaselineTopologyId() { DistributedMetaStorageVersion remoteVer = new DistributedMetaStorageVersion(joiningData.dVerId, joiningData.dVerHash); - DistributedMetaStorageHistoryItem[] remoteHist = DistributedMetaStorageHistoryItem.fromMessage(joiningData.hist); + DistributedMetaStorageHistoryItem[] remoteHist = DistributedMetaStorageHistoryItem.fromMessages(joiningData.hist); int remoteHistSize = remoteHist.length; @@ -755,7 +755,7 @@ private String validatePayload(DistributedMetaStorageHistoryItem[] remoteHist) { DistributedMetaStorageVersion locVer = ver; if (joiningData.dVerId > locVer.id()) { - DistributedMetaStorageHistoryItem[] hist = DistributedMetaStorageHistoryItem.fromMessage(joiningData.hist); + DistributedMetaStorageHistoryItem[] hist = DistributedMetaStorageHistoryItem.fromMessages(joiningData.hist); if (joiningData.dVerId - locVer.id() <= hist.length) { for (long v = locVer.id() + 1; v <= joiningData.dVerId; v++) { @@ -818,9 +818,7 @@ private String validatePayload(DistributedMetaStorageHistoryItem[] remoteHist) { if (locVer.id() - joiningData.dVerId <= histCache.size() && !dataBag.isJoiningNodeClient()) { DistributedMetaStorageHistoryItem[] updates = history(joiningData.dVerId + 1, locVer.id()); - var nodeData = new DistributedMetaStorageClusterNodeData(ver, null, null, updates); - - dataBag.addGridCommonData(COMPONENT_ID, nodeData); + dataBag.addGridCommonData(COMPONENT_ID, new DistributedMetaStorageClusterNodeData(ver, null, null, updates)); } else { DistributedMetaStorageVersion ver0 = ver; From f1523b571f1ca9161412f2acaa60e43017e02bcf Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Fri, 31 Jul 2026 16:29:07 +0300 Subject: [PATCH 13/25] review fixes --- ...DistributedMetaStorageClusterNodeData.java | 17 ++++-- .../DistributedMetaStorageImpl.java | 59 +++++++++---------- ...oryCachedDistributedMetaStorageBridge.java | 8 +-- ...achedDistributedMetaStorageBridgeTest.java | 27 +++++---- .../junits/common/GridCommonAbstractTest.java | 2 +- 5 files changed, 56 insertions(+), 57 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java index f9119780b02da..0d2846e2600f4 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java @@ -18,6 +18,7 @@ package org.apache.ignite.internal.processors.metastorage.persistence; import java.io.Externalizable; +import java.util.Map; import org.apache.ignite.internal.Order; import org.apache.ignite.internal.processors.cache.persistence.metastorage.MetaStorage; import org.apache.ignite.internal.util.tostring.GridToStringInclude; @@ -78,7 +79,7 @@ public DistributedMetaStorageClusterNodeData() { /** */ public DistributedMetaStorageClusterNodeData( DistributedMetaStorageVersion ver, - @Nullable DistributedMetaStorageKeyValuePair[] fullData, + @Nullable Map fullData, @Nullable DistributedMetaStorageHistoryItem[] hist, @Nullable DistributedMetaStorageHistoryItem[] updates ) { @@ -89,12 +90,16 @@ public DistributedMetaStorageClusterNodeData( dVerHash = ver.hash; if (fullData != null) { - fullDataKeys = new String[fullData.length]; - fullDataValsBytes = new byte[fullData.length][]; + fullDataKeys = new String[fullData.size()]; + fullDataValsBytes = new byte[fullData.size()][]; - for (int i = 0; i < fullData.length; ++i) { - fullDataKeys[i] = fullData[i].key; - fullDataValsBytes[i] = fullData[i].valBytes; + int i = 0; + + for (var e : fullData.entrySet()) { + fullDataKeys[i] = e.getKey(); + fullDataValsBytes[i] = e.getValue(); + + ++i; } } diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index 53e212ff66527..f1969225d73e6 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -23,6 +23,7 @@ import java.util.BitSet; import java.util.Collections; import java.util.List; +import java.util.Map; import java.util.Objects; import java.util.UUID; import java.util.concurrent.ConcurrentHashMap; @@ -823,7 +824,7 @@ private String validatePayload(DistributedMetaStorageHistoryItem[] remoteHist) { else { DistributedMetaStorageVersion ver0 = ver; - DistributedMetaStorageKeyValuePair[] fullData = bridge.localFullData(); + Map fullData = bridge.localFullData(); DistributedMetaStorageHistoryItem[] hist; @@ -922,7 +923,7 @@ private DistributedMetaStorageHistoryItem[] history(long startVer, long actualVe * {@link InMemoryCachedDistributedMetaStorageBridge#localFullData()} invoked on {@link #bridge}. */ @TestOnly - private DistributedMetaStorageKeyValuePair[] localFullData() { + private Map localFullData() { return bridge.localFullData(); } @@ -1310,44 +1311,38 @@ void clearHistoryCache() { private void notifyListenersBeforeReadyForWrite(String[] newDataKeys, byte[][] newDataVals) throws IgniteCheckedException { assert lock.isWriteLockedByCurrentThread(); - DistributedMetaStorageKeyValuePair[] oldData = bridge.localFullData(); + Map oldData = bridge.localFullData(); - int oldIdx = 0, newIdx = 0; + int newIdx = 0; - while (oldIdx < oldData.length && newIdx < newDataKeys.length) { - String oldKey = oldData[oldIdx].key; - byte[] oldValBytes = oldData[oldIdx].valBytes; + for (var oldE : oldData.entrySet()) { + String oldKey = oldE.getKey(); + byte[] oldValBytes = oldE.getValue(); - String newKey = newDataKeys[newIdx]; - byte[] newValBytes = newDataVals[newIdx]; - - int c = oldKey.compareTo(newKey); + if (newIdx < newDataKeys.length) { + String newKey = newDataKeys[newIdx]; + byte[] newValBytes = newDataVals[newIdx]; - if (c < 0) { - notifyListeners(oldKey, () -> unmarshal(marshaller, oldValBytes), () -> null); + int c = oldKey.compareTo(newKey); - ++oldIdx; - } - else if (c > 0) { - notifyListeners(newKey, () -> null, () -> unmarshal(marshaller, newValBytes)); - - ++newIdx; - } - else { - notifyListeners( - oldKey, - () -> unmarshal(marshaller, oldValBytes), - () -> unmarshal(marshaller, newValBytes)); + if (c < 0) + notifyListeners(oldKey, () -> unmarshal(marshaller, oldValBytes), () -> null); + else if (c > 0) { + notifyListeners(newKey, () -> null, () -> unmarshal(marshaller, newValBytes)); - ++oldIdx; + ++newIdx; + } + else { + notifyListeners( + oldKey, + () -> unmarshal(marshaller, oldValBytes), + () -> unmarshal(marshaller, newValBytes)); - ++newIdx; + ++newIdx; + } } - } - - for (; oldIdx < oldData.length; ++oldIdx) { - byte[] oldValBytes = oldData[oldIdx].valBytes; - notifyListeners(oldData[oldIdx].key, () -> unmarshal(marshaller, oldValBytes), () -> null); + else + notifyListeners(oldKey, () -> unmarshal(marshaller, oldValBytes), () -> null); } for (; newIdx < newDataKeys.length; ++newIdx) { diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridge.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridge.java index 4f35148b438e6..0147c410b49cb 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridge.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridge.java @@ -105,12 +105,10 @@ public void write(String globalKey, @Nullable byte[] valBytes) { * Returns all {@code } pairs currently stored in distributed metastorage. Values are not unmarshalled. * All keys are sorted in ascending order. * - * @return Array of all keys and values. + * @return All the keys and values. */ - public DistributedMetaStorageKeyValuePair[] localFullData() { - return cache.entrySet().stream().map( - entry -> new DistributedMetaStorageKeyValuePair(entry.getKey(), entry.getValue()) - ).toArray(DistributedMetaStorageKeyValuePair[]::new); + public Map localFullData() { + return cache; } /** */ diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridgeTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridgeTest.java index 4c5815e1dedb1..e223fae03067a 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridgeTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/persistence/InMemoryCachedDistributedMetaStorageBridgeTest.java @@ -19,6 +19,7 @@ import java.util.ArrayList; import java.util.List; +import java.util.Map; import org.apache.ignite.IgniteCheckedException; import org.apache.ignite.marshaller.jdk.JdkMarshaller; import org.junit.Before; @@ -33,12 +34,14 @@ import static org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageVersion.INITIAL_VERSION; import static org.apache.ignite.internal.processors.metastorage.persistence.DmsDataWriter.DUMMY_VALUE; import static org.apache.ignite.testframework.junits.common.GridCommonAbstractTest.TEST_JDK_MARSHALLER; +import static org.apache.ignite.testframework.junits.common.GridCommonAbstractTest.assertEqualsMaps; import static org.hamcrest.CoreMatchers.instanceOf; import static org.hamcrest.CoreMatchers.is; import static org.junit.Assert.assertArrayEquals; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNull; import static org.junit.Assert.assertThat; +import static org.junit.Assert.assertTrue; /** */ public class InMemoryCachedDistributedMetaStorageBridgeTest { @@ -101,13 +104,13 @@ public void testLocalFullData() throws Exception { bridge.write("key2", valBytes2); bridge.write("key1", valBytes1); - DistributedMetaStorageKeyValuePair[] exp = { - new DistributedMetaStorageKeyValuePair("key1", valBytes1), - new DistributedMetaStorageKeyValuePair("key2", valBytes2), - new DistributedMetaStorageKeyValuePair("key3", valBytes3) - }; + Map exp = Map.of( + "key1", valBytes1, + "key2", valBytes2, + "key3", valBytes3 + ); - assertArrayEquals(exp, bridge.localFullData()); + assertEqualsMaps(exp, bridge.localFullData()); } /** */ @@ -117,9 +120,7 @@ public void testWriteFullNodeData() throws Exception { bridge.writeFullNodeData(new DistributedMetaStorageClusterNodeData( DistributedMetaStorageVersion.INITIAL_VERSION, - new DistributedMetaStorageKeyValuePair[] { - new DistributedMetaStorageKeyValuePair("newKey", marshaller.marshal("newVal")) - }, + Map.of("newKey", marshaller.marshal("newVal")), DistributedMetaStorageHistoryItem.EMPTY_ARRAY, null )); @@ -140,7 +141,7 @@ public void testReadInitialDataAfterFailedCleanup() throws Exception { bridge.readInitialData(metastorage); - assertArrayEquals(DistributedMetaStorageKeyValuePair.EMPTY_ARRAY, bridge.localFullData()); + assertTrue(bridge.localFullData().isEmpty()); } /** */ @@ -156,7 +157,7 @@ public void testReadInitialData1() throws Exception { bridge.readInitialData(metastorage); - assertEquals(1, bridge.localFullData().length); + assertEquals(1, bridge.localFullData().size()); assertEquals("val1", bridge.read("key1")); } @@ -175,7 +176,7 @@ public void testReadInitialData2() throws Exception { bridge.readInitialData(metastorage); - assertEquals(1, bridge.localFullData().length); + assertEquals(1, bridge.localFullData().size()); assertEquals("val1", bridge.read("key1")); } @@ -195,7 +196,7 @@ public void testReadInitialData3() throws Exception { bridge.readInitialData(metastorage); - assertEquals(1, bridge.localFullData().length); + assertEquals(1, bridge.localFullData().size()); assertEquals("val1", bridge.read("key1")); } diff --git a/modules/core/src/test/java/org/apache/ignite/testframework/junits/common/GridCommonAbstractTest.java b/modules/core/src/test/java/org/apache/ignite/testframework/junits/common/GridCommonAbstractTest.java index 6432a528b3918..7c949b1527c0f 100755 --- a/modules/core/src/test/java/org/apache/ignite/testframework/junits/common/GridCommonAbstractTest.java +++ b/modules/core/src/test/java/org/apache/ignite/testframework/junits/common/GridCommonAbstractTest.java @@ -1890,7 +1890,7 @@ protected static void assertEqualsCollectionsIgnoringOrder(Collection exp * @param exp Expected. * @param act Actual. */ - protected static void assertEqualsMaps(Map exp, Map act) { + public static void assertEqualsMaps(Map exp, Map act) { if (exp.size() != act.size()) fail("Maps are not equal:\nExpected:\t" + exp + "\nActual:\t" + act); From 34574e93656db3f4c87295f198afe3146aad305e Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Mon, 3 Aug 2026 15:18:28 +0300 Subject: [PATCH 14/25] fix --- .../DistributedMetaStorageImpl.java | 57 +++++++++++-------- 1 file changed, 33 insertions(+), 24 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index f1969225d73e6..0123fe45ab9bf 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -22,6 +22,7 @@ import java.util.Arrays; import java.util.BitSet; import java.util.Collections; +import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.Objects; @@ -1313,42 +1314,50 @@ private void notifyListenersBeforeReadyForWrite(String[] newDataKeys, byte[][] n Map oldData = bridge.localFullData(); + Iterator> oldDataIt = oldData.entrySet().iterator(); + Map.Entry oldDataEntry = oldDataIt.hasNext() ? oldDataIt.next() : null; int newIdx = 0; - for (var oldE : oldData.entrySet()) { - String oldKey = oldE.getKey(); - byte[] oldValBytes = oldE.getValue(); + while (oldDataEntry != null && newIdx < newDataKeys.length) { + String oldKey = oldDataEntry.getKey(); + byte[] oldValBytes = oldDataEntry.getValue(); - if (newIdx < newDataKeys.length) { - String newKey = newDataKeys[newIdx]; - byte[] newValBytes = newDataVals[newIdx]; + String newKey = newDataKeys[newIdx]; + byte[] newValBytes = newDataVals[newIdx]; - int c = oldKey.compareTo(newKey); + int c = oldKey.compareTo(newKey); - if (c < 0) - notifyListeners(oldKey, () -> unmarshal(marshaller, oldValBytes), () -> null); - else if (c > 0) { - notifyListeners(newKey, () -> null, () -> unmarshal(marshaller, newValBytes)); + if (c < 0) { + notifyListeners(oldKey, () -> unmarshal(marshaller, oldValBytes), () -> null); - ++newIdx; - } - else { - notifyListeners( - oldKey, - () -> unmarshal(marshaller, oldValBytes), - () -> unmarshal(marshaller, newValBytes)); + oldDataEntry = oldDataIt.hasNext() ? oldDataIt.next() : null; + } + else if (c > 0) { + notifyListeners(newKey, () -> null, () -> unmarshal(marshaller, newValBytes)); - ++newIdx; - } + ++newIdx; } - else - notifyListeners(oldKey, () -> unmarshal(marshaller, oldValBytes), () -> null); + else { + notifyListeners(oldKey, () -> unmarshal(marshaller, oldValBytes), () -> unmarshal(marshaller, newValBytes)); + + oldDataEntry = oldDataIt.hasNext() ? oldDataIt.next() : null; + + ++newIdx; + } + } + + while (oldDataEntry != null) { + byte[] oldDataVal = oldDataEntry.getValue(); + + notifyListeners(oldDataEntry.getKey(), () -> unmarshal(marshaller, oldDataVal), () -> null); + + oldDataEntry = oldDataIt.hasNext() ? oldDataIt.next() : null; } for (; newIdx < newDataKeys.length; ++newIdx) { - byte[] newValBytes = newDataVals[newIdx]; + byte[] newDataVal = newDataVals[newIdx]; - notifyListeners(newDataKeys[newIdx], () -> null, () -> unmarshal(marshaller, newValBytes)); + notifyListeners(newDataKeys[newIdx], () -> null, () -> unmarshal(marshaller, newDataVal)); } } From c35718122b59919785bf1b63b5dd8e35e1b98e4f Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Mon, 3 Aug 2026 15:47:20 +0300 Subject: [PATCH 15/25] fixes --- ...DistributedMetaStorageClusterNodeData.java | 14 +--- .../DistributedMetaStorageKeyValuePair.java | 70 ------------------- .../resources/META-INF/classnames.properties | 1 - 3 files changed, 3 insertions(+), 82 deletions(-) delete mode 100644 modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageKeyValuePair.java diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java index 0d2846e2600f4..baa45c8702118 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java @@ -27,7 +27,7 @@ /** * Distributed metastorage data that cluster sends to joining node. To reduce messages number, contains plain representation - * of {@link DistributedMetaStorageVersion}, arrays of plain representations of {@link DistributedMetaStorageKeyValuePair}. + * of {@link DistributedMetaStorageVersion}, arrays of plain representations of Distributed MetaStorage's key-value pairs. * And wrapped {@link DistributedMetaStorageHistoryItem}s. The version and the full data holders are {@link Externalizable}s * persistent by {@link MetaStorage} with the dedicated code-generated serializers. Thus, we do not make them directly a {@link Message}. * @@ -45,20 +45,12 @@ public class DistributedMetaStorageClusterNodeData implements Message { @GridToStringInclude long dVerHash; - /** - * Array of the full data keys. - * - * @see DistributedMetaStorageKeyValuePair#key - */ + /** Array of the full data keys. */ @GridToStringInclude @Order(2) @Nullable String[] fullDataKeys; - /** - * Arrays of the full data bytes. - * - * @see DistributedMetaStorageKeyValuePair#valBytes - */ + /** Arrays of the full data bytes. */ @GridToStringInclude @Order(3) @Nullable byte[][] fullDataValsBytes; diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageKeyValuePair.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageKeyValuePair.java deleted file mode 100644 index 9adbe9b9069d0..0000000000000 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageKeyValuePair.java +++ /dev/null @@ -1,70 +0,0 @@ -/* - * 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.metastorage.persistence; - -import java.io.Serializable; -import java.util.Arrays; -import org.apache.ignite.internal.util.tostring.GridToStringInclude; -import org.apache.ignite.internal.util.typedef.internal.S; - -/** */ -@SuppressWarnings("PublicField") -final class DistributedMetaStorageKeyValuePair implements Serializable { - /** */ - private static final long serialVersionUID = 0L; - - /** */ - public static final DistributedMetaStorageKeyValuePair[] EMPTY_ARRAY = {}; - - /** */ - @GridToStringInclude - public final String key; - - /** */ - @GridToStringInclude - public final byte[] valBytes; - - /** */ - public DistributedMetaStorageKeyValuePair(String key, byte[] valBytes) { - this.key = key; - this.valBytes = valBytes; - } - - /** {@inheritDoc} */ - @Override public boolean equals(Object o) { - if (this == o) - return true; - - if (o == null || getClass() != o.getClass()) - return false; - - DistributedMetaStorageKeyValuePair pair = (DistributedMetaStorageKeyValuePair)o; - - return key.equals(pair.key) && Arrays.equals(valBytes, pair.valBytes); - } - - /** {@inheritDoc} */ - @Override public int hashCode() { - return 31 * key.hashCode() + Arrays.hashCode(valBytes); - } - - /** {@inheritDoc} */ - @Override public String toString() { - return S.toString(DistributedMetaStorageKeyValuePair.class, this); - } -} diff --git a/modules/core/src/main/resources/META-INF/classnames.properties b/modules/core/src/main/resources/META-INF/classnames.properties index 931395d4b7e35..2581d4f805506 100644 --- a/modules/core/src/main/resources/META-INF/classnames.properties +++ b/modules/core/src/main/resources/META-INF/classnames.properties @@ -1612,7 +1612,6 @@ org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaSto org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageClusterNodeData org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageHistoryItem org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageJoiningNodeData -org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageKeyValuePair org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageUpdateAckMessage org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageUpdateMessage org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageVersion From 16ad874e36b65349e2d60173c4bd23ade3bf4631 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Mon, 3 Aug 2026 15:54:23 +0300 Subject: [PATCH 16/25] fix --- .../persistence/DistributedMetaStorageClusterNodeData.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java index baa45c8702118..5ccf284b8543c 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageClusterNodeData.java @@ -59,7 +59,7 @@ public class DistributedMetaStorageClusterNodeData implements Message { @Order(4) @Nullable DistributedMetaStorageHistoryItemMessage[] hist; - /** Additional updates. Makes sence only if the full data is not {@code null}. */ + /** Additional updates. Makes sense only if the full data is not {@code null}. */ @Order(5) @Nullable DistributedMetaStorageHistoryItemMessage[] updates; From b310027d4d953082c8c1dfce9ddd71cd7ec3066e Mon Sep 17 00:00:00 2001 From: Vladimir Steshin Date: Mon, 3 Aug 2026 16:09:27 +0300 Subject: [PATCH 17/25] Update modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java Co-authored-by: Dmitry Werner --- .../metastorage/persistence/DistributedMetaStorageImpl.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index 0123fe45ab9bf..4086675ff8e67 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -973,8 +973,8 @@ private Map localFullData() { dataWriter.addUpdateTask(ver, newHist, nodeData.fullDataKeys, nodeData.fullDataValsBytes); if (nodeData.updates != null) { - for (DistributedMetaStorageHistoryItemMessage updateMsg : nodeData.updates) - completeWrite(new DistributedMetaStorageHistoryItem(updateMsg.keys, updateMsg.valBytes)); + for (DistributedMetaStorageHistoryItem item : DistributedMetaStorageHistoryItem.fromMessages(nodeData.updates)) + completeWrite(item); } } else if (!isClient && ver.id() > 0) { From f72e938628b3b4eeab718e589d5dc0449899324d Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Mon, 3 Aug 2026 16:17:36 +0300 Subject: [PATCH 18/25] fix --- .../persistence/DistributedMetaStorageImpl.java | 11 +++-------- 1 file changed, 3 insertions(+), 8 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index 4086675ff8e67..f0d938bc5ada8 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -956,17 +956,12 @@ private Map localFullData() { } if (nodeData.hist != null) { - newHist = new DistributedMetaStorageHistoryItem[nodeData.hist.length]; + newHist = DistributedMetaStorageHistoryItem.fromMessages(nodeData.hist); clearHistoryCache(); - for (int i = 0, len = nodeData.hist.length; i < len; i++) { - var histItem = new DistributedMetaStorageHistoryItem(nodeData.hist[i].keys, nodeData.hist[i].valBytes); - - newHist[i] = histItem; - - addToHistoryCache(ver.id() + i - (len - 1), histItem); - } + for (int i = 0, len = newHist.length; i < len; i++) + addToHistoryCache(ver.id() + i - (len - 1), newHist[i]); } if (isPersistenceEnabled && nodeData.fullDataKeys != null) From daee4b4f2d917b8c3e1d128fe946438f8c54cec6 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 10:07:32 +0300 Subject: [PATCH 19/25] Review fixes --- .../metastorage/persistence/DistributedMetaStorageImpl.java | 6 +++--- .../processors/metastorage/persistence/DmsDataWriter.java | 2 -- 2 files changed, 3 insertions(+), 5 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index f0d938bc5ada8..f16da4414dcfc 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -942,9 +942,6 @@ private Map localFullData() { DistributedMetaStorageClusterNodeData nodeData = data.commonData(); if (nodeData != null) { - // Cached unwrapped history. - DistributedMetaStorageHistoryItem[] newHist = null; - if (nodeData.fullDataKeys != null) { assert nodeData.fullDataValsBytes != null && nodeData.fullDataValsBytes.length == nodeData.fullDataKeys.length; @@ -955,6 +952,9 @@ private Map localFullData() { bridge.writeFullNodeData(nodeData); } + // Cached unwrapped history. + DistributedMetaStorageHistoryItem[] newHist = null; + if (nodeData.hist != null) { newHist = DistributedMetaStorageHistoryItem.fromMessages(nodeData.hist); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java index 1d8afc49ff10d..5ed14fa54a200 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java @@ -146,8 +146,6 @@ public void addUpdateTask( String[] newDataKeys, byte[][] newDataVals ) { - assert hist != null; - addToQueue(newDmsTask(() -> { metastorage.writeRaw(cleanupGuardKey(), DUMMY_VALUE); From 491ab409c6f714c23643825e454892cc9b45c861 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 11:37:39 +0300 Subject: [PATCH 20/25] + master --- .../java/org/apache/ignite/internal/CoreMessagesProvider.java | 3 --- 1 file changed, 3 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 fda27a83a41f6..df69f352eda47 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 @@ -218,9 +218,6 @@ import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageClusterNodeData; import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageHistoryItemMessage; import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageJoiningNodeData; -import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageClusterNodeData; -import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageHistoryItemMessage; -import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageJoiningNodeData; import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageUpdateAckMessage; import org.apache.ignite.internal.processors.metastorage.persistence.DistributedMetaStorageUpdateMessage; import org.apache.ignite.internal.processors.plugin.PluginsDataBagItem; From a01d228f3d53ccbc91d98a8305b5a17c404f54de Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 11:51:58 +0300 Subject: [PATCH 21/25] fix --- .../metastorage/persistence/DistributedMetaStorageImpl.java | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index f16da4414dcfc..fb5cd0a5675fa 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -964,8 +964,11 @@ private Map localFullData() { addToHistoryCache(ver.id() + i - (len - 1), newHist[i]); } - if (isPersistenceEnabled && nodeData.fullDataKeys != null) + if (isPersistenceEnabled && nodeData.fullDataKeys != null) { + assert newHist != null; + dataWriter.addUpdateTask(ver, newHist, nodeData.fullDataKeys, nodeData.fullDataValsBytes); + } if (nodeData.updates != null) { for (DistributedMetaStorageHistoryItem item : DistributedMetaStorageHistoryItem.fromMessages(nodeData.updates)) From f82333c2b3c968b263373ae4b30d5763bdba16b2 Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 11:55:31 +0300 Subject: [PATCH 22/25] fix --- .../persistence/DistributedMetaStorageImpl.java | 5 +---- .../metastorage/persistence/DmsDataWriter.java | 10 ++++++---- 2 files changed, 7 insertions(+), 8 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index fb5cd0a5675fa..f16da4414dcfc 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -964,11 +964,8 @@ private Map localFullData() { addToHistoryCache(ver.id() + i - (len - 1), newHist[i]); } - if (isPersistenceEnabled && nodeData.fullDataKeys != null) { - assert newHist != null; - + if (isPersistenceEnabled && nodeData.fullDataKeys != null) dataWriter.addUpdateTask(ver, newHist, nodeData.fullDataKeys, nodeData.fullDataValsBytes); - } if (nodeData.updates != null) { for (DistributedMetaStorageHistoryItem item : DistributedMetaStorageHistoryItem.fromMessages(nodeData.updates)) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java index 5ed14fa54a200..035271a2982ac 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java @@ -142,7 +142,7 @@ public void addUpdateTask(DistributedMetaStorageHistoryItem histItem) { /** */ public void addUpdateTask( DistributedMetaStorageVersion ver, - DistributedMetaStorageHistoryItem[] hist, + @Nullable DistributedMetaStorageHistoryItem[] hist, String[] newDataKeys, byte[][] newDataVals ) { @@ -154,10 +154,12 @@ public void addUpdateTask( for (int i = 0; i < newDataKeys.length; ++i) metastorage.writeRaw(localKey(newDataKeys[i]), newDataVals[i]); - for (int i = 0, len = hist.length; i < len; i++) { - long histItemVer = ver.id() + i - (len - 1); + if (hist != null) { + for (int i = 0, len = hist.length; i < len; i++) { + long histItemVer = ver.id() + i - (len - 1); - metastorage.write(historyItemKey(histItemVer), hist[i]); + metastorage.write(historyItemKey(histItemVer), hist[i]); + } } metastorage.write(versionKey(), ver); From d454c434e75d13e49f383153482758493a6d571b Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 11:56:57 +0300 Subject: [PATCH 23/25] fix --- .../persistence/DistributedMetaStorageImpl.java | 2 +- .../metastorage/persistence/DmsDataWriter.java | 10 ++++------ 2 files changed, 5 insertions(+), 7 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index f16da4414dcfc..7423934bdb441 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -953,7 +953,7 @@ private Map localFullData() { } // Cached unwrapped history. - DistributedMetaStorageHistoryItem[] newHist = null; + DistributedMetaStorageHistoryItem[] newHist = EMPTY_ARRAY; if (nodeData.hist != null) { newHist = DistributedMetaStorageHistoryItem.fromMessages(nodeData.hist); diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java index 035271a2982ac..5ed14fa54a200 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DmsDataWriter.java @@ -142,7 +142,7 @@ public void addUpdateTask(DistributedMetaStorageHistoryItem histItem) { /** */ public void addUpdateTask( DistributedMetaStorageVersion ver, - @Nullable DistributedMetaStorageHistoryItem[] hist, + DistributedMetaStorageHistoryItem[] hist, String[] newDataKeys, byte[][] newDataVals ) { @@ -154,12 +154,10 @@ public void addUpdateTask( for (int i = 0; i < newDataKeys.length; ++i) metastorage.writeRaw(localKey(newDataKeys[i]), newDataVals[i]); - if (hist != null) { - for (int i = 0, len = hist.length; i < len; i++) { - long histItemVer = ver.id() + i - (len - 1); + for (int i = 0, len = hist.length; i < len; i++) { + long histItemVer = ver.id() + i - (len - 1); - metastorage.write(historyItemKey(histItemVer), hist[i]); - } + metastorage.write(historyItemKey(histItemVer), hist[i]); } metastorage.write(versionKey(), ver); From ec731be6b56bc5ad665210d323e2bd6a2d54fe5b Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 16:29:59 +0300 Subject: [PATCH 24/25] test fixes --- .../DistributedMetaStorageImpl.java | 2 +- .../DistributedMetaStorageTest.java | 20 ++++++++----------- .../junits/common/GridCommonAbstractTest.java | 7 ++++++- 3 files changed, 15 insertions(+), 14 deletions(-) diff --git a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java index 7423934bdb441..575300772a905 100644 --- a/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java +++ b/modules/core/src/main/java/org/apache/ignite/internal/processors/metastorage/persistence/DistributedMetaStorageImpl.java @@ -924,7 +924,7 @@ private DistributedMetaStorageHistoryItem[] history(long startVer, long actualVe * {@link InMemoryCachedDistributedMetaStorageBridge#localFullData()} invoked on {@link #bridge}. */ @TestOnly - private Map localFullData() { + public Map localFullData() { return bridge.localFullData(); } diff --git a/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/DistributedMetaStorageTest.java b/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/DistributedMetaStorageTest.java index 49f576ef60322..4a88aede8fc01 100644 --- a/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/DistributedMetaStorageTest.java +++ b/modules/core/src/test/java/org/apache/ignite/internal/processors/metastorage/DistributedMetaStorageTest.java @@ -17,9 +17,7 @@ package org.apache.ignite.internal.processors.metastorage; -import java.lang.reflect.Method; -import java.util.Arrays; -import java.util.Comparator; +import java.util.TreeMap; import java.util.UUID; import java.util.concurrent.Callable; import java.util.concurrent.ThreadLocalRandom; @@ -561,9 +559,9 @@ protected DistributedMetaStorage metastorage(int i) { * Assert that two nodes have the same internal state in {@link DistributedMetaStorage}. */ protected void assertDistributedMetastoragesAreEqual(IgniteEx ignite1, IgniteEx ignite2) throws Exception { - DistributedMetaStorage distributedMetastorage1 = ignite1.context().distributedMetastorage(); + DistributedMetaStorageImpl distributedMetastorage1 = (DistributedMetaStorageImpl)ignite1.context().distributedMetastorage(); - DistributedMetaStorage distributedMetastorage2 = ignite2.context().distributedMetastorage(); + DistributedMetaStorageImpl distributedMetastorage2 = (DistributedMetaStorageImpl)ignite2.context().distributedMetastorage(); Object ver1 = U.field(distributedMetastorage1, "ver"); @@ -577,17 +575,15 @@ protected void assertDistributedMetastoragesAreEqual(IgniteEx ignite1, IgniteEx assertEquals(histCache1, histCache2); - Method fullDataMtd = U.findNonPublicMethod(DistributedMetaStorageImpl.class, "localFullData"); + var fullData1 = distributedMetastorage1.localFullData(); - Object[] fullData1 = (Object[])fullDataMtd.invoke(distributedMetastorage1); + var fullData2 = distributedMetastorage2.localFullData(); - Object[] fullData2 = (Object[])fullDataMtd.invoke(distributedMetastorage2); - - assertEqualsCollections(Arrays.asList(fullData1), Arrays.asList(fullData2)); + assertEqualsMaps(fullData1, fullData2); // Also check that arrays are sorted. - Arrays.sort(fullData1, Comparator.comparing(o -> U.field(o, "key"))); + fullData1 = new TreeMap<>(fullData1); - assertEqualsCollections(Arrays.asList(fullData1), Arrays.asList(fullData2)); + assertEqualsMaps(fullData1, fullData2); } } diff --git a/modules/core/src/test/java/org/apache/ignite/testframework/junits/common/GridCommonAbstractTest.java b/modules/core/src/test/java/org/apache/ignite/testframework/junits/common/GridCommonAbstractTest.java index 7c949b1527c0f..6753d0861e6e6 100755 --- a/modules/core/src/test/java/org/apache/ignite/testframework/junits/common/GridCommonAbstractTest.java +++ b/modules/core/src/test/java/org/apache/ignite/testframework/junits/common/GridCommonAbstractTest.java @@ -1897,8 +1897,13 @@ public static void assertEqualsMaps(Map exp, Map act) { for (Map.Entry e : exp.entrySet()) { if (!act.containsKey(e.getKey())) fail("Maps are not equal (missing key " + e.getKey() + "):\nExpected:\t" + exp + "\nActual:\t" + act); - else if (!Objects.equals(e.getValue(), act.get(e.getKey()))) + + try { + assertEqualsArraysAware(e.getValue(), act.get(e.getKey())); + } + catch (AssertionError ignored) { fail("Maps are not equal (key " + e.getKey() + "):\nExpected:\t" + exp + "\nActual:\t" + act); + } } } From 7de595461455a980454162871d7d80d0cfe53f0b Mon Sep 17 00:00:00 2001 From: Steshin Vladimir Date: Tue, 4 Aug 2026 16:36:08 +0300 Subject: [PATCH 25/25] test fixes --- .../junits/common/GridCommonAbstractTest.java | 11 +++++------ 1 file changed, 5 insertions(+), 6 deletions(-) diff --git a/modules/core/src/test/java/org/apache/ignite/testframework/junits/common/GridCommonAbstractTest.java b/modules/core/src/test/java/org/apache/ignite/testframework/junits/common/GridCommonAbstractTest.java index 6753d0861e6e6..55e53c7d44cb1 100755 --- a/modules/core/src/test/java/org/apache/ignite/testframework/junits/common/GridCommonAbstractTest.java +++ b/modules/core/src/test/java/org/apache/ignite/testframework/junits/common/GridCommonAbstractTest.java @@ -1898,12 +1898,11 @@ public static void assertEqualsMaps(Map exp, Map act) { if (!act.containsKey(e.getKey())) fail("Maps are not equal (missing key " + e.getKey() + "):\nExpected:\t" + exp + "\nActual:\t" + act); - try { - assertEqualsArraysAware(e.getValue(), act.get(e.getKey())); - } - catch (AssertionError ignored) { - fail("Maps are not equal (key " + e.getKey() + "):\nExpected:\t" + exp + "\nActual:\t" + act); - } + assertEqualsArraysAware( + "Maps are not equal (key " + e.getKey() + "):\nExpected:\t" + exp + "\nActual:\t" + act, + e.getValue(), + act.get(e.getKey()) + ); } }