summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorivanmorozov333 <[email protected]>2024-06-10 15:27:48 +0300
committerGitHub <[email protected]>2024-06-10 15:27:48 +0300
commit4cef62d2b9199a9d72a23c68179ef46d8abe05ef (patch)
tree2063ab957fc45884715acc02cbc6051b988401ed
parent1d951c4cf7f550a09cbf91b3b09b16e8511cc441 (diff)
fixes for resharding tests (#5360)
-rw-r--r--ydb/core/kqp/ut/olap/blobs_sharing_ut.cpp36
-rw-r--r--ydb/core/tx/columnshard/blobs_action/abstract/blob_set.cpp11
-rw-r--r--ydb/core/tx/columnshard/blobs_action/abstract/blob_set.h22
-rw-r--r--ydb/core/tx/columnshard/blobs_action/abstract/gc.h4
-rw-r--r--ydb/core/tx/columnshard/blobs_action/abstract/gc_actor.h37
-rw-r--r--ydb/core/tx/columnshard/blobs_action/abstract/remove.h4
-rw-r--r--ydb/core/tx/columnshard/blobs_action/abstract/storage.h1
-rw-r--r--ydb/core/tx/columnshard/blobs_action/abstract/storages_manager.h4
-rw-r--r--ydb/core/tx/columnshard/blobs_action/bs/blob_manager.h4
-rw-r--r--ydb/core/tx/columnshard/blobs_action/bs/gc_actor.h2
-rw-r--r--ydb/core/tx/columnshard/blobs_action/bs/storage.h4
-rw-r--r--ydb/core/tx/columnshard/blobs_action/events/delete_blobs.h4
-rw-r--r--ydb/core/tx/columnshard/blobs_action/protos/events.proto7
-rw-r--r--ydb/core/tx/columnshard/blobs_action/tier/gc_actor.h2
-rw-r--r--ydb/core/tx/columnshard/blobs_action/tier/gc_info.h5
-rw-r--r--ydb/core/tx/columnshard/blobs_action/tier/remove.cpp23
-rw-r--r--ydb/core/tx/columnshard/blobs_action/tier/remove.h20
-rw-r--r--ydb/core/tx/columnshard/blobs_action/tier/storage.h5
-rw-r--r--ydb/core/tx/columnshard/blobs_action/transaction/tx_remove_blobs.cpp8
-rw-r--r--ydb/core/tx/columnshard/blobs_action/transaction/tx_remove_blobs.h23
-rw-r--r--ydb/core/tx/columnshard/columnshard_impl.cpp14
-rw-r--r--ydb/core/tx/columnshard/data_locks/locks/abstract.h13
-rw-r--r--ydb/core/tx/columnshard/data_locks/locks/composite.h10
-rw-r--r--ydb/core/tx/columnshard/data_locks/locks/list.h8
-rw-r--r--ydb/core/tx/columnshard/data_locks/locks/snapshot.h17
-rw-r--r--ydb/core/tx/columnshard/data_locks/manager/manager.cpp14
-rw-r--r--ydb/core/tx/columnshard/data_locks/manager/manager.h4
-rw-r--r--ydb/core/tx/columnshard/data_sharing/common/session/common.cpp62
-rw-r--r--ydb/core/tx/columnshard/data_sharing/common/session/common.h28
-rw-r--r--ydb/core/tx/columnshard/data_sharing/destination/events/transfer.cpp12
-rw-r--r--ydb/core/tx/columnshard/data_sharing/destination/events/transfer.h2
-rw-r--r--ydb/core/tx/columnshard/data_sharing/destination/session/destination.cpp4
-rw-r--r--ydb/core/tx/columnshard/data_sharing/destination/transactions/tx_finish_from_source.cpp2
-rw-r--r--ydb/core/tx/columnshard/data_sharing/destination/transactions/tx_start_from_initiator.cpp1
-rw-r--r--ydb/core/tx/columnshard/data_sharing/manager/sessions.cpp23
-rw-r--r--ydb/core/tx/columnshard/data_sharing/manager/sessions.h13
-rw-r--r--ydb/core/tx/columnshard/data_sharing/manager/shared_blobs.cpp20
-rw-r--r--ydb/core/tx/columnshard/data_sharing/manager/shared_blobs.h32
-rw-r--r--ydb/core/tx/columnshard/data_sharing/source/session/cursor.cpp14
-rw-r--r--ydb/core/tx/columnshard/data_sharing/source/session/cursor.h6
-rw-r--r--ydb/core/tx/columnshard/data_sharing/source/session/source.cpp14
-rw-r--r--ydb/core/tx/columnshard/data_sharing/source/session/source.h10
-rw-r--r--ydb/core/tx/columnshard/data_sharing/source/transactions/tx_data_ack_to_source.cpp2
-rw-r--r--ydb/core/tx/columnshard/data_sharing/source/transactions/tx_start_to_source.cpp1
-rw-r--r--ydb/core/tx/columnshard/hooks/testing/controller.cpp2
-rw-r--r--ydb/core/tx/columnshard/hooks/testing/controller.h4
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);
}