diff options
| author | ivanmorozov333 <[email protected]> | 2024-06-10 15:27:48 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2024-06-10 15:27:48 +0300 |
| commit | 4cef62d2b9199a9d72a23c68179ef46d8abe05ef (patch) | |
| tree | 2063ab957fc45884715acc02cbc6051b988401ed | |
| parent | 1d951c4cf7f550a09cbf91b3b09b16e8511cc441 (diff) | |
fixes for resharding tests (#5360)
46 files changed, 372 insertions, 186 deletions
diff --git a/ydb/core/kqp/ut/olap/blobs_sharing_ut.cpp b/ydb/core/kqp/ut/olap/blobs_sharing_ut.cpp index de619f580a7..0ee5db09f4f 100644 --- a/ydb/core/kqp/ut/olap/blobs_sharing_ut.cpp +++ b/ydb/core/kqp/ut/olap/blobs_sharing_ut.cpp @@ -96,7 +96,7 @@ Y_UNIT_TEST_SUITE(KqpOlapBlobsSharing) { Controller->SetPeriodicWakeupActivationPeriod(TDuration::Seconds(1)); Controller->SetReadTimeoutClean(TDuration::Seconds(1)); - Tests::NCommon::TLoggerInit(Kikimr).SetComponents({NKikimrServices::TX_COLUMNSHARD}, "CS").Initialize(); + Tests::NCommon::TLoggerInit(Kikimr).SetComponents({ NKikimrServices::TX_COLUMNSHARD }, "CS").Initialize(); Helper.CreateTestOlapTable(ShardsCount, ShardsCount); ShardIds = Controller->GetShardActualIds(); @@ -190,14 +190,13 @@ Y_UNIT_TEST_SUITE(KqpOlapBlobsSharing) { AFL_VERIFY(!Controller->IsTrivialLinks()); } }; - Y_UNIT_TEST(BlobsSharingSplit1_1) { auto settings = TKikimrSettings().SetWithSampleTables(false); TKikimrRunner kikimr(settings); TSharingDataTestCase tester(4, kikimr); tester.AddRecords(800000); Sleep(TDuration::Seconds(1)); - tester.Execute(0, {1}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); + tester.Execute(0, { 1 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); } Y_UNIT_TEST(BlobsSharingSplit1_1_clean) { @@ -207,7 +206,7 @@ Y_UNIT_TEST_SUITE(KqpOlapBlobsSharing) { tester.AddRecords(80000); CompareYson(tester.GetHelper().GetQueryResult("SELECT COUNT(*) FROM `/Root/olapStore12/olapTable`"), R"([[80000u;]])"); Sleep(TDuration::Seconds(1)); - tester.Execute(0, {1}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); + tester.Execute(0, { 1 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); CompareYson(tester.GetHelper().GetQueryResult("SELECT COUNT(*) FROM `/Root/olapStore12/olapTable`"), R"([[119928u;]])"); tester.AddRecords(80000, 0.8); tester.WaitNormalization(); @@ -222,7 +221,7 @@ Y_UNIT_TEST_SUITE(KqpOlapBlobsSharing) { tester.AddRecords(80000); CompareYson(tester.GetHelper().GetQueryResult("SELECT COUNT(*) FROM `/Root/olapStore12/olapTable`"), R"([[80000u;]])"); Sleep(TDuration::Seconds(1)); - tester.Execute(0, {1}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); + tester.Execute(0, { 1 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); CompareYson(tester.GetHelper().GetQueryResult("SELECT COUNT(*) FROM `/Root/olapStore12/olapTable`"), R"([[119928u;]])"); tester.AddRecords(80000, 0.8); tester.WaitNormalization(); @@ -235,7 +234,7 @@ Y_UNIT_TEST_SUITE(KqpOlapBlobsSharing) { TSharingDataTestCase tester(4, kikimr); tester.AddRecords(800000); Sleep(TDuration::Seconds(1)); - tester.Execute(0, {1, 2, 3}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); + tester.Execute(0, { 1, 2, 3 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); } Y_UNIT_TEST(BlobsSharingSplit1_3_1) { @@ -244,10 +243,10 @@ Y_UNIT_TEST_SUITE(KqpOlapBlobsSharing) { TSharingDataTestCase tester(4, kikimr); tester.AddRecords(800000); Sleep(TDuration::Seconds(1)); - tester.Execute(1, {0}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); - tester.Execute(2, {0}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); - tester.Execute(3, {0}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); - tester.Execute(0, {1, 2, 3}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); + tester.Execute(1, { 0 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); + tester.Execute(2, { 0 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); + tester.Execute(3, { 0 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); + tester.Execute(0, { 1, 2, 3 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); } Y_UNIT_TEST(BlobsSharingSplit1_3_2_1_clean) { @@ -256,13 +255,13 @@ Y_UNIT_TEST_SUITE(KqpOlapBlobsSharing) { TSharingDataTestCase tester(4, kikimr); tester.AddRecords(800000); Sleep(TDuration::Seconds(1)); - tester.Execute(1, {0}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); - tester.Execute(2, {0}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); - tester.Execute(3, {0}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); + tester.Execute(1, { 0 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); + tester.Execute(2, { 0 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); + tester.Execute(3, { 0 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); tester.AddRecords(800000, 0.9); Sleep(TDuration::Seconds(1)); - tester.Execute(3, {2}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); - tester.Execute(0, {1, 2}, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), {0}); + tester.Execute(3, { 2 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); + tester.Execute(0, { 1, 2 }, false, NOlap::TSnapshot(TInstant::Now().MilliSeconds(), 1232123), { 0 }); tester.WaitNormalization(); } @@ -308,8 +307,7 @@ Y_UNIT_TEST_SUITE(KqpOlapBlobsSharing) { public: TReshardingTest() - : Kikimr(TKikimrSettings().SetWithSampleTables(false)) - { + : Kikimr(TKikimrSettings().SetWithSampleTables(false)) { } @@ -387,7 +385,7 @@ Y_UNIT_TEST_SUITE(KqpOlapBlobsSharing) { UNIT_ASSERT_VALUES_EQUAL_C(alterResult.GetStatus(), NYdb::EStatus::SUCCESS, alterResult.GetIssues().ToString()); } WaitResharding(); - csController->WaitCleaning(TDuration::Seconds(5)); + // csController->WaitCleaning(TDuration::Seconds(5)); CheckCount(230000); AFL_VERIFY(count + portionsCount == csController->GetShardingFiltersCount().Val())("count", count)("val", csController->GetShardingFiltersCount().Val()); @@ -409,7 +407,5 @@ Y_UNIT_TEST_SUITE(KqpOlapBlobsSharing) { Y_UNIT_TEST(TableReshardingModuloN) { TReshardingTest().SetShardingType("HASH_FUNCTION_MODULO_N").Execute(); } - } - } diff --git a/ydb/core/tx/columnshard/blobs_action/abstract/blob_set.cpp b/ydb/core/tx/columnshard/blobs_action/abstract/blob_set.cpp index b0952746b9b..0e4ee2a2063 100644 --- a/ydb/core/tx/columnshard/blobs_action/abstract/blob_set.cpp +++ b/ydb/core/tx/columnshard/blobs_action/abstract/blob_set.cpp @@ -3,6 +3,7 @@ #include <ydb/core/tx/columnshard/blobs_action/protos/blobs.pb.h> #include <util/generic/refcount.h> +#include <util/string/join.h> namespace NKikimr::NOlap { @@ -54,4 +55,14 @@ NKikimr::TConclusionStatus TTabletsByBlob::DeserializeFromProto(const NKikimrCol return TConclusionStatus::Success(); } +TString TTabletsByBlob::DebugString() const { + TStringBuilder sb; + for (auto&& i : Data) { + sb << "["; + sb << i.first.ToStringNew() << ":" << JoinSeq(",", i.second); + sb << "];"; + } + return sb; +} + } diff --git a/ydb/core/tx/columnshard/blobs_action/abstract/blob_set.h b/ydb/core/tx/columnshard/blobs_action/abstract/blob_set.h index 7692dad22eb..d14d7b2dee1 100644 --- a/ydb/core/tx/columnshard/blobs_action/abstract/blob_set.h +++ b/ydb/core/tx/columnshard/blobs_action/abstract/blob_set.h @@ -261,6 +261,8 @@ public: } return true; } + + TString DebugString() const; }; class TBlobsByTablet { @@ -336,6 +338,14 @@ public: std::swap(result, resultLocal); } + bool Contains(const TTabletId tabletId, const TUnifiedBlobId& blobId) const { + auto it = Data.find(tabletId); + if (it == Data.end()) { + return false; + } + return it->second.contains(blobId); + } + const THashSet<TUnifiedBlobId>* Find(const TTabletId tabletId) const { auto it = Data.find(tabletId); if (it == Data.end()) { @@ -384,6 +394,14 @@ public: } } + ui32 GetSize() const { + ui32 result = 0; + for (auto&& i : Data) { + result += i.second.size(); + } + return result; + } + bool IsEmpty() const { return Data.empty(); } @@ -423,6 +441,10 @@ public: return Sharing.IsEmpty() && Direct.IsEmpty() && Borrowed.IsEmpty(); } + bool HasSharingOnly() const { + return !Sharing.IsEmpty() && Direct.IsEmpty() && Borrowed.IsEmpty(); + } + class TIterator { private: const TBlobsCategories* Owner; diff --git a/ydb/core/tx/columnshard/blobs_action/abstract/gc.h b/ydb/core/tx/columnshard/blobs_action/abstract/gc.h index 2d0a5e44732..19e2da1b39b 100644 --- a/ydb/core/tx/columnshard/blobs_action/abstract/gc.h +++ b/ydb/core/tx/columnshard/blobs_action/abstract/gc.h @@ -32,6 +32,10 @@ protected: virtual void RemoveBlobIdFromDB(const TTabletId tabletId, const TUnifiedBlobId& blobId, TBlobManagerDb& dbBlobs) = 0; virtual bool DoIsEmpty() const = 0; public: + void AddSharedBlobToNextIteration(const TUnifiedBlobId& blobId, const TTabletId ownerTabletId) { + BlobsToRemove.RemoveSharing(ownerTabletId, blobId); + } + void OnExecuteTxAfterCleaning(NColumnShard::TColumnShard& self, TBlobManagerDb& dbBlobs); void OnCompleteTxAfterCleaning(NColumnShard::TColumnShard& self, const std::shared_ptr<IBlobsGCAction>& taskAction); diff --git a/ydb/core/tx/columnshard/blobs_action/abstract/gc_actor.h b/ydb/core/tx/columnshard/blobs_action/abstract/gc_actor.h index 7acd09a7c16..62008388c94 100644 --- a/ydb/core/tx/columnshard/blobs_action/abstract/gc_actor.h +++ b/ydb/core/tx/columnshard/blobs_action/abstract/gc_actor.h @@ -1,6 +1,7 @@ #pragma once #include "blob_set.h" #include "common.h" +#include "gc.h" #include <ydb/core/base/tablet_pipecache.h> @@ -16,13 +17,30 @@ private: const TString OperatorId; TBlobsByTablet BlobIdsByTablets; const TTabletId SelfTabletId; + std::shared_ptr<IBlobsGCAction> GCAction; virtual void DoOnSharedRemovingFinished() = 0; void OnSharedRemovingFinished() { SharedRemovingFinished = true; DoOnSharedRemovingFinished(); } void Handle(NEvents::TEvDeleteSharedBlobsFinished::TPtr& ev) { - AFL_VERIFY(BlobIdsByTablets.Remove((TTabletId)ev->Get()->Record.GetTabletId())); + const TTabletId sourceTabletId = (TTabletId)ev->Get()->Record.GetTabletId(); + auto* blobIds = BlobIdsByTablets.Find(sourceTabletId); + AFL_VERIFY(blobIds); + switch (ev->Get()->Record.GetStatus()) { + case NKikimrColumnShardBlobOperationsProto::TEvDeleteSharedBlobsFinished::Success: + AFL_VERIFY(BlobIdsByTablets.Remove(sourceTabletId)); + break; + case NKikimrColumnShardBlobOperationsProto::TEvDeleteSharedBlobsFinished::DestinationCurrenlyLocked: + for (auto&& i : *blobIds) { + GCAction->AddSharedBlobToNextIteration(i, sourceTabletId); + } + AFL_VERIFY(BlobIdsByTablets.Remove(sourceTabletId)); + break; + case NKikimrColumnShardBlobOperationsProto::TEvDeleteSharedBlobsFinished::NeedRetry: + SendToTablet(*blobIds, sourceTabletId); + break; + } if (BlobIdsByTablets.IsEmpty()) { AFL_VERIFY(!SharedRemovingFinished); OnSharedRemovingFinished(); @@ -31,19 +49,24 @@ private: void Handle(NActors::TEvents::TEvUndelivered::TPtr& ev) { auto* blobIds = BlobIdsByTablets.Find((TTabletId)ev->Cookie); AFL_VERIFY(blobIds); - auto evResend = std::make_unique<NEvents::TEvDeleteSharedBlobs>(TBase::SelfId(), ev->Cookie, OperatorId, *blobIds); + SendToTablet(*blobIds, (TTabletId)ev->Cookie); + } + + void SendToTablet(const THashSet<TUnifiedBlobId>& blobIds, const TTabletId tabletId) const { + auto ev = std::make_unique<NEvents::TEvDeleteSharedBlobs>(TBase::SelfId(), (ui64)SelfTabletId, OperatorId, blobIds); NActors::TActivationContext::AsActorContext().Send(MakePipePerNodeCacheID(false), - new TEvPipeCache::TEvForward(evResend.release(), ev->Cookie, true), IEventHandle::FlagTrackDelivery, ev->Cookie); + new TEvPipeCache::TEvForward(ev.release(), (ui64)tabletId, true), IEventHandle::FlagTrackDelivery, (ui64)tabletId); } protected: bool SharedRemovingFinished = false; public: - TSharedBlobsCollectionActor(const TString& operatorId, const TTabletId selfTabletId, const TBlobsByTablet& blobIds) + TSharedBlobsCollectionActor(const TString& operatorId, const TTabletId selfTabletId, const TBlobsByTablet& blobIds, const std::shared_ptr<IBlobsGCAction>& gcAction) : OperatorId(operatorId) , BlobIdsByTablets(blobIds) , SelfTabletId(selfTabletId) + , GCAction(gcAction) { - + AFL_VERIFY(GCAction); } STFUNC(StateWork) { @@ -60,9 +83,7 @@ public: OnSharedRemovingFinished(); } else { for (auto&& i : BlobIdsByTablets) { - auto ev = std::make_unique<NEvents::TEvDeleteSharedBlobs>(TBase::SelfId(), (ui64)SelfTabletId, OperatorId, i.second); - NActors::TActivationContext::AsActorContext().Send(MakePipePerNodeCacheID(false), - new TEvPipeCache::TEvForward(ev.release(), (ui64)i.first, true), IEventHandle::FlagTrackDelivery, (ui64)i.first); + SendToTablet(i.second, i.first); } } } diff --git a/ydb/core/tx/columnshard/blobs_action/abstract/remove.h b/ydb/core/tx/columnshard/blobs_action/abstract/remove.h index e48ae7c2bb5..8de1aa69cc1 100644 --- a/ydb/core/tx/columnshard/blobs_action/abstract/remove.h +++ b/ydb/core/tx/columnshard/blobs_action/abstract/remove.h @@ -32,6 +32,10 @@ public: } + TTabletId GetSelfTabletId() const { + return SelfTabletId; + } + void DeclareRemove(const TTabletId tabletId, const TUnifiedBlobId& blobId); void DeclareSelfRemove(const TUnifiedBlobId& blobId); void OnExecuteTxAfterRemoving(TBlobManagerDb& dbBlobs, const bool blobsWroteSuccessfully) { diff --git a/ydb/core/tx/columnshard/blobs_action/abstract/storage.h b/ydb/core/tx/columnshard/blobs_action/abstract/storage.h index 19a3655c17b..778a06b45ce 100644 --- a/ydb/core/tx/columnshard/blobs_action/abstract/storage.h +++ b/ydb/core/tx/columnshard/blobs_action/abstract/storage.h @@ -78,6 +78,7 @@ public: } virtual TTabletsByBlob GetBlobsToDelete() const = 0; + virtual bool HasToDelete(const TUnifiedBlobId& blobId, const TTabletId initiatorTabletId) const = 0; virtual std::shared_ptr<IBlobInUseTracker> GetBlobsTracker() const = 0; virtual ~IBlobsStorageOperator() = default; diff --git a/ydb/core/tx/columnshard/blobs_action/abstract/storages_manager.h b/ydb/core/tx/columnshard/blobs_action/abstract/storages_manager.h index cc0e8f606d3..e1e07a41206 100644 --- a/ydb/core/tx/columnshard/blobs_action/abstract/storages_manager.h +++ b/ydb/core/tx/columnshard/blobs_action/abstract/storages_manager.h @@ -38,9 +38,7 @@ public: } bool LoadIdempotency(NTable::TDatabase& database); - bool HasBlobsToDelete() const; - void Stop(); std::shared_ptr<IBlobsStorageOperator> GetDefaultOperator() const { @@ -51,7 +49,7 @@ public: return GetDefaultOperator(); } - const THashMap<TString, std::shared_ptr<IBlobsStorageOperator>>& GetStorages() { + const THashMap<TString, std::shared_ptr<IBlobsStorageOperator>>& GetStorages() const { AFL_VERIFY(Initialized); return Constructed; } diff --git a/ydb/core/tx/columnshard/blobs_action/bs/blob_manager.h b/ydb/core/tx/columnshard/blobs_action/bs/blob_manager.h index 8a5ff68eb9f..8f207724c67 100644 --- a/ydb/core/tx/columnshard/blobs_action/bs/blob_manager.h +++ b/ydb/core/tx/columnshard/blobs_action/bs/blob_manager.h @@ -177,6 +177,10 @@ private: public: TBlobManager(TIntrusivePtr<TTabletStorageInfo> tabletInfo, const ui32 gen, const TTabletId selfTabletId); + bool HasToDelete(const TUnifiedBlobId& blobId, const TTabletId tabletId) const { + return BlobsToDelete.Contains(tabletId, blobId) || BlobsToDeleteDelayed.Contains(tabletId, blobId); + } + TTabletsByBlob GetBlobsToDeleteAll() const { auto result = BlobsToDelete; result.Add(BlobsToDeleteDelayed); diff --git a/ydb/core/tx/columnshard/blobs_action/bs/gc_actor.h b/ydb/core/tx/columnshard/blobs_action/bs/gc_actor.h index 8324ff75817..cddb4244816 100644 --- a/ydb/core/tx/columnshard/blobs_action/bs/gc_actor.h +++ b/ydb/core/tx/columnshard/blobs_action/bs/gc_actor.h @@ -20,7 +20,7 @@ private: } public: TGarbageCollectionActor(const std::shared_ptr<TGCTask>& task, const NActors::TActorId& tabletActorId, const TTabletId selfTabletId) - : TBase(task->GetStorageId(), selfTabletId, task->GetBlobsToRemove().GetBorrowed()) + : TBase(task->GetStorageId(), selfTabletId, task->GetBlobsToRemove().GetBorrowed(), task) , TabletActorId(tabletActorId) , GCTask(task) { diff --git a/ydb/core/tx/columnshard/blobs_action/bs/storage.h b/ydb/core/tx/columnshard/blobs_action/bs/storage.h index de22b03d818..1f726c13863 100644 --- a/ydb/core/tx/columnshard/blobs_action/bs/storage.h +++ b/ydb/core/tx/columnshard/blobs_action/bs/storage.h @@ -31,6 +31,10 @@ public: TOperator(const TString& storageId, const NActors::TActorId& tabletActorId, const TIntrusivePtr<TTabletStorageInfo>& tabletInfo, const ui64 generation, const std::shared_ptr<NDataSharing::TStorageSharedBlobsManager>& sharedBlobs); + virtual bool HasToDelete(const TUnifiedBlobId& blobId, const TTabletId tabletId) const override { + return Manager->HasToDelete(blobId, tabletId); + } + virtual TTabletsByBlob GetBlobsToDelete() const override { return Manager->GetBlobsToDeleteAll(); } diff --git a/ydb/core/tx/columnshard/blobs_action/events/delete_blobs.h b/ydb/core/tx/columnshard/blobs_action/events/delete_blobs.h index 50e37bf7fe0..45800257f0f 100644 --- a/ydb/core/tx/columnshard/blobs_action/events/delete_blobs.h +++ b/ydb/core/tx/columnshard/blobs_action/events/delete_blobs.h @@ -23,9 +23,9 @@ struct TEvDeleteSharedBlobs: public NActors::TEventPB<TEvDeleteSharedBlobs, NKik struct TEvDeleteSharedBlobsFinished: public NActors::TEventPB<TEvDeleteSharedBlobsFinished, NKikimrColumnShardBlobOperationsProto::TEvDeleteSharedBlobsFinished, TEvColumnShard::EvDeleteSharedBlobsFinished> { TEvDeleteSharedBlobsFinished() = default; - TEvDeleteSharedBlobsFinished(const TTabletId tabletId) - { + TEvDeleteSharedBlobsFinished(const TTabletId tabletId, const NKikimrColumnShardBlobOperationsProto::TEvDeleteSharedBlobsFinished::EStatus status) { Record.SetTabletId((ui64)tabletId); + Record.SetStatus(status); } }; diff --git a/ydb/core/tx/columnshard/blobs_action/protos/events.proto b/ydb/core/tx/columnshard/blobs_action/protos/events.proto index 34ae1fc57e8..1b008b1366a 100644 --- a/ydb/core/tx/columnshard/blobs_action/protos/events.proto +++ b/ydb/core/tx/columnshard/blobs_action/protos/events.proto @@ -11,4 +11,11 @@ message TEvDeleteSharedBlobs { message TEvDeleteSharedBlobsFinished { optional uint64 TabletId = 1; + enum EStatus { + Success = 0; + NeedRetry = 1; + DestinationCurrenlyLocked = 2; + } + + optional EStatus Status = 2; } diff --git a/ydb/core/tx/columnshard/blobs_action/tier/gc_actor.h b/ydb/core/tx/columnshard/blobs_action/tier/gc_actor.h index 78ee58e9152..4834b0e3911 100644 --- a/ydb/core/tx/columnshard/blobs_action/tier/gc_actor.h +++ b/ydb/core/tx/columnshard/blobs_action/tier/gc_actor.h @@ -23,7 +23,7 @@ private: } public: TGarbageCollectionActor(const std::shared_ptr<TGCTask>& task, const NActors::TActorId& tabletActorId, const TTabletId selfTabletId) - : TBase(task->GetStorageId(), selfTabletId, task->GetBlobsToRemove().GetBorrowed()) + : TBase(task->GetStorageId(), selfTabletId, task->GetBlobsToRemove().GetBorrowed(), task) , TabletActorId(tabletActorId) , GCTask(task) { diff --git a/ydb/core/tx/columnshard/blobs_action/tier/gc_info.h b/ydb/core/tx/columnshard/blobs_action/tier/gc_info.h index 4c710080115..564da152453 100644 --- a/ydb/core/tx/columnshard/blobs_action/tier/gc_info.h +++ b/ydb/core/tx/columnshard/blobs_action/tier/gc_info.h @@ -11,6 +11,10 @@ private: YDB_ACCESSOR_DEF(std::deque<TUnifiedBlobId>, DraftBlobIdsToRemove); YDB_ACCESSOR_DEF(TTabletsByBlob, BlobsToDeleteInFuture); public: + bool HasToDelete(const TUnifiedBlobId& blobId, const TTabletId tabletId) const { + return BlobsToDelete.Contains(tabletId, blobId) || BlobsToDeleteInFuture.Contains(tabletId, blobId); + } + virtual void OnBlobFree(const TUnifiedBlobId& blobId) override { BlobsToDeleteInFuture.ExtractBlobTo(blobId, BlobsToDelete); } @@ -30,6 +34,7 @@ public: while (BlobsToDelete.ExtractFrontTo(deleteBlobIdsLocal) && count < blobsCountLimit) { ++count; } + AFL_DEBUG(NKikimrServices::TX_COLUMNSHARD)("event", "extract_blobs_to_gc")("blob_ids", deleteBlobIdsLocal.DebugString()); std::swap(deleteBlobIdsLocal, deleteBlobIds); return true; } diff --git a/ydb/core/tx/columnshard/blobs_action/tier/remove.cpp b/ydb/core/tx/columnshard/blobs_action/tier/remove.cpp index 5391f513279..f3ad2b8ec19 100644 --- a/ydb/core/tx/columnshard/blobs_action/tier/remove.cpp +++ b/ydb/core/tx/columnshard/blobs_action/tier/remove.cpp @@ -1,5 +1,28 @@ #include "remove.h" +#include <util/string/join.h> namespace NKikimr::NOlap::NBlobOperations::NTier { +void TDeclareRemovingAction::DoOnCompleteTxAfterRemoving(const bool blobsWroteSuccessfully) { + if (blobsWroteSuccessfully) { + for (auto&& i : GetDeclaredBlobs()) { + if (GCInfo->IsBlobInUsage(i.first)) { + AFL_VERIFY(GCInfo->MutableBlobsToDeleteInFuture().Add(i.first, i.second)); + } else { + AFL_DEBUG(NKikimrServices::TX_COLUMNSHARD)("event", "blob_to_delete") + ("blob_id", i.first)("tablet_ids", JoinSeq(",", i.second)); + AFL_VERIFY(GCInfo->MutableBlobsToDelete().Add(i.first, i.second)); + } + } + } +} + +void TDeclareRemovingAction::DoOnExecuteTxAfterRemoving(TBlobManagerDb& dbBlobs, const bool blobsWroteSuccessfully) { + if (blobsWroteSuccessfully) { + for (auto i = GetDeclaredBlobs().GetIterator(); i.IsValid(); ++i) { + dbBlobs.AddTierBlobToDelete(GetStorageId(), i.GetBlobId(), i.GetTabletId()); + } + } +} + } diff --git a/ydb/core/tx/columnshard/blobs_action/tier/remove.h b/ydb/core/tx/columnshard/blobs_action/tier/remove.h index 9dbd6fb0be9..c19ed975d3a 100644 --- a/ydb/core/tx/columnshard/blobs_action/tier/remove.h +++ b/ydb/core/tx/columnshard/blobs_action/tier/remove.h @@ -15,24 +15,8 @@ protected: } - virtual void DoOnExecuteTxAfterRemoving(TBlobManagerDb& dbBlobs, const bool blobsWroteSuccessfully) { - if (blobsWroteSuccessfully) { - for (auto i = GetDeclaredBlobs().GetIterator(); i.IsValid(); ++i) { - dbBlobs.AddTierBlobToDelete(GetStorageId(), i.GetBlobId(), i.GetTabletId()); - } - } - } - virtual void DoOnCompleteTxAfterRemoving(const bool blobsWroteSuccessfully) { - if (blobsWroteSuccessfully) { - for (auto&& i : GetDeclaredBlobs()) { - if (GCInfo->IsBlobInUsage(i.first)) { - AFL_VERIFY(GCInfo->MutableBlobsToDeleteInFuture().Add(i.first, i.second)); - } else { - AFL_VERIFY(GCInfo->MutableBlobsToDelete().Add(i.first, i.second)); - } - } - } - } + virtual void DoOnExecuteTxAfterRemoving(TBlobManagerDb& dbBlobs, const bool blobsWroteSuccessfully); + virtual void DoOnCompleteTxAfterRemoving(const bool blobsWroteSuccessfully); public: TDeclareRemovingAction(const TString& storageId, const TTabletId selfTabletId, const std::shared_ptr<NBlobOperations::TRemoveDeclareCounters>& counters, const std::shared_ptr<TGCInfo>& gcInfo) : TBase(storageId, selfTabletId, counters) diff --git a/ydb/core/tx/columnshard/blobs_action/tier/storage.h b/ydb/core/tx/columnshard/blobs_action/tier/storage.h index 0dc8c300529..cee81ff3720 100644 --- a/ydb/core/tx/columnshard/blobs_action/tier/storage.h +++ b/ydb/core/tx/columnshard/blobs_action/tier/storage.h @@ -49,6 +49,11 @@ public: virtual std::shared_ptr<IBlobInUseTracker> GetBlobsTracker() const override { return GCInfo; } + + virtual bool HasToDelete(const TUnifiedBlobId& blobId, const TTabletId tabletId) const override { + return GCInfo->HasToDelete(blobId, tabletId); + } + }; } diff --git a/ydb/core/tx/columnshard/blobs_action/transaction/tx_remove_blobs.cpp b/ydb/core/tx/columnshard/blobs_action/transaction/tx_remove_blobs.cpp index 74ad4c82b7e..da04c45bbc8 100644 --- a/ydb/core/tx/columnshard/blobs_action/transaction/tx_remove_blobs.cpp +++ b/ydb/core/tx/columnshard/blobs_action/transaction/tx_remove_blobs.cpp @@ -8,6 +8,8 @@ bool TTxRemoveSharedBlobs::Execute(TTransactionContext& txc, const TActorContext NActors::TLogContextGuard logGuard = NActors::TLogContextBuilder::Build(NKikimrServices::TX_COLUMNSHARD)("tablet_id", Self->TabletID())("tx_state", "execute"); NOlap::TBlobManagerDb blobManagerDb(txc.DB); RemoveAction->OnExecuteTxAfterRemoving(blobManagerDb, true); + + Manager->RemoveSharedBlobsDB(txc, SharingBlobIds); return true; } @@ -15,8 +17,12 @@ void TTxRemoveSharedBlobs::Complete(const TActorContext& ctx) { TMemoryProfileGuard mpg("TTxRemoveSharedBlobs::Complete"); NActors::TLogContextGuard logGuard = NActors::TLogContextBuilder::Build(NKikimrServices::TX_COLUMNSHARD)("tablet_id", Self->TabletID())("tx_state", "complete"); RemoveAction->OnCompleteTxAfterRemoving(true); + Manager->RemoveSharedBlobs(SharingBlobIds); + + ctx.Send(InitiatorActorId, new NOlap::NBlobOperations::NEvents::TEvDeleteSharedBlobsFinished((NOlap::TTabletId)Self->TabletID(), + NKikimrColumnShardBlobOperationsProto::TEvDeleteSharedBlobsFinished::Success)); - ctx.Send(InitiatorActorId, new NOlap::NBlobOperations::NEvents::TEvDeleteSharedBlobsFinished((NOlap::TTabletId)Self->TabletID())); + Self->GetStoragesManager()->GetSharedBlobsManager()->FinishExternalModification(); } } diff --git a/ydb/core/tx/columnshard/blobs_action/transaction/tx_remove_blobs.h b/ydb/core/tx/columnshard/blobs_action/transaction/tx_remove_blobs.h index d716ba89860..937174875fb 100644 --- a/ydb/core/tx/columnshard/blobs_action/transaction/tx_remove_blobs.h +++ b/ydb/core/tx/columnshard/blobs_action/transaction/tx_remove_blobs.h @@ -1,6 +1,8 @@ #pragma once #include <ydb/core/tx/columnshard/columnshard_impl.h> #include <ydb/core/tx/columnshard/engines/writer/indexed_blob_constructor.h> +#include <ydb/core/tx/columnshard/blobs_action/abstract/blob_set.h> +#include <ydb/core/tx/columnshard/data_sharing/manager/shared_blobs.h> namespace NKikimr::NColumnShard { @@ -9,7 +11,8 @@ private: std::shared_ptr<NOlap::IBlobsDeclareRemovingAction> RemoveAction; const ui32 TabletTxNo; const NActors::TActorId InitiatorActorId; - + NOlap::TTabletsByBlob SharingBlobIds; + const std::shared_ptr<NOlap::NDataSharing::TStorageSharedBlobsManager> Manager; TStringBuilder TxPrefix() const { return TStringBuilder() << "TxWrite[" << ToString(TabletTxNo) << "] "; } @@ -18,12 +21,24 @@ private: return TStringBuilder() << " at tablet " << Self->TabletID(); } public: - TTxRemoveSharedBlobs(TColumnShard* self, const std::shared_ptr<NOlap::IBlobsDeclareRemovingAction>& removeAction, const NActors::TActorId initiatorActorId) + TTxRemoveSharedBlobs(TColumnShard* self, + const NOlap::TTabletsByBlob& sharingBlobIds, const NActors::TActorId initiatorActorId, + const TString& storageId) : TBase(self) - , RemoveAction(removeAction) , TabletTxNo(++Self->TabletTxCounter) , InitiatorActorId(initiatorActorId) - {} + , SharingBlobIds(sharingBlobIds) + , Manager(Self->GetStoragesManager()->GetSharedBlobsManager()->GetStorageManagerVerified(storageId)) + { + Self->GetStoragesManager()->GetSharedBlobsManager()->StartExternalModification(); + RemoveAction = Self->GetStoragesManager()->GetOperatorVerified(storageId)->StartDeclareRemovingAction(NOlap::NBlobOperations::EConsumer::CLEANUP_SHARED_BLOBS); + auto categories = Manager->BuildRemoveCategories(SharingBlobIds); + for (auto it = categories.GetDirect().GetIterator(); it.IsValid(); ++it) { + RemoveAction->DeclareRemove(it.GetTabletId(), it.GetBlobId()); + } + AFL_VERIFY(categories.GetBorrowed().IsEmpty()); + AFL_VERIFY(categories.GetSharing().GetSize() == SharingBlobIds.GetSize()); + } bool Execute(TTransactionContext& txc, const TActorContext& ctx) override; void Complete(const TActorContext& ctx) override; diff --git a/ydb/core/tx/columnshard/columnshard_impl.cpp b/ydb/core/tx/columnshard/columnshard_impl.cpp index ecfa5d7450c..a5615460916 100644 --- a/ydb/core/tx/columnshard/columnshard_impl.cpp +++ b/ydb/core/tx/columnshard/columnshard_impl.cpp @@ -1070,15 +1070,23 @@ void TColumnShard::Handle(TAutoPtr<TEventHandle<NOlap::NBackground::TEvExecuteGe } void TColumnShard::Handle(NOlap::NBlobOperations::NEvents::TEvDeleteSharedBlobs::TPtr& ev, const TActorContext& ctx) { + if (SharingSessionsManager->IsSharingInProgress()) { + ctx.Send(NActors::ActorIdFromProto(ev->Get()->Record.GetSourceActorId()), + new NOlap::NBlobOperations::NEvents::TEvDeleteSharedBlobsFinished((NOlap::TTabletId)TabletID(), + NKikimrColumnShardBlobOperationsProto::TEvDeleteSharedBlobsFinished::DestinationCurrenlyLocked)); + AFL_DEBUG(NKikimrServices::TX_COLUMNSHARD)("event", "sharing_in_progress"); + return; + } + AFL_NOTICE(NKikimrServices::TX_COLUMNSHARD)("process", "BlobsSharing")("event", "TEvDeleteSharedBlobs"); NActors::TLogContextGuard gLogging = NActors::TLogContextBuilder::Build(NKikimrServices::TX_COLUMNSHARD)("tablet_id", TabletID())("event", "TEvDeleteSharedBlobs"); - auto removeAction = StoragesManager->GetOperator(ev->Get()->Record.GetStorageId())->StartDeclareRemovingAction(NOlap::NBlobOperations::EConsumer::CLEANUP_SHARED_BLOBS); + NOlap::TTabletsByBlob blobs; for (auto&& i : ev->Get()->Record.GetBlobIds()) { auto blobId = NOlap::TUnifiedBlobId::BuildFromString(i, nullptr); AFL_VERIFY(!!blobId)("problem", blobId.GetErrorMessage()); - removeAction->DeclareRemove((NOlap::TTabletId)ev->Get()->Record.GetSourceTabletId(), *blobId); + AFL_VERIFY(blobs.Add((NOlap::TTabletId)ev->Get()->Record.GetSourceTabletId(), blobId.DetachResult())); } - Execute(new TTxRemoveSharedBlobs(this, removeAction, NActors::ActorIdFromProto(ev->Get()->Record.GetSourceActorId())), ctx); + Execute(new TTxRemoveSharedBlobs(this, blobs, NActors::ActorIdFromProto(ev->Get()->Record.GetSourceActorId()), ev->Get()->Record.GetStorageId()), ctx); } void TColumnShard::Handle(NMetadata::NProvider::TEvRefreshSubscriberData::TPtr& ev) { diff --git a/ydb/core/tx/columnshard/data_locks/locks/abstract.h b/ydb/core/tx/columnshard/data_locks/locks/abstract.h index 2e301d62a51..826d363ad85 100644 --- a/ydb/core/tx/columnshard/data_locks/locks/abstract.h +++ b/ydb/core/tx/columnshard/data_locks/locks/abstract.h @@ -2,6 +2,7 @@ #include <ydb/library/accessor/accessor.h> #include <util/generic/string.h> +#include <util/generic/hash_set.h> #include <optional> #include <memory> @@ -19,8 +20,8 @@ private: YDB_READONLY_DEF(TString, LockName); YDB_READONLY_FLAG(ReadOnly, false); protected: - virtual std::optional<TString> DoIsLocked(const TPortionInfo& portion) const = 0; - virtual std::optional<TString> DoIsLocked(const TGranuleMeta& granule) const = 0; + virtual std::optional<TString> DoIsLocked(const TPortionInfo& portion, const THashSet<TString>& excludedLocks = {}) const = 0; + virtual std::optional<TString> DoIsLocked(const TGranuleMeta& granule, const THashSet<TString>& excludedLocks = {}) const = 0; virtual bool DoIsEmpty() const = 0; public: ILock(const TString& lockName, const bool isReadOnly = false) @@ -32,17 +33,17 @@ public: virtual ~ILock() = default; - std::optional<TString> IsLocked(const TPortionInfo& portion, const bool readOnly = false) const { + std::optional<TString> IsLocked(const TPortionInfo& portion, const THashSet<TString>& excludedLocks = {}, const bool readOnly = false) const { if (IsReadOnly() && readOnly) { return {}; } - return DoIsLocked(portion); + return DoIsLocked(portion, excludedLocks); } - std::optional<TString> IsLocked(const TGranuleMeta& g, const bool readOnly = false) const { + std::optional<TString> IsLocked(const TGranuleMeta& g, const THashSet<TString>& excludedLocks = {}, const bool readOnly = false) const { if (IsReadOnly() && readOnly) { return {}; } - return DoIsLocked(g); + return DoIsLocked(g, excludedLocks); } bool IsEmpty() const { return DoIsEmpty(); diff --git a/ydb/core/tx/columnshard/data_locks/locks/composite.h b/ydb/core/tx/columnshard/data_locks/locks/composite.h index fd23ee15f85..3c57da6047b 100644 --- a/ydb/core/tx/columnshard/data_locks/locks/composite.h +++ b/ydb/core/tx/columnshard/data_locks/locks/composite.h @@ -8,16 +8,22 @@ private: using TBase = ILock; std::vector<std::shared_ptr<ILock>> Locks; protected: - virtual std::optional<TString> DoIsLocked(const TPortionInfo& portion) const override { + virtual std::optional<TString> DoIsLocked(const TPortionInfo& portion, const THashSet<TString>& excludedLocks) const override { for (auto&& i : Locks) { + if (excludedLocks.contains(i->GetLockName())) { + continue; + } if (auto lockName = i->IsLocked(portion)) { return lockName; } } return {}; } - virtual std::optional<TString> DoIsLocked(const TGranuleMeta& granule) const override { + virtual std::optional<TString> DoIsLocked(const TGranuleMeta& granule, const THashSet<TString>& excludedLocks) const override { for (auto&& i : Locks) { + if (excludedLocks.contains(i->GetLockName())) { + continue; + } if (auto lockName = i->IsLocked(granule)) { return lockName; } diff --git a/ydb/core/tx/columnshard/data_locks/locks/list.h b/ydb/core/tx/columnshard/data_locks/locks/list.h index e74251e30b4..512386e985b 100644 --- a/ydb/core/tx/columnshard/data_locks/locks/list.h +++ b/ydb/core/tx/columnshard/data_locks/locks/list.h @@ -11,13 +11,13 @@ private: THashSet<TPortionAddress> Portions; THashSet<ui64> Granules; protected: - virtual std::optional<TString> DoIsLocked(const TPortionInfo& portion) const override { + virtual std::optional<TString> DoIsLocked(const TPortionInfo& portion, const THashSet<TString>& /*excludedLocks*/) const override { if (Portions.contains(portion.GetAddress())) { return GetLockName(); } return {}; } - virtual std::optional<TString> DoIsLocked(const TGranuleMeta& granule) const override { + virtual std::optional<TString> DoIsLocked(const TGranuleMeta& granule, const THashSet<TString>& /*excludedLocks*/) const override { if (Granules.contains(granule.GetPathId())) { return GetLockName(); } @@ -70,13 +70,13 @@ private: using TBase = ILock; THashSet<ui64> Tables; protected: - virtual std::optional<TString> DoIsLocked(const TPortionInfo& portion) const override { + virtual std::optional<TString> DoIsLocked(const TPortionInfo& portion, const THashSet<TString>& /*excludedLocks*/) const override { if (Tables.contains(portion.GetPathId())) { return GetLockName(); } return {}; } - virtual std::optional<TString> DoIsLocked(const TGranuleMeta& granule) const override { + virtual std::optional<TString> DoIsLocked(const TGranuleMeta& granule, const THashSet<TString>& /*excludedLocks*/) const override { if (Tables.contains(granule.GetPathId())) { return GetLockName(); } diff --git a/ydb/core/tx/columnshard/data_locks/locks/snapshot.h b/ydb/core/tx/columnshard/data_locks/locks/snapshot.h index c1f6e10b06e..78edc72599f 100644 --- a/ydb/core/tx/columnshard/data_locks/locks/snapshot.h +++ b/ydb/core/tx/columnshard/data_locks/locks/snapshot.h @@ -9,27 +9,30 @@ class TSnapshotLock: public ILock { private: using TBase = ILock; const TSnapshot SnapshotBarrier; - const THashSet<TTabletId> PathIds; + const THashSet<ui64> PathIds; protected: - virtual std::optional<TString> DoIsLocked(const TPortionInfo& portion) const override { - if (PathIds.contains((TTabletId)portion.GetPathId()) && portion.RecordSnapshotMin() <= SnapshotBarrier) { + virtual std::optional<TString> DoIsLocked(const TPortionInfo& portion, const THashSet<TString>& /*excludedLocks*/) const override { + if (PathIds.contains(portion.GetPathId()) && portion.RecordSnapshotMin() <= SnapshotBarrier) { return GetLockName(); } return {}; } - virtual std::optional<TString> DoIsLocked(const TGranuleMeta& granule) const override { - if (PathIds.contains((TTabletId)granule.GetPathId())) { + virtual bool DoIsEmpty() const override { + return PathIds.empty(); + } + virtual std::optional<TString> DoIsLocked(const TGranuleMeta& granule, const THashSet<TString>& /*excludedLocks*/) const override { + if (PathIds.contains(granule.GetPathId())) { return GetLockName(); } return {}; } public: - TSnapshotLock(const TString& lockName, const TSnapshot& snapshotBarrier, const THashSet<TTabletId>& pathIds, const bool readOnly = false) + TSnapshotLock(const TString& lockName, const TSnapshot& snapshotBarrier, const THashSet<ui64>& pathIds, const bool readOnly = false) : TBase(lockName, readOnly) , SnapshotBarrier(snapshotBarrier) , PathIds(pathIds) { - + AFL_VERIFY(SnapshotBarrier.Valid()); } }; diff --git a/ydb/core/tx/columnshard/data_locks/manager/manager.cpp b/ydb/core/tx/columnshard/data_locks/manager/manager.cpp index 53cb884fefb..8de6300b85a 100644 --- a/ydb/core/tx/columnshard/data_locks/manager/manager.cpp +++ b/ydb/core/tx/columnshard/data_locks/manager/manager.cpp @@ -15,18 +15,24 @@ void TManager::UnregisterLock(const TString& processId) { AFL_VERIFY(ProcessLocks.erase(processId))("process_id", processId); } -std::optional<TString> TManager::IsLocked(const TPortionInfo& portion) const { +std::optional<TString> TManager::IsLocked(const TPortionInfo& portion, const THashSet<TString>& excludedLocks) const { for (auto&& i : ProcessLocks) { - if (auto lockName = i.second->IsLocked(portion)) { + if (excludedLocks.contains(i.first)) { + continue; + } + if (auto lockName = i.second->IsLocked(portion, excludedLocks)) { return lockName; } } return {}; } -std::optional<TString> TManager::IsLocked(const TGranuleMeta& granule) const { +std::optional<TString> TManager::IsLocked(const TGranuleMeta& granule, const THashSet<TString>& excludedLocks) const { for (auto&& i : ProcessLocks) { - if (auto lockName = i.second->IsLocked(granule)) { + if (excludedLocks.contains(i.first)) { + continue; + } + if (auto lockName = i.second->IsLocked(granule, excludedLocks)) { return lockName; } } diff --git a/ydb/core/tx/columnshard/data_locks/manager/manager.h b/ydb/core/tx/columnshard/data_locks/manager/manager.h index ee0e9ea31ea..b59a0bdb1e8 100644 --- a/ydb/core/tx/columnshard/data_locks/manager/manager.h +++ b/ydb/core/tx/columnshard/data_locks/manager/manager.h @@ -41,8 +41,8 @@ public: [[nodiscard]] std::shared_ptr<TGuard> RegisterLock(Args&&... args) { return RegisterLock(std::make_shared<TLock>(args...)); } - std::optional<TString> IsLocked(const TPortionInfo& portion) const; - std::optional<TString> IsLocked(const TGranuleMeta& granule) const; + std::optional<TString> IsLocked(const TPortionInfo& portion, const THashSet<TString>& excludedLocks = {}) const; + std::optional<TString> IsLocked(const TGranuleMeta& granule, const THashSet<TString>& excludedLocks = {}) const; }; diff --git a/ydb/core/tx/columnshard/data_sharing/common/session/common.cpp b/ydb/core/tx/columnshard/data_sharing/common/session/common.cpp index bb22359d925..002f2243096 100644 --- a/ydb/core/tx/columnshard/data_sharing/common/session/common.cpp +++ b/ydb/core/tx/columnshard/data_sharing/common/session/common.cpp @@ -1,7 +1,7 @@ #include "common.h" #include <ydb/core/tx/columnshard/columnshard_impl.h> -#include <ydb/core/tx/columnshard/data_locks/locks/list.h> +#include <ydb/core/tx/columnshard/data_locks/locks/snapshot.h> #include <ydb/core/tx/columnshard/engines/column_engine_logs.h> #include <util/string/builder.h> @@ -12,55 +12,51 @@ TString TCommonSession::DebugString() const { return TStringBuilder() << "{id=" << SessionId << ";context=" << TransferContext.DebugString() << ";}"; } -bool TCommonSession::Start(const NColumnShard::TColumnShard& shard) { +bool TCommonSession::TryStart(const NColumnShard::TColumnShard& shard) { const NActors::TLogContextGuard lGuard = NActors::TLogContextBuilder::Build()("info", Info); AFL_DEBUG(NKikimrServices::TX_COLUMNSHARD)("info", "Start"); - AFL_VERIFY(!IsStartingFlag); - IsStartingFlag = true; - AFL_VERIFY(!IsStartedFlag); + AFL_VERIFY(State == EState::Prepared); + + AFL_VERIFY(!!LockGuard); const auto& index = shard.GetIndexAs<TColumnEngineForLogs>(); THashMap<ui64, std::vector<std::shared_ptr<TPortionInfo>>> portionsByPath; - std::vector<std::shared_ptr<TPortionInfo>> portionsLock; - THashMap<TString, THashSet<TUnifiedBlobId>> local; + THashSet<TString> StoragesIds; for (auto&& i : GetPathIdsForStart()) { -// const auto insertTableSnapshot = shard.GetInsertTable().GetMinCommittedSnapshot(i); -// const auto shardSnapshot = shard.GetLastCompletedTx(); -// if (shard.GetInsertTable().GetMinCommittedSnapshot(i).value_or(shard.GetLastPlannedSnapshot()) <= GetSnapshotBarrier()) { -// AFL_WARN(NKikimrServices::TX_COLUMNSHARD)("insert_table_snapshot", insertTableSnapshot)("last_completed_tx", shardSnapshot)("barrier", GetSnapshotBarrier()); -// IsStartingFlag = false; -// return false; -// } auto& portionsVector = portionsByPath[i]; const auto& g = index.GetGranuleVerified(i); for (auto&& p : g.GetPortionsOlderThenSnapshot(GetSnapshotBarrier())) { - if (shard.GetDataLocksManager()->IsLocked(*p.second)) { - IsStartingFlag = false; + if (shard.GetDataLocksManager()->IsLocked(*p.second, { "sharing_session:" + GetSessionId() })) { return false; } portionsVector.emplace_back(p.second); - portionsLock.emplace_back(p.second); } } - IsStartedFlag = DoStart(shard, portionsByPath); - if (IsFinishedFlag) { - IsStartedFlag = false; - } - if (IsStartedFlag) { - AFL_VERIFY(!LockGuard); - LockGuard = shard.GetDataLocksManager()->RegisterLock<NDataLocks::TListPortionsLock>("sharing_session:" + GetSessionId(), portionsLock, true); + if (shard.GetStoragesManager()->GetSharedBlobsManager()->HasExternalModifications()) { + return false; } - IsStartingFlag = false; - return IsStartedFlag; + + AFL_VERIFY(DoStart(shard, portionsByPath)); + State = EState::InProgress; + return true; } -void TCommonSession::Finish(const std::shared_ptr<NDataLocks::TManager>& dataLocksManager) { - AFL_VERIFY(!IsFinishedFlag); - IsFinishedFlag = true; - if (IsStartedFlag) { - AFL_VERIFY(LockGuard); - LockGuard->Release(*dataLocksManager); - } +void TCommonSession::PrepareToStart(const NColumnShard::TColumnShard& shard) { + const NActors::TLogContextGuard lGuard = NActors::TLogContextBuilder::Build()("info", Info); + AFL_VERIFY(State == EState::Created); + State = EState::Prepared; + AFL_VERIFY(!LockGuard); + LockGuard = shard.GetDataLocksManager()->RegisterLock<NDataLocks::TSnapshotLock>("sharing_session:" + GetSessionId(), + TransferContext.GetSnapshotBarrierVerified(), GetPathIdsForStart(), true); + shard.GetSharingSessionsManager()->StartSharingSession(); +} + +void TCommonSession::Finish(const NColumnShard::TColumnShard& shard, const std::shared_ptr<NDataLocks::TManager>& dataLocksManager) { + AFL_VERIFY(State == EState::InProgress); + State = EState::Finished; + shard.GetSharingSessionsManager()->FinishSharingSession(); + AFL_VERIFY(LockGuard); + LockGuard->Release(*dataLocksManager); } }
\ No newline at end of file diff --git a/ydb/core/tx/columnshard/data_sharing/common/session/common.h b/ydb/core/tx/columnshard/data_sharing/common/session/common.h index d2e06fa9889..717636bba28 100644 --- a/ydb/core/tx/columnshard/data_sharing/common/session/common.h +++ b/ydb/core/tx/columnshard/data_sharing/common/session/common.h @@ -22,6 +22,13 @@ namespace NKikimr::NOlap::NDataSharing { class TCommonSession { private: + enum class EState { + Created, + Prepared, + InProgress, + Finished + }; + static ui64 GetNextRuntimeId() { static TAtomicCounter Counter = 0; return (ui64)Counter.Inc(); @@ -31,9 +38,7 @@ private: const TString Info; YDB_READONLY(ui64, RuntimeId, GetNextRuntimeId()); std::shared_ptr<NDataLocks::TManager::TGuard> LockGuard; - bool IsStartedFlag = false; - bool IsStartingFlag = false; - bool IsFinishedFlag = false; + EState State = EState::Created; protected: TTransferContext TransferContext; virtual bool DoStart(const NColumnShard::TColumnShard& shard, const THashMap<ui64, std::vector<std::shared_ptr<TPortionInfo>>>& portions) = 0; @@ -57,24 +62,25 @@ public: return TransferContext; } - bool IsFinished() const { - return IsFinishedFlag; + bool IsReadyForStarting() const { + return State == EState::Created; } - bool IsStarted() const { - return IsStartedFlag; + bool IsPrepared() const { + return State == EState::Prepared; } - bool IsStarting() const { - return IsStartingFlag; + bool IsInProgress() const { + return State == EState::InProgress; } bool IsEqualTo(const TCommonSession& item) const { return SessionId == item.SessionId && TransferContext.IsEqualTo(item.TransferContext); } - bool Start(const NColumnShard::TColumnShard& shard); - void Finish(const std::shared_ptr<NDataLocks::TManager>& dataLocksManager); + void PrepareToStart(const NColumnShard::TColumnShard& shard); + bool TryStart(const NColumnShard::TColumnShard& shard); + void Finish(const NColumnShard::TColumnShard& shard, const std::shared_ptr<NDataLocks::TManager>& dataLocksManager); const TSnapshot& GetSnapshotBarrier() const { return TransferContext.GetSnapshotBarrierVerified(); diff --git a/ydb/core/tx/columnshard/data_sharing/destination/events/transfer.cpp b/ydb/core/tx/columnshard/data_sharing/destination/events/transfer.cpp index 55d7e8e9ea7..cd90a735332 100644 --- a/ydb/core/tx/columnshard/data_sharing/destination/events/transfer.cpp +++ b/ydb/core/tx/columnshard/data_sharing/destination/events/transfer.cpp @@ -6,18 +6,21 @@ namespace NKikimr::NOlap::NDataSharing::NEvents { THashMap<NKikimr::NOlap::TTabletId, NKikimr::NOlap::NDataSharing::TTaskForTablet> TPathIdData::BuildLinkTabletTasks( - const std::shared_ptr<TSharedBlobsManager>& sharedBlobs, const TTabletId selfTabletId, const TTransferContext& context, const TVersionedIndex& index) { + const std::shared_ptr<IStoragesManager>& storages, const TTabletId selfTabletId, const TTransferContext& context, const TVersionedIndex& index) { THashMap<TString, THashSet<TUnifiedBlobId>> blobIds; for (auto&& i : Portions) { auto schema = i.GetSchema(index); i.FillBlobIdsByStorage(blobIds, schema->GetIndexInfo()); } + const std::shared_ptr<TSharedBlobsManager> sharedBlobs = storages->GetSharedBlobsManager(); + THashMap<TString, THashMap<TUnifiedBlobId, TBlobSharing>> blobsInfo; for (auto&& i : blobIds) { - auto storageManager = sharedBlobs->GetStorageManagerVerified(i.first); - auto storeCategories = storageManager->BuildStoreCategories(i.second); + auto sharingManager = sharedBlobs->GetStorageManagerVerified(i.first); + auto storageManager = storages->GetOperatorVerified(i.first); + auto storeCategories = sharingManager->BuildStoreCategories(i.second); auto& blobs = blobsInfo[i.first]; for (auto it = storeCategories.GetDirect().GetIterator(); it.IsValid(); ++it) { auto itSharing = blobs.find(it.GetBlobId()); @@ -30,6 +33,9 @@ THashMap<NKikimr::NOlap::TTabletId, NKikimr::NOlap::NDataSharing::TTaskForTablet if (itSharing == blobs.end()) { itSharing = blobs.emplace(it.GetBlobId(), TBlobSharing(i.first, it.GetBlobId())).first; } + if (storageManager->HasToDelete(it.GetBlobId(), it.GetTabletId())) { + continue; + } itSharing->second.AddShared(it.GetTabletId()); } for (auto it = storeCategories.GetBorrowed().GetIterator(); it.IsValid(); ++it) { diff --git a/ydb/core/tx/columnshard/data_sharing/destination/events/transfer.h b/ydb/core/tx/columnshard/data_sharing/destination/events/transfer.h index f69e08334d8..715cee95fb1 100644 --- a/ydb/core/tx/columnshard/data_sharing/destination/events/transfer.h +++ b/ydb/core/tx/columnshard/data_sharing/destination/events/transfer.h @@ -49,7 +49,7 @@ public: std::vector<TPortionInfo> DetachPortions() { return std::move(Portions); } - THashMap<TTabletId, TTaskForTablet> BuildLinkTabletTasks(const std::shared_ptr<TSharedBlobsManager>& sharedBlobs, const TTabletId selfTabletId, + THashMap<TTabletId, TTaskForTablet> BuildLinkTabletTasks(const std::shared_ptr<IStoragesManager>& storages, const TTabletId selfTabletId, const TTransferContext& context, const TVersionedIndex& index); void InitPortionIds(ui64* lastPortionId, const std::optional<ui64> pathId = {}) { diff --git a/ydb/core/tx/columnshard/data_sharing/destination/session/destination.cpp b/ydb/core/tx/columnshard/data_sharing/destination/session/destination.cpp index b6601280346..ed82c91a84e 100644 --- a/ydb/core/tx/columnshard/data_sharing/destination/session/destination.cpp +++ b/ydb/core/tx/columnshard/data_sharing/destination/session/destination.cpp @@ -25,7 +25,7 @@ NKikimr::TConclusionStatus TDestinationSession::DataReceived(THashMap<ui64, NEve } ui32 TDestinationSession::GetSourcesInProgressCount() const { - AFL_VERIFY(IsStarted() || IsStarting()); + AFL_VERIFY(IsInProgress()); AFL_VERIFY(Cursors.size()); ui32 result = 0; for (auto&& [_, cursor] : Cursors) { @@ -37,7 +37,7 @@ ui32 TDestinationSession::GetSourcesInProgressCount() const { } void TDestinationSession::SendCurrentCursorAck(const NColumnShard::TColumnShard& shard, const std::optional<TTabletId> tabletId) { - AFL_VERIFY(IsStarted() || IsStarting()); + AFL_VERIFY(IsInProgress() || IsPrepared()); bool found = false; for (auto&& [_, cursor] : Cursors) { if (tabletId && *tabletId != cursor.GetTabletId()) { diff --git a/ydb/core/tx/columnshard/data_sharing/destination/transactions/tx_finish_from_source.cpp b/ydb/core/tx/columnshard/data_sharing/destination/transactions/tx_finish_from_source.cpp index f09641f6696..a6a180963bf 100644 --- a/ydb/core/tx/columnshard/data_sharing/destination/transactions/tx_finish_from_source.cpp +++ b/ydb/core/tx/columnshard/data_sharing/destination/transactions/tx_finish_from_source.cpp @@ -26,7 +26,7 @@ void TTxFinishFromSource::DoComplete(const TActorContext& ctx) { Self->GetProgressTxController().FinishProposeOnComplete(*Session->GetTransferContext().GetTxId(), ctx); } NYDBTest::TControllers::GetColumnShardController()->OnDataSharingFinished(Self->TabletID(), Session->GetSessionId()); - Session->Finish(Self->GetDataLocksManager()); + Session->Finish(*Self, Self->GetDataLocksManager()); Session->GetInitiatorController().Finished(Session->GetSessionId()); } } diff --git a/ydb/core/tx/columnshard/data_sharing/destination/transactions/tx_start_from_initiator.cpp b/ydb/core/tx/columnshard/data_sharing/destination/transactions/tx_start_from_initiator.cpp index c07a010e6fc..e773f5320d5 100644 --- a/ydb/core/tx/columnshard/data_sharing/destination/transactions/tx_start_from_initiator.cpp +++ b/ydb/core/tx/columnshard/data_sharing/destination/transactions/tx_start_from_initiator.cpp @@ -27,7 +27,6 @@ bool TTxConfirmFromInitiator::DoExecute(NTabletFlatExecutor::TTransactionContext } void TTxConfirmFromInitiator::DoComplete(const TActorContext& /*ctx*/) { - Session->Start(*Self); Session->GetInitiatorController().ConfirmSuccess(Session->GetSessionId()); } diff --git a/ydb/core/tx/columnshard/data_sharing/manager/sessions.cpp b/ydb/core/tx/columnshard/data_sharing/manager/sessions.cpp index 8749354e13f..18a30ac7606 100644 --- a/ydb/core/tx/columnshard/data_sharing/manager/sessions.cpp +++ b/ydb/core/tx/columnshard/data_sharing/manager/sessions.cpp @@ -10,15 +10,26 @@ namespace NKikimr::NOlap::NDataSharing { void TSessionsManager::Start(const NColumnShard::TColumnShard& shard) const { NActors::TLogContextGuard logGuard = NActors::TLogContextBuilder::Build()("sessions", "start")("tablet_id", shard.TabletID()); for (auto&& i : SourceSessions) { - if (!i.second->IsStarted()) { - i.second->Start(shard); + if (i.second->IsReadyForStarting()) { + i.second->PrepareToStart(shard); } } for (auto&& i : DestSessions) { - if (!i.second->IsStarted() && i.second->IsConfirmed()) { - i.second->Start(shard); + if (i.second->IsReadyForStarting() && i.second->IsConfirmed()) { + i.second->PrepareToStart(shard); + } + } + + for (auto&& i : SourceSessions) { + if (i.second->IsPrepared()) { + i.second->TryStart(shard); + } + } + for (auto&& i : DestSessions) { + if (i.second->IsPrepared() && i.second->IsConfirmed()) { + i.second->TryStart(shard); if (!i.second->GetSourcesInProgressCount()) { - i.second->Finish(shard.GetDataLocksManager()); + i.second->Finish(shard, shard.GetDataLocksManager()); } } } @@ -31,7 +42,7 @@ void TSessionsManager::InitializeEventsExchange(const NColumnShard::TColumnShard if (sessionCookie && *sessionCookie != i.second->GetRuntimeId()) { continue; } - i.second->ActualizeDestination(shard.GetDataLocksManager()); + i.second->ActualizeDestination(shard, shard.GetDataLocksManager()); } for (auto&& i : DestSessions) { if (sessionCookie && *sessionCookie != i.second->GetRuntimeId()) { diff --git a/ydb/core/tx/columnshard/data_sharing/manager/sessions.h b/ydb/core/tx/columnshard/data_sharing/manager/sessions.h index 691b42ad8bb..a2e5efdaa60 100644 --- a/ydb/core/tx/columnshard/data_sharing/manager/sessions.h +++ b/ydb/core/tx/columnshard/data_sharing/manager/sessions.h @@ -12,9 +12,22 @@ class TSessionsManager { private: THashMap<TString, std::shared_ptr<TSourceSession>> SourceSessions; THashMap<TString, std::shared_ptr<TDestinationSession>> DestSessions; + TAtomicCounter SharingSessions; public: TSessionsManager() = default; + void StartSharingSession() { + SharingSessions.Inc(); + } + + void FinishSharingSession() { + AFL_VERIFY(SharingSessions.Dec() >= 0); + } + + bool IsSharingInProgress() const { + return SharingSessions.Val(); + } + void Start(const NColumnShard::TColumnShard& shard) const; std::shared_ptr<TSourceSession> GetSourceSession(const TString& sessionId) const { diff --git a/ydb/core/tx/columnshard/data_sharing/manager/shared_blobs.cpp b/ydb/core/tx/columnshard/data_sharing/manager/shared_blobs.cpp index 40387937a93..5188af5478f 100644 --- a/ydb/core/tx/columnshard/data_sharing/manager/shared_blobs.cpp +++ b/ydb/core/tx/columnshard/data_sharing/manager/shared_blobs.cpp @@ -87,16 +87,32 @@ void TStorageSharedBlobsManager::CASBorrowedBlobsDB(NTabletFlatExecutor::TTransa NIceDb::TNiceDb db(txc.DB); for (auto&& i : blobIds) { auto it = BorrowedBlobIds.find(i); - AFL_VERIFY(it != BorrowedBlobIds.end())("blob_id", i.ToStringNew()); - AFL_VERIFY(it->second == tabletIdFrom || it->second == tabletIdTo); if (tabletIdTo == SelfTabletId) { + AFL_VERIFY(it == BorrowedBlobIds.end() || it->second == tabletIdFrom); db.Table<NColumnShard::Schema::BorrowedBlobIds>().Key(StorageId, i.ToStringNew()).Delete(); } else { + AFL_VERIFY(it != BorrowedBlobIds.end())("blob_id", i.ToStringNew()); db.Table<NColumnShard::Schema::BorrowedBlobIds>().Key(StorageId, i.ToStringNew()).Update(NIceDb::TUpdate<NColumnShard::Schema::BorrowedBlobIds::TabletId>((ui64)tabletIdTo)); } } } +void TStorageSharedBlobsManager::CASBorrowedBlobs(const TTabletId tabletIdFrom, const TTabletId tabletIdTo, const THashSet<TUnifiedBlobId>& blobIds) { + for (auto&& i : blobIds) { + auto it = BorrowedBlobIds.find(i); + if (tabletIdTo == SelfTabletId) { + AFL_VERIFY(it == BorrowedBlobIds.end() || it->second == tabletIdFrom); + if (it != BorrowedBlobIds.end()) { + BorrowedBlobIds.erase(it); + } + } else { + AFL_VERIFY(it != BorrowedBlobIds.end()); + AFL_VERIFY(it->second == tabletIdFrom || it->second == tabletIdTo); + it->second = tabletIdTo; + } + } +} + void TStorageSharedBlobsManager::OnTransactionExecuteAfterCleaning(const TBlobsCategories& removeTask, NTable::TDatabase& db) { TBlobManagerDb dbBlobs(db); for (auto&& i : removeTask.GetSharing()) { diff --git a/ydb/core/tx/columnshard/data_sharing/manager/shared_blobs.h b/ydb/core/tx/columnshard/data_sharing/manager/shared_blobs.h index f04f094165d..ec91655327c 100644 --- a/ydb/core/tx/columnshard/data_sharing/manager/shared_blobs.h +++ b/ydb/core/tx/columnshard/data_sharing/manager/shared_blobs.h @@ -66,7 +66,7 @@ public: return result; } - TBlobsCategories BuildRemoveCategories(TTabletsByBlob&& blobs) const { + TBlobsCategories BuildRemoveCategories(const TTabletsByBlob& blobs) const { TBlobsCategories result(SelfTabletId); for (auto it = blobs.GetIterator(); it.IsValid(); ++it) { CheckRemoveBlobId(it.GetTabletId(), it.GetBlobId(), result); @@ -124,6 +124,9 @@ public: void AddBorrowedBlobs(const TTabletByBlob& blobIds) { for (auto&& i : blobIds) { + if (i.second == SelfTabletId) { + continue; + } auto infoInsert = BorrowedBlobIds.emplace(i.first, i.second); if (!infoInsert.second) { AFL_VERIFY(infoInsert.first->second == i.second)("before", infoInsert.first->second)("after", i.second); @@ -133,24 +136,14 @@ public: void CASBorrowedBlobsDB(NTabletFlatExecutor::TTransactionContext& txc, const TTabletId tabletIdFrom, const TTabletId tabletIdTo, const THashSet<TUnifiedBlobId>& blobIds); - void CASBorrowedBlobs(const TTabletId tabletIdFrom, const TTabletId tabletIdTo, const THashSet<TUnifiedBlobId>& blobIds) { - for (auto&& i : blobIds) { - auto it = BorrowedBlobIds.find(i); - AFL_VERIFY(it != BorrowedBlobIds.end()); - AFL_VERIFY(it->second == tabletIdFrom || it->second == tabletIdTo); - if (it->second == SelfTabletId) { - BorrowedBlobIds.erase(it); - } else { - it->second = tabletIdTo; - } - } - } + void CASBorrowedBlobs(const TTabletId tabletIdFrom, const TTabletId tabletIdTo, const THashSet<TUnifiedBlobId>& blobIds); [[nodiscard]] bool UpsertSharedBlobOnLoad(const TUnifiedBlobId& blobId, const TTabletId tabletId) { return SharedBlobIds.Add(tabletId, blobId); } [[nodiscard]] bool UpsertBorrowedBlobOnLoad(const TUnifiedBlobId& blobId, const TTabletId ownerTabletId) { + AFL_VERIFY(ownerTabletId != SelfTabletId); return BorrowedBlobIds.emplace(blobId, ownerTabletId).second; } @@ -167,6 +160,7 @@ class TSharedBlobsManager { private: const TTabletId SelfTabletId; THashMap<TString, std::shared_ptr<TStorageSharedBlobsManager>> Storages; + TAtomicCounter ExternalModificationsCount; public: TSharedBlobsManager(const TTabletId tabletId) : SelfTabletId(tabletId) @@ -174,6 +168,18 @@ public: } + void StartExternalModification() { + ExternalModificationsCount.Inc(); + } + + void FinishExternalModification() { + AFL_VERIFY(ExternalModificationsCount.Dec() >= 0); + } + + bool HasExternalModifications() const { + return ExternalModificationsCount.Val(); + } + bool IsTrivialLinks() const { for (auto&& i : Storages) { if (!i.second->IsTrivialLinks()) { diff --git a/ydb/core/tx/columnshard/data_sharing/source/session/cursor.cpp b/ydb/core/tx/columnshard/data_sharing/source/session/cursor.cpp index 1ca8fac2ace..1072d6ff1cb 100644 --- a/ydb/core/tx/columnshard/data_sharing/source/session/cursor.cpp +++ b/ydb/core/tx/columnshard/data_sharing/source/session/cursor.cpp @@ -5,7 +5,7 @@ namespace NKikimr::NOlap::NDataSharing { -void TSourceCursor::BuildSelection(const std::shared_ptr<TSharedBlobsManager>& sharedBlobsManager, const TVersionedIndex& index) { +void TSourceCursor::BuildSelection(const std::shared_ptr<IStoragesManager>& storagesManager, const TVersionedIndex& index) { THashMap<ui64, NEvents::TPathIdData> result; auto itCurrentPath = PortionsForSend.find(StartPathId); AFL_VERIFY(itCurrentPath != PortionsForSend.end()); @@ -41,7 +41,7 @@ void TSourceCursor::BuildSelection(const std::shared_ptr<TSharedBlobsManager>& s THashMap<TTabletId, TTaskForTablet> tabletTasksResult; for (auto&& i : result) { - THashMap<TTabletId, TTaskForTablet> tabletTasks = i.second.BuildLinkTabletTasks(sharedBlobsManager, SelfTabletId, TransferContext, index); + THashMap<TTabletId, TTaskForTablet> tabletTasks = i.second.BuildLinkTabletTasks(storagesManager, SelfTabletId, TransferContext, index); for (auto&& t : tabletTasks) { auto it = tabletTasksResult.find(t.first); if (it == tabletTasksResult.end()) { @@ -56,7 +56,7 @@ void TSourceCursor::BuildSelection(const std::shared_ptr<TSharedBlobsManager>& s std::swap(Selected, result); } -bool TSourceCursor::Next(const std::shared_ptr<TSharedBlobsManager>& sharedBlobsManager, const TVersionedIndex& index) { +bool TSourceCursor::Next(const std::shared_ptr<IStoragesManager>& storagesManager, const TVersionedIndex& index) { PreviousSelected = std::move(Selected); LinksModifiedTablets.clear(); Selected.clear(); @@ -71,7 +71,7 @@ bool TSourceCursor::Next(const std::shared_ptr<TSharedBlobsManager>& sharedBlobs NextPathId = {}; NextPortionId = {}; ++PackIdx; - BuildSelection(sharedBlobsManager, index); + BuildSelection(storagesManager, index); AFL_VERIFY(IsValid()); return true; } @@ -141,7 +141,7 @@ TSourceCursor::TSourceCursor(const TTabletId selfTabletId, const std::set<ui64>& { } -bool TSourceCursor::Start(const std::shared_ptr<TSharedBlobsManager>& sharedBlobsManager, const THashMap<ui64, std::vector<std::shared_ptr<TPortionInfo>>>& portions, const TVersionedIndex& index) { +bool TSourceCursor::Start(const std::shared_ptr<IStoragesManager>& storagesManager, const THashMap<ui64, std::vector<std::shared_ptr<TPortionInfo>>>& portions, const TVersionedIndex& index) { AFL_VERIFY(!IsStartedFlag); std::map<ui64, std::map<ui32, std::shared_ptr<TPortionInfo>>> local; std::vector<std::shared_ptr<TPortionInfo>> portionsLock; @@ -170,9 +170,9 @@ bool TSourceCursor::Start(const std::shared_ptr<TSharedBlobsManager>& sharedBlob NextPathId = PortionsForSend.begin()->first; NextPortionId = PortionsForSend.begin()->second.begin()->first; - AFL_VERIFY(Next(sharedBlobsManager, index)); + AFL_VERIFY(Next(storagesManager, index)); } else { - BuildSelection(sharedBlobsManager, index); + BuildSelection(storagesManager, index); } IsStartedFlag = true; return true; diff --git a/ydb/core/tx/columnshard/data_sharing/source/session/cursor.h b/ydb/core/tx/columnshard/data_sharing/source/session/cursor.h index 7e527758bf0..3f4cdba86c1 100644 --- a/ydb/core/tx/columnshard/data_sharing/source/session/cursor.h +++ b/ydb/core/tx/columnshard/data_sharing/source/session/cursor.h @@ -31,7 +31,7 @@ private: THashMap<ui64, TString> PathPortionHashes; bool IsStartedFlag = false; YDB_ACCESSOR(bool, StaticSaved, false); - void BuildSelection(const std::shared_ptr<TSharedBlobsManager>& sharedBlobsManager, const TVersionedIndex& index); + void BuildSelection(const std::shared_ptr<IStoragesManager>& storagesManager, const TVersionedIndex& index); public: bool IsAckDataReceived() const { return AckReceivedForPackIdx == PackIdx; @@ -88,7 +88,7 @@ public: return Links; } - bool Next(const std::shared_ptr<TSharedBlobsManager>& sharedBlobsManager, const TVersionedIndex& index); + bool Next(const std::shared_ptr<IStoragesManager>& storagesManager, const TVersionedIndex& index); bool IsValid() { return Selected.size(); @@ -96,7 +96,7 @@ public: TSourceCursor(const TTabletId selfTabletId, const std::set<ui64>& pathIds, const TTransferContext transferContext); - bool Start(const std::shared_ptr<TSharedBlobsManager>& sharedBlobsManager, const THashMap<ui64, std::vector<std::shared_ptr<TPortionInfo>>>& portions, const TVersionedIndex& index); + bool Start(const std::shared_ptr<IStoragesManager>& storagesManager, const THashMap<ui64, std::vector<std::shared_ptr<TPortionInfo>>>& portions, const TVersionedIndex& index); NKikimrColumnShardDataSharingProto::TSourceSession::TCursorDynamic SerializeDynamicToProto() const; NKikimrColumnShardDataSharingProto::TSourceSession::TCursorStatic SerializeStaticToProto() const; diff --git a/ydb/core/tx/columnshard/data_sharing/source/session/source.cpp b/ydb/core/tx/columnshard/data_sharing/source/session/source.cpp index 3e9ad429c18..7c3e244ade1 100644 --- a/ydb/core/tx/columnshard/data_sharing/source/session/source.cpp +++ b/ydb/core/tx/columnshard/data_sharing/source/session/source.cpp @@ -48,7 +48,7 @@ TConclusion<std::unique_ptr<NTabletFlatExecutor::ITransaction>> TSourceSession:: return ackResult; } if (Cursor->IsReadyForNext()) { - Cursor->Next(self->GetStoragesManager()->GetSharedBlobsManager(), self->GetIndexAs<TColumnEngineForLogs>().GetVersionedIndex()); + Cursor->Next(self->GetStoragesManager(), self->GetIndexAs<TColumnEngineForLogs>().GetVersionedIndex()); return std::unique_ptr<NTabletFlatExecutor::ITransaction>(new TTxDataAckToSource(self, selfPtr, "ack_to_source_on_ack_data")); } else { return std::unique_ptr<NTabletFlatExecutor::ITransaction>(new TTxWriteSourceCursor(self, selfPtr, "write_source_cursor_on_ack_data")); @@ -61,15 +61,15 @@ TConclusion<std::unique_ptr<NTabletFlatExecutor::ITransaction>> TSourceSession:: return ackResult; } if (Cursor->IsReadyForNext()) { - Cursor->Next(self->GetStoragesManager()->GetSharedBlobsManager(), self->GetIndexAs<TColumnEngineForLogs>().GetVersionedIndex()); + Cursor->Next(self->GetStoragesManager(), self->GetIndexAs<TColumnEngineForLogs>().GetVersionedIndex()); return std::unique_ptr<NTabletFlatExecutor::ITransaction>(new TTxDataAckToSource(self, selfPtr, "ack_to_source_on_ack_links")); } else { return std::unique_ptr<NTabletFlatExecutor::ITransaction>(new TTxWriteSourceCursor(self, selfPtr, "write_source_cursor_on_ack_links")); } } -void TSourceSession::ActualizeDestination(const std::shared_ptr<NDataLocks::TManager>& dataLocksManager) { - AFL_VERIFY(IsStarted() || IsStarting()); +void TSourceSession::ActualizeDestination(const NColumnShard::TColumnShard& shard, const std::shared_ptr<NDataLocks::TManager>& dataLocksManager) { + AFL_VERIFY(IsInProgress() || IsPrepared()); AFL_VERIFY(Cursor); if (Cursor->IsValid()) { if (!Cursor->IsAckDataReceived()) { @@ -93,14 +93,14 @@ void TSourceSession::ActualizeDestination(const std::shared_ptr<NDataLocks::TMan auto ev = std::make_unique<NEvents::TEvFinishedFromSource>(GetSessionId(), SelfTabletId); NActors::TActivationContext::AsActorContext().Send(MakePipePerNodeCacheID(false), new TEvPipeCache::TEvForward(ev.release(), (ui64)DestinationTabletId, true), IEventHandle::FlagTrackDelivery, GetRuntimeId()); - Finish(dataLocksManager); + Finish(shard, dataLocksManager); } } bool TSourceSession::DoStart(const NColumnShard::TColumnShard& shard, const THashMap<ui64, std::vector<std::shared_ptr<TPortionInfo>>>& portions) { AFL_VERIFY(Cursor); - if (Cursor->Start(shard.GetStoragesManager()->GetSharedBlobsManager(), portions, shard.GetIndexAs<TColumnEngineForLogs>().GetVersionedIndex())) { - ActualizeDestination(shard.GetDataLocksManager()); + if (Cursor->Start(shard.GetStoragesManager(), portions, shard.GetIndexAs<TColumnEngineForLogs>().GetVersionedIndex())) { + ActualizeDestination(shard, shard.GetDataLocksManager()); return true; } else { return false; diff --git a/ydb/core/tx/columnshard/data_sharing/source/session/source.h b/ydb/core/tx/columnshard/data_sharing/source/session/source.h index 319302fe92d..903fc61c178 100644 --- a/ydb/core/tx/columnshard/data_sharing/source/session/source.h +++ b/ydb/core/tx/columnshard/data_sharing/source/session/source.h @@ -58,21 +58,21 @@ public: AFL_VERIFY(!!Cursor); return Cursor; } - - bool TryNextCursor(const ui32 packIdx, const std::shared_ptr<TSharedBlobsManager>& sharedBlobsManager, const TVersionedIndex& index) { +/* + bool TryNextCursor(const ui32 packIdx, const std::shared_ptr<IStoragesManager>& storagesManager, const TVersionedIndex& index) { AFL_VERIFY(Cursor); if (packIdx != Cursor->GetPackIdx()) { return false; } - Cursor->Next(sharedBlobsManager, index); + Cursor->Next(storagesManager, index); return true; } - +*/ [[nodiscard]] TConclusion<std::unique_ptr<NTabletFlatExecutor::ITransaction>> AckFinished(NColumnShard::TColumnShard* self, const std::shared_ptr<TSourceSession>& selfPtr); [[nodiscard]] TConclusion<std::unique_ptr<NTabletFlatExecutor::ITransaction>> AckData(NColumnShard::TColumnShard* self, const ui32 receivedPackIdx, const std::shared_ptr<TSourceSession>& selfPtr); [[nodiscard]] TConclusion<std::unique_ptr<NTabletFlatExecutor::ITransaction>> AckLinks(NColumnShard::TColumnShard* self, const TTabletId tabletId, const ui32 packIdx, const std::shared_ptr<TSourceSession>& selfPtr); - void ActualizeDestination(const std::shared_ptr<NDataLocks::TManager>& dataLocksManager); + void ActualizeDestination(const NColumnShard::TColumnShard& shard, const std::shared_ptr<NDataLocks::TManager>& dataLocksManager); NKikimrColumnShardDataSharingProto::TSourceSession SerializeDataToProto() const { NKikimrColumnShardDataSharingProto::TSourceSession result; diff --git a/ydb/core/tx/columnshard/data_sharing/source/transactions/tx_data_ack_to_source.cpp b/ydb/core/tx/columnshard/data_sharing/source/transactions/tx_data_ack_to_source.cpp index 146a93bc706..5a9bb1cf127 100644 --- a/ydb/core/tx/columnshard/data_sharing/source/transactions/tx_data_ack_to_source.cpp +++ b/ydb/core/tx/columnshard/data_sharing/source/transactions/tx_data_ack_to_source.cpp @@ -34,7 +34,7 @@ bool TTxDataAckToSource::DoExecute(NTabletFlatExecutor::TTransactionContext& txc } void TTxDataAckToSource::DoComplete(const TActorContext& /*ctx*/) { - Session->ActualizeDestination(Self->GetDataLocksManager()); + Session->ActualizeDestination(*Self, Self->GetDataLocksManager()); } }
\ No newline at end of file diff --git a/ydb/core/tx/columnshard/data_sharing/source/transactions/tx_start_to_source.cpp b/ydb/core/tx/columnshard/data_sharing/source/transactions/tx_start_to_source.cpp index 687b2e3661e..31e5d68768e 100644 --- a/ydb/core/tx/columnshard/data_sharing/source/transactions/tx_start_to_source.cpp +++ b/ydb/core/tx/columnshard/data_sharing/source/transactions/tx_start_to_source.cpp @@ -14,7 +14,6 @@ bool TTxStartToSource::DoExecute(NTabletFlatExecutor::TTransactionContext& txc, void TTxStartToSource::DoComplete(const TActorContext& /*ctx*/) { AFL_DEBUG(NKikimrServices::TX_COLUMNSHARD)("info", "TTxStartToSource::Complete"); AFL_VERIFY(Sessions->emplace(Session->GetSessionId(), Session).second); - Session->Start(*Self); } }
\ No newline at end of file diff --git a/ydb/core/tx/columnshard/hooks/testing/controller.cpp b/ydb/core/tx/columnshard/hooks/testing/controller.cpp index 4943addc756..e47dc08dcd6 100644 --- a/ydb/core/tx/columnshard/hooks/testing/controller.cpp +++ b/ydb/core/tx/columnshard/hooks/testing/controller.cpp @@ -58,9 +58,11 @@ void TController::CheckInvariants(const ::NKikimr::NColumnShard::TColumnShard& s auto manager = shard.GetStoragesManager()->GetOperatorVerified(i.first); const NOlap::TTabletsByBlob blobs = manager->GetBlobsToDelete(); for (auto b = blobs.GetIterator(); b.IsValid(); ++b) { + Cerr << shard.TabletID() << " SHARING_REMOVE_LOCAL:" << b.GetBlobId().ToStringNew() << " FROM " << b.GetTabletId() << Endl; i.second.RemoveSharing(b.GetTabletId(), b.GetBlobId()); } for (auto b = blobs.GetIterator(); b.IsValid(); ++b) { + Cerr << shard.TabletID() << " BORROWED_REMOVE_LOCAL:" << b.GetBlobId().ToStringNew() << " FROM " << b.GetTabletId() << Endl; i.second.RemoveBorrowed(b.GetTabletId(), b.GetBlobId()); } } diff --git a/ydb/core/tx/columnshard/hooks/testing/controller.h b/ydb/core/tx/columnshard/hooks/testing/controller.h index 94bf6436a97..c8211afb544 100644 --- a/ydb/core/tx/columnshard/hooks/testing/controller.h +++ b/ydb/core/tx/columnshard/hooks/testing/controller.h @@ -182,10 +182,8 @@ protected: } } virtual void DoOnDataSharingStarted(const ui64 /*tabletId*/, const TString& sessionId) override { + // dont check here. on finish only TGuard<TMutex> g(Mutex); - if (SharingIds.empty()) { - CheckInvariants(); - } SharingIds.emplace(sessionId); } |
