From 7dc9d05148aef292c84c1cf0375b78d818c02584 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Wed, 29 Jul 2026 11:34:59 +0800 Subject: [PATCH 1/2] [Fix] Keep heartbeat alive when metadata lease is fenced --- .../iotdb/db/i18n/DataNodeSchemaMessages.java | 3 + .../iotdb/db/i18n/DataNodeSchemaMessages.java | 3 + .../impl/DataNodeInternalRPCServiceImpl.java | 31 +++++++-- .../iotdb/db/schemaengine/SchemaEngine.java | 28 ++++++-- .../DataNodeInternalRPCServiceImplTest.java | 68 +++++++++++++++++++ 5 files changed, 122 insertions(+), 11 deletions(-) diff --git a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeSchemaMessages.java b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeSchemaMessages.java index 4e0e09fe0b704..dd1c4af94f647 100644 --- a/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeSchemaMessages.java +++ b/iotdb-core/datanode/src/main/i18n/en/org/apache/iotdb/db/i18n/DataNodeSchemaMessages.java @@ -36,6 +36,9 @@ public final class DataNodeSchemaMessages { "SchemaRegion(id = {}) has been deleted, skiped"; public static final String FAILED_TO_GET_TABLE_FOR_TIMESERIES_COUNT = "Failed to get table {}.{} when calculating the time series number. Maybe the cluster is restarting or the table is being dropped."; + public static final String + LOG_METADATA_LEASE_IS_FENCED_SKIP_REPORTING_SCHEMA_USAGE_IN_THIS_HEARTBEAT_81D36975 = + "Metadata lease is fenced. Skip reporting schema usage in this heartbeat."; public static final String TREE_VIEW_TABLE_CANNOT_BE_WRITTEN_OR_DELETED = "The table %s.%s is a view from tree, cannot be written or deleted from"; public static final String PEER_IS_SHUTTING_DOWN = "Peer is shutting down now."; diff --git a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeSchemaMessages.java b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeSchemaMessages.java index e60932b120a7f..411044b14afb7 100644 --- a/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeSchemaMessages.java +++ b/iotdb-core/datanode/src/main/i18n/zh/org/apache/iotdb/db/i18n/DataNodeSchemaMessages.java @@ -34,6 +34,9 @@ public final class DataNodeSchemaMessages { public static final String SCHEMA_REGION_ALREADY_DELETED = "SchemaRegion(id = {}) 已被删除,已跳过"; public static final String FAILED_TO_GET_TABLE_FOR_TIMESERIES_COUNT = "计算时间序列数量时获取表 {}.{} 失败,可能是集群正在重启或表正在被删除。"; + public static final String + LOG_METADATA_LEASE_IS_FENCED_SKIP_REPORTING_SCHEMA_USAGE_IN_THIS_HEARTBEAT_81D36975 = + "元数据租约已被隔离,本次心跳跳过上报 schema 使用量。"; public static final String TREE_VIEW_TABLE_CANNOT_BE_WRITTEN_OR_DELETED = "表 %s.%s 是树模型视图,不能写入或删除"; public static final String PEER_IS_SHUTTING_DOWN = "节点正在关闭中。"; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java index ebf9845aeaec0..9f62a18c73a0b 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/protocol/thrift/impl/DataNodeInternalRPCServiceImpl.java @@ -68,6 +68,7 @@ import org.apache.iotdb.commons.enums.DataPartitionTableGeneratorState; import org.apache.iotdb.commons.exception.IllegalPathException; import org.apache.iotdb.commons.exception.MetadataException; +import org.apache.iotdb.commons.exception.MetadataLeaseFencedException; import org.apache.iotdb.commons.i18n.PipeMessages; import org.apache.iotdb.commons.partition.DataPartitionTable; import org.apache.iotdb.commons.partition.DatabaseScopedDataPartitionTable; @@ -117,6 +118,7 @@ import org.apache.iotdb.db.consensus.SchemaRegionConsensusImpl; import org.apache.iotdb.db.exception.StorageEngineException; import org.apache.iotdb.db.i18n.DataNodeMiscMessages; +import org.apache.iotdb.db.i18n.DataNodeSchemaMessages; import org.apache.iotdb.db.partition.DataPartitionTableGenerator; import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent; import org.apache.iotdb.db.protocol.client.ConfigNodeInfo; @@ -2355,19 +2357,38 @@ public TDataNodeHeartbeatResp getDataNodeHeartBeat(TDataNodeHeartbeatReq req) th if (commonConfig.getStatusReason() != null) { resp.setStatusReason(commonConfig.getStatusReason()); } + MetadataLeaseFencedException schemaUsageCollectionException = null; if (req.getSchemaRegionIds() != null) { spaceQuotaManager.updateSpaceQuotaUsage(req.getSpaceQuotaUsage()); - resp.setRegionDeviceUsageMap( - schemaEngine.countDeviceNumBySchemaRegion(req.getSchemaRegionIds())); - resp.setRegionSeriesUsageMap( - schemaEngine.countTimeSeriesNumBySchemaRegion(req.getSchemaRegionIds())); + try { + resp.setRegionDeviceUsageMap( + schemaEngine.countDeviceNumBySchemaRegion(req.getSchemaRegionIds())); + resp.setRegionSeriesUsageMap( + schemaEngine.countTimeSeriesNumBySchemaRegion(req.getSchemaRegionIds())); + } catch (final MetadataLeaseFencedException e) { + schemaUsageCollectionException = e; + } } if (req.getDataRegionIds() != null) { spaceQuotaManager.setDataRegionIds(req.getDataRegionIds()); resp.setRegionDisk(spaceQuotaManager.getRegionDisk()); } // Update schema quota if necessary - SchemaEngine.getInstance().updateAndFillSchemaCountMap(req, resp); + try { + SchemaEngine.getInstance().updateAndFillSchemaCountMap(req, resp); + } catch (final MetadataLeaseFencedException e) { + if (schemaUsageCollectionException == null) { + schemaUsageCollectionException = e; + } + } + if (schemaUsageCollectionException != null) { + resp.unsetRegionDeviceUsageMap(); + resp.unsetRegionSeriesUsageMap(); + LOGGER.warn( + DataNodeSchemaMessages + .LOG_METADATA_LEASE_IS_FENCED_SKIP_REPORTING_SCHEMA_USAGE_IN_THIS_HEARTBEAT_81D36975, + schemaUsageCollectionException); + } // Update pipe meta if necessary if (req.isNeedPipeMetaList()) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/SchemaEngine.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/SchemaEngine.java index 295813dc42066..9b634b443a47c 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/SchemaEngine.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/schemaengine/SchemaEngine.java @@ -447,11 +447,21 @@ public void updateAndFillSchemaCountMap( schemaQuotaManager.updateRemain( req.getTimeSeriesQuotaRemain(), req.isSetDeviceQuotaRemain() ? req.getDeviceQuotaRemain() : -1); + // Build both maps on snapshots and publish them together only after all counts succeed, so a + // failure cannot leave a partially updated heartbeat response. + Map regionDeviceUsageMap = + resp.getRegionDeviceUsageMap() == null + ? null + : new HashMap<>(resp.getRegionDeviceUsageMap()); + Map regionSeriesUsageMap = + resp.getRegionSeriesUsageMap() == null + ? null + : new HashMap<>(resp.getRegionSeriesUsageMap()); if (schemaQuotaManager.isDeviceLimit()) { - if (resp.getRegionDeviceUsageMap() == null) { - resp.setRegionDeviceUsageMap(new HashMap<>()); + if (regionDeviceUsageMap == null) { + regionDeviceUsageMap = new HashMap<>(); } - final Map tmp = resp.getRegionDeviceUsageMap(); + final Map tmp = regionDeviceUsageMap; SchemaRegionConsensusImpl.getInstance().getAllConsensusGroupIds().stream() .filter( consensusGroupId -> @@ -467,10 +477,10 @@ public void updateAndFillSchemaCountMap( .orElse(0L))); } if (schemaQuotaManager.isMeasurementLimit()) { - if (resp.getRegionSeriesUsageMap() == null) { - resp.setRegionSeriesUsageMap(new HashMap<>()); + if (regionSeriesUsageMap == null) { + regionSeriesUsageMap = new HashMap<>(); } - final Map tmp = resp.getRegionSeriesUsageMap(); + final Map tmp = regionSeriesUsageMap; SchemaRegionConsensusImpl.getInstance().getAllConsensusGroupIds().stream() .filter( consensusGroupId -> @@ -490,6 +500,12 @@ public void updateAndFillSchemaCountMap( .map(this::getTimeSeriesNumber4Quota) .orElse(0L))); } + if (regionDeviceUsageMap != null) { + resp.setRegionDeviceUsageMap(regionDeviceUsageMap); + } + if (regionSeriesUsageMap != null) { + resp.setRegionSeriesUsageMap(regionSeriesUsageMap); + } } public ISchemaEngineStatistics getSchemaEngineStatistics() { diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/DataNodeInternalRPCServiceImplTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/DataNodeInternalRPCServiceImplTest.java index ca68525243d12..f7dcb64295ad2 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/DataNodeInternalRPCServiceImplTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/DataNodeInternalRPCServiceImplTest.java @@ -28,6 +28,7 @@ import org.apache.iotdb.commons.consensus.DataRegionId; import org.apache.iotdb.commons.consensus.SchemaRegionId; import org.apache.iotdb.commons.exception.MetadataException; +import org.apache.iotdb.commons.exception.MetadataLeaseFencedException; import org.apache.iotdb.commons.path.MeasurementPath; import org.apache.iotdb.commons.path.PartialPath; import org.apache.iotdb.commons.queryengine.plan.planner.plan.node.PlanNodeId; @@ -50,10 +51,14 @@ import org.apache.iotdb.db.queryengine.plan.planner.plan.node.metadata.write.CreateMultiTimeSeriesNode; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.metadata.write.CreateTimeSeriesNode; import org.apache.iotdb.db.schemaengine.SchemaEngine; +import org.apache.iotdb.db.schemaengine.lease.MetadataLeaseManager; +import org.apache.iotdb.db.schemaengine.rescon.DataNodeSchemaQuotaManager; import org.apache.iotdb.db.service.DataNode.DataNodeContext; import org.apache.iotdb.db.storageengine.dataregion.DataRegion; import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource; import org.apache.iotdb.db.utils.EnvironmentUtils; +import org.apache.iotdb.mpp.rpc.thrift.TDataNodeHeartbeatReq; +import org.apache.iotdb.mpp.rpc.thrift.TDataNodeHeartbeatResp; import org.apache.iotdb.mpp.rpc.thrift.TPlanNode; import org.apache.iotdb.mpp.rpc.thrift.TSendBatchPlanNodeReq; import org.apache.iotdb.mpp.rpc.thrift.TSendBatchPlanNodeResp; @@ -73,6 +78,8 @@ import java.io.File; import java.io.IOException; +import java.lang.reflect.Field; +import java.lang.reflect.Method; import java.nio.ByteBuffer; import java.util.ArrayList; import java.util.Collections; @@ -81,6 +88,7 @@ import java.util.Map; import java.util.Objects; import java.util.Optional; +import java.util.concurrent.atomic.AtomicBoolean; import static org.mockito.Mockito.when; @@ -229,6 +237,66 @@ public void testRejectLoad4NonActiveImpl() throws LoadFileException { .setActive(true); } + @Test + public void testHeartbeatSkipsSchemaUsageWhenMetadataLeaseIsFenced() throws Exception { + final MetadataLeaseManager leaseManager = MetadataLeaseManager.getInstance(); + final Field hasPullTaskNowRefField = + MetadataLeaseManager.class.getDeclaredField("hasPullTaskNowRef"); + hasPullTaskNowRefField.setAccessible(true); + final AtomicBoolean hasPullTaskNowRef = + (AtomicBoolean) hasPullTaskNowRefField.get(leaseManager); + final Map tableDeviceNumberMap = + SchemaEngine.getInstance() + .getSchemaRegion(new SchemaRegionId(0)) + .getSchemaRegionStatistics() + .getTable2DevicesNumMap(); + + try { + hasPullTaskNowRef.set(true); + tableDeviceNumberMap.put("table", 1L); + leaseManager.recoveryLeaseForTest(false); + final Method checkLeaseStatus = + MetadataLeaseManager.class.getDeclaredMethod("checkLeaseStatus"); + checkLeaseStatus.setAccessible(true); + checkLeaseStatus.invoke(leaseManager); + Assert.assertTrue(leaseManager.isFenced()); + + final TDataNodeHeartbeatReq req = + new TDataNodeHeartbeatReq() + .setHeartbeatTimestamp(1L) + .setNeedJudgeLeader(false) + .setNeedSamplingLoad(false) + .setTimeSeriesQuotaRemain(10L) + .setDeviceQuotaRemain(10L) + .setLogicalClock(0L); + final Map previousDeviceCountMap = new HashMap<>(); + previousDeviceCountMap.put(-1, 1L); + final Map previousSeriesCountMap = new HashMap<>(); + previousSeriesCountMap.put(-1, 2L); + final TDataNodeHeartbeatResp schemaCountResp = + new TDataNodeHeartbeatResp() + .setRegionDeviceUsageMap(new HashMap<>(previousDeviceCountMap)) + .setRegionSeriesUsageMap(new HashMap<>(previousSeriesCountMap)); + + Assert.assertThrows( + MetadataLeaseFencedException.class, + () -> SchemaEngine.getInstance().updateAndFillSchemaCountMap(req, schemaCountResp)); + Assert.assertEquals(previousDeviceCountMap, schemaCountResp.getRegionDeviceUsageMap()); + Assert.assertEquals(previousSeriesCountMap, schemaCountResp.getRegionSeriesUsageMap()); + + final TDataNodeHeartbeatResp resp = dataNodeInternalRPCServiceImpl.getDataNodeHeartBeat(req); + + Assert.assertEquals(req.getHeartbeatTimestamp(), resp.getHeartbeatTimestamp()); + Assert.assertFalse(resp.isSetRegionDeviceUsageMap()); + Assert.assertFalse(resp.isSetRegionSeriesUsageMap()); + } finally { + tableDeviceNumberMap.remove("table"); + hasPullTaskNowRef.set(false); + leaseManager.recoveryLeaseForTest(true); + DataNodeSchemaQuotaManager.getInstance().updateRemain(-1, -1); + } + } + @Test public void testCreateAlignedTimeSeries() throws MetadataException { CreateAlignedTimeSeriesNode createAlignedTimeSeriesNode = From 48b537436d72351357c590390b2b8471a406e556 Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Thu, 30 Jul 2026 16:56:11 +0800 Subject: [PATCH 2/2] Stabilize metadata lease heartbeat test --- .../iotdb/db/service/DataNodeInternalRPCServiceImplTest.java | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/DataNodeInternalRPCServiceImplTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/DataNodeInternalRPCServiceImplTest.java index f7dcb64295ad2..c6059d39d049e 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/DataNodeInternalRPCServiceImplTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/service/DataNodeInternalRPCServiceImplTest.java @@ -68,6 +68,7 @@ import org.apache.tsfile.enums.TSDataType; import org.apache.tsfile.file.metadata.enums.CompressionType; import org.apache.tsfile.file.metadata.enums.TSEncoding; +import org.awaitility.Awaitility; import org.junit.After; import org.junit.AfterClass; import org.junit.Assert; @@ -88,6 +89,7 @@ import java.util.Map; import java.util.Objects; import java.util.Optional; +import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import static org.mockito.Mockito.when; @@ -252,6 +254,9 @@ public void testHeartbeatSkipsSchemaUsageWhenMetadataLeaseIsFenced() throws Exce .getTable2DevicesNumMap(); try { + Awaitility.await() + .atMost(10, TimeUnit.SECONDS) + .until(() -> SchemaRegionConsensusImpl.getInstance().isLeader(new SchemaRegionId(0))); hasPullTaskNowRef.set(true); tableDeviceNumberMap.put("table", 1L); leaseManager.recoveryLeaseForTest(false);