summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorNikolay Shestakov <[email protected]>2026-07-02 12:47:55 +0500
committerGitHub <[email protected]>2026-07-02 07:47:55 +0000
commit72b60f79942f3cdd2d18ac2d080c23d4878b5902 (patch)
tree2fb67515de99c54a8b20026ed03d1da84fa4c343
parent4142c3987052c2013e6442f4d0f09e8b5fe3ff82 (diff)
Fixed yexception: Invalid partition count in PQ alter operation at schemeshard_info_types.h:1762 (#45018)
-rw-r--r--ydb/core/tx/schemeshard/schemeshard__init.cpp11
-rw-r--r--ydb/core/tx/schemeshard/schemeshard__operation_alter_pq.cpp94
-rw-r--r--ydb/core/tx/schemeshard/schemeshard_info_types.h8
-rw-r--r--ydb/core/tx/schemeshard/ut_topic_splitmerge/ut_topic_splitmerge.cpp285
4 files changed, 372 insertions, 26 deletions
diff --git a/ydb/core/tx/schemeshard/schemeshard__init.cpp b/ydb/core/tx/schemeshard/schemeshard__init.cpp
index 4a3a424ac37..d93f3ec18bf 100644
--- a/ydb/core/tx/schemeshard/schemeshard__init.cpp
+++ b/ydb/core/tx/schemeshard/schemeshard__init.cpp
@@ -2992,7 +2992,16 @@ struct TSchemeShard::TTxInit : public TTransactionBase<TSchemeShard> {
auto it = Self->Topics.find(pathId);
Y_ABORT_UNLESS(it != Self->Topics.end());
- alterData->TotalPartitionCount = it->second->GetTotalPartitionCountWithAlter();
+ alterData->TotalPartitionCount = 0;
+ alterData->ActivePartitionCount = 0;
+ for (const auto& [_, partition] : it->second->Partitions) {
+ if (partition->AlterVersion <= alterData->AlterVersion) {
+ ++alterData->TotalPartitionCount;
+ if (partition->Status == NKikimrPQ::ETopicPartitionStatus::Active) {
+ ++alterData->ActivePartitionCount;
+ }
+ }
+ }
alterData->BalancerTabletID = it->second->BalancerTabletID;
alterData->BalancerShardIdx = it->second->BalancerShardIdx;
it->second->AlterData = alterData;
diff --git a/ydb/core/tx/schemeshard/schemeshard__operation_alter_pq.cpp b/ydb/core/tx/schemeshard/schemeshard__operation_alter_pq.cpp
index 41b11d55a68..5484219fafb 100644
--- a/ydb/core/tx/schemeshard/schemeshard__operation_alter_pq.cpp
+++ b/ydb/core/tx/schemeshard/schemeshard__operation_alter_pq.cpp
@@ -3,6 +3,7 @@
#include "schemeshard_impl.h"
#include "schemeshard_pq_helpers.h" // for PQGroupReserve
+#include <library/cpp/containers/absl_flat_hash/flat_hash_set.h>
#include <library/cpp/containers/top_keeper/top_keeper.h>
#include <ydb/core/base/subdomain.h>
@@ -47,6 +48,63 @@ std::expected<void, std::string> ValidateKeyRangeSequence(const auto& partitions
return NKikimr::NPQ::ValidateKeyRangeSequence(bounds);
}
+size_t CountTopicTotalPartitions(const TTopicInfo::TPtr& topic) {
+ return topic->Partitions.size();
+}
+
+size_t CountTopicActivePartitions(const TTopicInfo::TPtr& topic) {
+ size_t count = 0;
+ for (const auto& [_, partition] : topic->Partitions) {
+ if (partition->Status == NKikimrPQ::ETopicPartitionStatus::Active) {
+ ++count;
+ }
+ }
+ return count;
+}
+
+size_t ComputeAlterActivePartitionCount(
+ const TTopicInfo::TPtr& topic,
+ const TTopicInfo::TPtr& alterData)
+{
+ size_t count = CountTopicActivePartitions(topic);
+ const auto& partitionsToAdd = alterData->PartitionsToAdd;
+ const size_t partitionsToAddCount = partitionsToAdd.size();
+
+ size_t parentCount = 0;
+ absl::flat_hash_set<ui32> addedPartitionIds;
+ addedPartitionIds.reserve(partitionsToAddCount * 2);
+ for (const auto& partition : partitionsToAdd) {
+ parentCount += partition.ParentPartitionIds.size();
+ addedPartitionIds.insert(partition.PartitionId);
+ }
+
+ absl::flat_hash_set<ui32> deactivatedParents;
+ deactivatedParents.reserve(parentCount * 2);
+ for (const auto& partition : partitionsToAdd) {
+ ++count;
+ for (const ui32 parentId : partition.ParentPartitionIds) {
+ if (deactivatedParents.emplace(parentId).second) {
+ const auto parentIt = topic->Partitions.find(parentId);
+ if (parentIt != topic->Partitions.end()
+ && parentIt->second->Status == NKikimrPQ::ETopicPartitionStatus::Active) {
+ --count;
+ } else if (addedPartitionIds.contains(parentId)) {
+ --count;
+ }
+ }
+ }
+ }
+ return count;
+}
+
+void ComputeAlterPartitionCounts(
+ const TTopicInfo::TPtr& topic,
+ TTopicInfo::TPtr alterData)
+{
+ alterData->TotalPartitionCount = CountTopicTotalPartitions(topic) + alterData->PartitionsToAdd.size();
+ alterData->ActivePartitionCount = ComputeAlterActivePartitionCount(topic, alterData);
+}
+
class TAlterPQ: public TSubOperation {
// Make sure we make decisions using a consistent runtime value
static TTxState::ETxState NextState() {
@@ -505,7 +563,9 @@ public:
<< ", first new shardIdx " << startShardIdx
<< " hasBalancer " << hasBalancer);
- ReassignIds(pqGroup);
+ if (!pqGroup->AlterData->PartitionsToAdd.empty()) {
+ ReassignIds(pqGroup);
+ }
return shardsToCreate > 0;
}
@@ -638,8 +698,6 @@ public:
return result;
}
- alterData->ActivePartitionCount = topic->ActivePartitionCount;
-
bool splitMergeEnabled = NKikimr::NPQ::SplitMergeEnabled(tabletConfig)
&& NKikimr::NPQ::SplitMergeEnabled(newTabletConfig);
@@ -735,14 +793,12 @@ public:
}
alterData->PartitionsToAdd.emplace_back(partitionIdIndex.value(), partitionIdIndex.value() + 1, range);
alterData->TotalGroupCount += 1;
- ++alterData->ActivePartitionCount;
}
}
}
for (const auto& split : alter.GetSplit()) {
alterData->TotalGroupCount += 2;
- ++alterData->ActivePartitionCount;
const auto splittedPartitionId = split.GetPartition();
if (!topic->Partitions.contains(splittedPartitionId)) {
@@ -837,7 +893,6 @@ public:
}
for (const auto& merge : alter.GetMerge()) {
alterData->TotalGroupCount += 1;
- --alterData->ActivePartitionCount;
const auto partitionId = merge.GetPartition();
if (!topic->Partitions.contains(partitionId)) {
@@ -916,10 +971,11 @@ public:
}
if (alter.HasPQTabletConfig() && alter.GetPQTabletConfig().HasPartitionStrategy()) {
+ size_t activePartitionCount = CountTopicActivePartitions(topic);
auto requestedMinPartitionCount = alter.GetPQTabletConfig().GetPartitionStrategy().GetMinPartitionCount();
- if (requestedMinPartitionCount > alterData->ActivePartitionCount) {
+ if (requestedMinPartitionCount > activePartitionCount) {
// select exisisting active partitions for split
- auto numPartitionsToSplit = requestedMinPartitionCount - alterData->ActivePartitionCount;
+ auto numPartitionsToSplit = requestedMinPartitionCount - activePartitionCount;
struct TPartitionsComparer {
bool operator()(const TTopicTabletInfo::TTopicPartitionInfo* lhs, const TTopicTabletInfo::TTopicPartitionInfo* rhs) const {
return lhs->ParentPartitionIds.size() < rhs->ParentPartitionIds.size(); // for now simply sort by number of parents
@@ -963,7 +1019,7 @@ public:
container.emplace_back(childPartitionId.value(), childPartitionId.value() + 1, range, parents);
}
alterData->TotalGroupCount += 2;
- ++alterData->ActivePartitionCount;
+ ++activePartitionCount;
return {};
};
@@ -997,7 +1053,7 @@ public:
// repeat splitting
size_t startIdx = 0;
- while (requestedMinPartitionCount > alterData->ActivePartitionCount) {
+ while (requestedMinPartitionCount > activePartitionCount) {
TVector<NKikimr::NSchemeShard::TTopicInfo::TPartitionToAdd> partitionsToAdd;
auto endIdx = alterData->PartitionsToAdd.size();
for (size_t i = startIdx; i < endIdx; ++i) {
@@ -1014,7 +1070,7 @@ public:
result->SetError(NKikimrScheme::StatusInvalidParameter, errStr);
return result;
}
- if (requestedMinPartitionCount <= alterData->ActivePartitionCount) {
+ if (requestedMinPartitionCount <= activePartitionCount) {
break;
}
}
@@ -1048,9 +1104,14 @@ public:
}
}
- alterData->TotalPartitionCount = topic->TotalPartitionCount + alterData->PartitionsToAdd.size();
- if (!splitMergeEnabled) {
- alterData->ActivePartitionCount = alterData->TotalPartitionCount;
+ ComputeAlterPartitionCounts(topic, alterData);
+
+ if (!(0 < alterData->ActivePartitionCount && alterData->ActivePartitionCount <= alterData->TotalPartitionCount)) {
+ errStr = TStringBuilder()
+ << "Invalid active partition count: " << alterData->ActivePartitionCount
+ << " vs total: " << alterData->TotalPartitionCount;
+ result->SetError(NKikimrScheme::StatusInvalidParameter, errStr);
+ return result;
}
alterData->NextPartitionId = topic->NextPartitionId;
@@ -1090,10 +1151,11 @@ public:
}
const auto& stats = topic->Stats;
+ const auto topicActivePartitionCount = CountTopicActivePartitions(topic);
const PQGroupReserve reserve(newTabletConfig, alterData->ActivePartitionCount);
const PQGroupReserve reserveForCheckLimit(newTabletConfig, alterData->ActivePartitionCount + involvedPartitions.size());
- const PQGroupReserve oldReserve(tabletConfig, topic->ActivePartitionCount);
- const PQGroupReserve oldReserveForCheckLimit(tabletConfig, topic->ActivePartitionCount, stats.DataSize);
+ const PQGroupReserve oldReserve(tabletConfig, topicActivePartitionCount);
+ const PQGroupReserve oldReserveForCheckLimit(tabletConfig, topicActivePartitionCount, stats.DataSize);
const ui64 storageToReserve = reserveForCheckLimit.Storage > oldReserveForCheckLimit.Storage ? reserveForCheckLimit.Storage - oldReserveForCheckLimit.Storage : 0;
diff --git a/ydb/core/tx/schemeshard/schemeshard_info_types.h b/ydb/core/tx/schemeshard/schemeshard_info_types.h
index 4613fdb87bf..09611b7de8d 100644
--- a/ydb/core/tx/schemeshard/schemeshard_info_types.h
+++ b/ydb/core/tx/schemeshard/schemeshard_info_types.h
@@ -1746,14 +1746,6 @@ struct TTopicInfo : TSimpleRefCount<TTopicInfo> {
bool HasBalancer() const { return bool(BalancerTabletID); }
- ui32 GetTotalPartitionCountWithAlter() const {
- ui32 res = 0;
- for (const auto& shard : Shards) {
- res += shard.second->PartsCount();
- }
- return res;
- }
-
ui32 ExpectedShardCount() const {
Y_ENSURE(TotalPartitionCount);
diff --git a/ydb/core/tx/schemeshard/ut_topic_splitmerge/ut_topic_splitmerge.cpp b/ydb/core/tx/schemeshard/ut_topic_splitmerge/ut_topic_splitmerge.cpp
index b7c072b4502..738aabd486d 100644
--- a/ydb/core/tx/schemeshard/ut_topic_splitmerge/ut_topic_splitmerge.cpp
+++ b/ydb/core/tx/schemeshard/ut_topic_splitmerge/ut_topic_splitmerge.cpp
@@ -955,6 +955,202 @@ Y_UNIT_TEST_SUITE(TSchemeShardTopicSplitMergeTest) {
} // Y_UNIT_TEST(SplitByMinPartitionCountWithTwoIter)
+ Y_UNIT_TEST(SplitByMinPartitionCountToSixteen) {
+ TTestBasicRuntime runtime;
+ TTestEnv env = CreateTestEnv(runtime);
+
+ ui64 txId = 100;
+
+ auto countActivePartitions = [](const auto& topic) {
+ ui32 activePartitions = 0;
+ for (const auto& partition : topic.GetPartitions()) {
+ if (partition.GetStatus() == NKikimrPQ::ETopicPartitionStatus::Active) {
+ ++activePartitions;
+ }
+ }
+ return activePartitions;
+ };
+
+ CreateSubDomain(runtime, env, ++txId);
+ CreateTopic(runtime, env, ++txId, 1);
+
+ ::NKikimrSchemeOp::TPersQueueGroupDescription scheme;
+ scheme.SetName("Topic1");
+ scheme.MutablePQTabletConfig()->MutablePartitionConfig();
+ scheme.MutablePQTabletConfig()->MutablePartitionStrategy()->SetMaxPartitionCount(100);
+ scheme.MutablePQTabletConfig()->MutablePartitionStrategy()->SetMinPartitionCount(16);
+ scheme.MutablePQTabletConfig()->MutablePartitionStrategy()->SetPartitionStrategyType(
+ ::NKikimrPQ::TPQTabletConfig_TPartitionStrategyType::TPQTabletConfig_TPartitionStrategyType_CAN_SPLIT_AND_MERGE);
+
+ TStringBuilder sb;
+ sb << scheme;
+ const TString schemeStr = sb.substr(1, sb.size() - 2);
+
+ AsyncAlterPQGroup(runtime, ++txId, "/MyRoot/USER_1", schemeStr);
+
+ env.TestWaitNotification(runtime, txId);
+
+ auto topic = DescribeTopic(runtime);
+ UNIT_ASSERT_VALUES_EQUAL(topic.GetPartitions().size(), 31);
+ UNIT_ASSERT_VALUES_EQUAL(countActivePartitions(topic), 16);
+
+ ModifyTopic(runtime, env, txId, [&](auto& alterScheme) {
+ alterScheme.MutablePQTabletConfig()->MutablePartitionStrategy()->SetMaxPartitionCount(100);
+ alterScheme.MutablePQTabletConfig()->MutablePartitionStrategy()->SetMinPartitionCount(32);
+ alterScheme.MutablePQTabletConfig()->MutablePartitionStrategy()->SetPartitionStrategyType(
+ ::NKikimrPQ::TPQTabletConfig_TPartitionStrategyType::TPQTabletConfig_TPartitionStrategyType_CAN_SPLIT_AND_MERGE);
+ });
+
+ topic = DescribeTopic(runtime);
+ UNIT_ASSERT_VALUES_EQUAL(topic.GetPartitions().size(), 63);
+ UNIT_ASSERT_VALUES_EQUAL(countActivePartitions(topic), 32);
+ } // Y_UNIT_TEST(SplitByMinPartitionCountToSixteen)
+
+ Y_UNIT_TEST(SplitByMinPartitionCountActivePartitionCountWithoutReboot) {
+ TTestBasicRuntime runtime;
+ TTestEnv env = CreateTestEnv(runtime);
+
+ ui64 txId = 100;
+
+ constexpr ui64 lifetimeSeconds = 3600;
+ constexpr ui64 writeSpeed = 1024;
+ const ui64 partitionReserveSize = lifetimeSeconds * writeSpeed;
+ const ui32 expectedActivePartitions = 16;
+
+ const auto AssertTopicReserve = [&](ui64 expectedActiveCount) {
+ TestDescribeResult(DescribePath(runtime, "/MyRoot/USER_1"),
+ {NLs::Finished,
+ NLs::TopicReservedStorage(expectedActiveCount * partitionReserveSize)});
+ };
+
+ CreateSubDomain(runtime, env, ++txId);
+
+ TestCreatePQGroup(runtime, ++txId, "/MyRoot/USER_1", TStringBuilder() << R"(
+ Name: "Topic1"
+ TotalGroupCount: 1
+ PartitionPerTablet: 7
+ PQTabletConfig {
+ PartitionConfig {
+ LifetimeSeconds: )" << lifetimeSeconds << R"(
+ WriteSpeedInBytesPerSecond : )" << writeSpeed << R"(
+ }
+ MeteringMode: METERING_MODE_RESERVED_CAPACITY
+ PartitionStrategy {
+ PartitionStrategyType: CAN_SPLIT_AND_MERGE
+ }
+ }
+ )");
+ env.TestWaitNotification(runtime, txId);
+ AssertTopicReserve(1);
+
+ ::NKikimrSchemeOp::TPersQueueGroupDescription scheme;
+ scheme.SetName("Topic1");
+ scheme.MutablePQTabletConfig()->MutablePartitionConfig()->SetLifetimeSeconds(lifetimeSeconds);
+ scheme.MutablePQTabletConfig()->MutablePartitionConfig()->SetWriteSpeedInBytesPerSecond(writeSpeed);
+ scheme.MutablePQTabletConfig()->SetMeteringMode(
+ ::NKikimrPQ::TPQTabletConfig::METERING_MODE_RESERVED_CAPACITY);
+ scheme.MutablePQTabletConfig()->MutablePartitionStrategy()->SetMaxPartitionCount(100);
+ scheme.MutablePQTabletConfig()->MutablePartitionStrategy()->SetMinPartitionCount(expectedActivePartitions);
+ scheme.MutablePQTabletConfig()->MutablePartitionStrategy()->SetPartitionStrategyType(
+ ::NKikimrPQ::TPQTabletConfig_TPartitionStrategyType::TPQTabletConfig_TPartitionStrategyType_CAN_SPLIT_AND_MERGE);
+
+ TStringBuilder sb;
+ sb << scheme;
+ const TString schemeStr = sb.substr(1, sb.size() - 2);
+
+ AsyncAlterPQGroup(runtime, ++txId, "/MyRoot/USER_1", schemeStr);
+ env.TestWaitNotification(runtime, txId);
+
+ auto topic = DescribeTopic(runtime);
+ UNIT_ASSERT_VALUES_EQUAL(topic.GetPartitions().size(), 31);
+
+ ui32 activePartitions = 0;
+ for (const auto& partition : topic.GetPartitions()) {
+ if (partition.GetStatus() == NKikimrPQ::ETopicPartitionStatus::Active) {
+ ++activePartitions;
+ }
+ }
+ UNIT_ASSERT_VALUES_EQUAL(activePartitions, expectedActivePartitions);
+
+ // ComputeAlterActivePartitionCount must deduct intermediate parents from PartitionsToAdd.
+ // Without reboot, reserve is based on the stored ActivePartitionCount, not partition statuses.
+ AssertTopicReserve(expectedActivePartitions);
+ } // Y_UNIT_TEST(SplitByMinPartitionCountActivePartitionCountWithoutReboot)
+
+ Y_UNIT_TEST(CreateRootLevelSiblingPartitions) {
+ TTestBasicRuntime runtime;
+ TTestEnv env = CreateTestEnv(runtime);
+
+ ui64 txId = 100;
+
+ auto countActivePartitions = [](const auto& topic) {
+ ui32 activePartitions = 0;
+ for (const auto& partition : topic.GetPartitions()) {
+ if (partition.GetStatus() == NKikimrPQ::ETopicPartitionStatus::Active) {
+ ++activePartitions;
+ }
+ }
+ return activePartitions;
+ };
+
+ TString bound0((char*)bound_1_3, sizeof(bound_1_3));
+ TString bound1((char*)bound_2_3, sizeof(bound_2_3));
+ const unsigned char partitionBoundaryBytes[] = {0xBF};
+ const unsigned char splitBoundaryBytes[] = {0xB5};
+ TString partitionBoundary((char*)partitionBoundaryBytes, sizeof(partitionBoundaryBytes));
+ TString splitBoundary((char*)splitBoundaryBytes, sizeof(splitBoundaryBytes));
+
+ CreateSubDomain(runtime, env, ++txId);
+ CreateTopic(runtime, env, ++txId, 1);
+
+ ModifyTopic(runtime, env, txId, [&](auto& scheme) {
+ auto* boundary0 = scheme.AddRootPartitionBoundaries();
+ boundary0->SetPartition(0);
+ boundary0->MutableKeyRange()->SetToBound(bound0);
+ auto* boundary1 = scheme.AddRootPartitionBoundaries();
+ boundary1->SetPartition(1);
+ boundary1->SetCreatePartition(true);
+ boundary1->MutableKeyRange()->SetFromBound(bound0);
+ boundary1->MutableKeyRange()->SetToBound(bound1);
+ auto* boundary2 = scheme.AddRootPartitionBoundaries();
+ boundary2->SetPartition(2);
+ boundary2->SetCreatePartition(true);
+ boundary2->MutableKeyRange()->SetFromBound(bound1);
+ });
+
+ auto topic = DescribeTopic(runtime);
+ UNIT_ASSERT_VALUES_EQUAL(topic.GetPartitions().size(), 3);
+ UNIT_ASSERT_VALUES_EQUAL(countActivePartitions(topic), 3);
+
+ ModifyTopic(runtime, env, txId, [&](auto& scheme) {
+ auto* boundary0 = scheme.AddRootPartitionBoundaries();
+ boundary0->SetPartition(0);
+ boundary0->MutableKeyRange()->SetToBound(bound0);
+ auto* boundary1 = scheme.AddRootPartitionBoundaries();
+ boundary1->SetPartition(1);
+ boundary1->MutableKeyRange()->SetFromBound(bound0);
+ boundary1->MutableKeyRange()->SetToBound(bound1);
+ auto* boundary2 = scheme.AddRootPartitionBoundaries();
+ boundary2->SetPartition(2);
+ boundary2->MutableKeyRange()->SetFromBound(bound1);
+ boundary2->MutableKeyRange()->SetToBound(partitionBoundary);
+ auto* boundary3 = scheme.AddRootPartitionBoundaries();
+ boundary3->SetPartition(3);
+ boundary3->SetCreatePartition(true);
+ boundary3->MutableKeyRange()->SetFromBound(partitionBoundary);
+ });
+
+ topic = DescribeTopic(runtime);
+ UNIT_ASSERT_VALUES_EQUAL(topic.GetPartitions().size(), 4);
+ UNIT_ASSERT_VALUES_EQUAL(countActivePartitions(topic), 4);
+
+ SplitPartitionTo(runtime, env, txId, 2, splitBoundary, {3, 4}, true);
+
+ topic = DescribeTopic(runtime);
+ UNIT_ASSERT_VALUES_EQUAL(topic.GetPartitions().size(), 5);
+ UNIT_ASSERT_VALUES_EQUAL(countActivePartitions(topic), 5);
+ } // Y_UNIT_TEST(CreateRootLevelSiblingPartitions)
+
} // Y_UNIT_TEST_SUITE(TSchemeShardTopicSplitMergeTest)
Y_UNIT_TEST_SUITE(TSchemeShardTopicSplitMergePrescribedPartitionsTest) {
@@ -1396,5 +1592,92 @@ Y_UNIT_TEST_SUITE(TSchemeShardTopicSplitMergePrescribedPartitionsTest) {
ValidatePartitionChildren(partition2, {});
ValidatePartitionChildren(partition3, {});
- } // Y_UNIT_TEST
+ } // Y_UNIT_TEST(SplitWithExistingPartitionWithPartialOverlapAndCreateRootLevelSibling)
+
+ Y_UNIT_TEST(AlterTopicConfigAfterSplitAlterInterruptedByReboot) {
+ TTestBasicRuntime runtime;
+ TTestEnv env = CreateTestEnv(runtime);
+
+ ui64 txId = 100;
+
+ CreateSubDomain(runtime, env, ++txId);
+ CreateTopic(runtime, env, ++txId, 3);
+
+ const unsigned char b[] = {0x7F};
+ TString boundary((char*)b, sizeof(b));
+
+ ::NKikimrSchemeOp::TPersQueueGroupDescription scheme;
+ scheme.SetName("Topic1");
+ scheme.MutablePQTabletConfig()->MutablePartitionConfig();
+ auto* split = scheme.AddSplit();
+ split->SetPartition(1);
+ split->SetSplitBoundary(boundary);
+
+ TStringBuilder sb;
+ sb << scheme;
+ const TString schemeStr = sb.substr(1, sb.size() - 2);
+
+ AsyncAlterPQGroup(runtime, ++txId, "/MyRoot/USER_1", schemeStr);
+
+ RebootTablet(runtime, TTestTxConfig::SchemeShard, runtime.AllocateEdgeActor());
+
+ env.TestWaitNotification(runtime, txId);
+
+ auto topic = DescribeTopic(runtime);
+ UNIT_ASSERT_VALUES_EQUAL(topic.GetPartitions().size(), 5);
+
+ ui32 activePartitions = 0;
+ for (const auto& partition : topic.GetPartitions()) {
+ if (partition.GetStatus() == NKikimrPQ::ETopicPartitionStatus::Active) {
+ ++activePartitions;
+ }
+ }
+ UNIT_ASSERT_VALUES_EQUAL(activePartitions, 4);
+
+ ModifyTopic(runtime, env, txId, [&](auto& alter) {
+ alter.MutablePQTabletConfig()->MutablePartitionConfig()->SetLifetimeSeconds(7200);
+ });
+ } // Y_UNIT_TEST(AlterTopicConfigAfterSplitAlterInterruptedByReboot)
+
+ Y_UNIT_TEST(AlterTopicConfigAfterMergeAlterInterruptedByReboot) {
+ TTestBasicRuntime runtime;
+ TTestEnv env = CreateTestEnv(runtime);
+
+ ui64 txId = 100;
+
+ CreateSubDomain(runtime, env, ++txId);
+ CreateTopic(runtime, env, ++txId, 3);
+
+ ::NKikimrSchemeOp::TPersQueueGroupDescription scheme;
+ scheme.SetName("Topic1");
+ scheme.MutablePQTabletConfig()->MutablePartitionConfig();
+ auto* merge = scheme.AddMerge();
+ merge->SetPartition(0);
+ merge->SetAdjacentPartition(1);
+
+ TStringBuilder sb;
+ sb << scheme;
+ const TString schemeStr = sb.substr(1, sb.size() - 2);
+
+ AsyncAlterPQGroup(runtime, ++txId, "/MyRoot/USER_1", schemeStr);
+
+ RebootTablet(runtime, TTestTxConfig::SchemeShard, runtime.AllocateEdgeActor());
+
+ env.TestWaitNotification(runtime, txId);
+
+ auto topic = DescribeTopic(runtime);
+ UNIT_ASSERT_VALUES_EQUAL(topic.GetPartitions().size(), 4);
+
+ ui32 activePartitions = 0;
+ for (const auto& partition : topic.GetPartitions()) {
+ if (partition.GetStatus() == NKikimrPQ::ETopicPartitionStatus::Active) {
+ ++activePartitions;
+ }
+ }
+ UNIT_ASSERT_VALUES_EQUAL(activePartitions, 2);
+
+ ModifyTopic(runtime, env, txId, [&](auto& alter) {
+ alter.MutablePQTabletConfig()->MutablePartitionConfig()->SetLifetimeSeconds(7200);
+ });
+ } // Y_UNIT_TEST(AlterTopicConfigAfterMergeAlterInterruptedByReboot)
}