diff options
| author | Aleksei Borzenkov <[email protected]> | 2024-01-26 08:57:21 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2024-01-26 08:57:21 +0300 |
| commit | a8d82c5c4e82ecf2055f2f0f07e87aeeff5e0306 (patch) | |
| tree | e765e513dc64b4d29d0826e30df9753a6d656984 | |
| parent | 79d70ca8feff2f57b30bb609a0e5c2e5846d4870 (diff) | |
Support table changes observer KIKIMR-20853 (#1310)
| -rw-r--r-- | ydb/core/tablet_flat/CMakeLists.darwin-arm64.txt | 1 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/CMakeLists.darwin-x86_64.txt | 1 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/CMakeLists.linux-aarch64.txt | 1 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/CMakeLists.linux-x86_64.txt | 1 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/CMakeLists.windows-x86_64.txt | 1 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/flat_database.cpp | 5 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/flat_database.h | 3 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/flat_executor.cpp | 52 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/flat_executor.h | 2 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/flat_table.cpp | 12 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/flat_table.h | 4 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/flat_table_observer.cpp | 1 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/flat_table_observer.h | 31 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/tablet_flat_executor.cpp | 8 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/tablet_flat_executor.h | 3 | ||||
| -rw-r--r-- | ydb/core/tablet_flat/ya.make | 2 |
16 files changed, 118 insertions, 10 deletions
diff --git a/ydb/core/tablet_flat/CMakeLists.darwin-arm64.txt b/ydb/core/tablet_flat/CMakeLists.darwin-arm64.txt index c6d0640e5ff..1ec29730846 100644 --- a/ydb/core/tablet_flat/CMakeLists.darwin-arm64.txt +++ b/ydb/core/tablet_flat/CMakeLists.darwin-arm64.txt @@ -124,6 +124,7 @@ target_sources(ydb-core-tablet_flat PRIVATE ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_part.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_misc.cpp + ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_observer.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/probes.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/shared_handle.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/shared_sausagecache.cpp diff --git a/ydb/core/tablet_flat/CMakeLists.darwin-x86_64.txt b/ydb/core/tablet_flat/CMakeLists.darwin-x86_64.txt index c6d0640e5ff..1ec29730846 100644 --- a/ydb/core/tablet_flat/CMakeLists.darwin-x86_64.txt +++ b/ydb/core/tablet_flat/CMakeLists.darwin-x86_64.txt @@ -124,6 +124,7 @@ target_sources(ydb-core-tablet_flat PRIVATE ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_part.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_misc.cpp + ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_observer.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/probes.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/shared_handle.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/shared_sausagecache.cpp diff --git a/ydb/core/tablet_flat/CMakeLists.linux-aarch64.txt b/ydb/core/tablet_flat/CMakeLists.linux-aarch64.txt index 8dd469d1970..6a78f839ed8 100644 --- a/ydb/core/tablet_flat/CMakeLists.linux-aarch64.txt +++ b/ydb/core/tablet_flat/CMakeLists.linux-aarch64.txt @@ -125,6 +125,7 @@ target_sources(ydb-core-tablet_flat PRIVATE ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_part.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_misc.cpp + ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_observer.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/probes.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/shared_handle.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/shared_sausagecache.cpp diff --git a/ydb/core/tablet_flat/CMakeLists.linux-x86_64.txt b/ydb/core/tablet_flat/CMakeLists.linux-x86_64.txt index 8dd469d1970..6a78f839ed8 100644 --- a/ydb/core/tablet_flat/CMakeLists.linux-x86_64.txt +++ b/ydb/core/tablet_flat/CMakeLists.linux-x86_64.txt @@ -125,6 +125,7 @@ target_sources(ydb-core-tablet_flat PRIVATE ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_part.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_misc.cpp + ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_observer.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/probes.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/shared_handle.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/shared_sausagecache.cpp diff --git a/ydb/core/tablet_flat/CMakeLists.windows-x86_64.txt b/ydb/core/tablet_flat/CMakeLists.windows-x86_64.txt index c6d0640e5ff..1ec29730846 100644 --- a/ydb/core/tablet_flat/CMakeLists.windows-x86_64.txt +++ b/ydb/core/tablet_flat/CMakeLists.windows-x86_64.txt @@ -124,6 +124,7 @@ target_sources(ydb-core-tablet_flat PRIVATE ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_part.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_misc.cpp + ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/flat_table_observer.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/probes.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/shared_handle.cpp ${CMAKE_SOURCE_DIR}/ydb/core/tablet_flat/shared_sausagecache.cpp diff --git a/ydb/core/tablet_flat/flat_database.cpp b/ydb/core/tablet_flat/flat_database.cpp index ab46be9de03..5a47f8c31d1 100644 --- a/ydb/core/tablet_flat/flat_database.cpp +++ b/ydb/core/tablet_flat/flat_database.cpp @@ -495,6 +495,11 @@ const TDbStats& TDatabase::Counters() const noexcept return DatabaseImpl->Stats; } +void TDatabase::SetTableObserver(ui32 table, TIntrusivePtr<ITableObserver> ptr) noexcept +{ + Require(table)->SetTableObserver(std::move(ptr)); +} + TDatabase::TChg TDatabase::Head(ui32 table) const noexcept { if (table == Max<ui32>()) { diff --git a/ydb/core/tablet_flat/flat_database.h b/ydb/core/tablet_flat/flat_database.h index 85f81af257f..e35f58fbd32 100644 --- a/ydb/core/tablet_flat/flat_database.h +++ b/ydb/core/tablet_flat/flat_database.h @@ -6,6 +6,7 @@ #include "flat_dbase_change.h" #include "flat_dbase_misc.h" #include "flat_iterator.h" +#include "flat_table_observer.h" #include "util_basics.h" namespace NKikimr { @@ -57,6 +58,8 @@ public: TDatabase(TDatabaseImpl *databaseImpl = nullptr) noexcept; ~TDatabase(); + void SetTableObserver(ui32 table, TIntrusivePtr<ITableObserver> ptr) noexcept; + /* Returns durable monotonic change number for table or entire database on default (table = Max<ui32>()). Serial is incremented for each successful Commit(). AHTUNG: Serial may go to the past in case of diff --git a/ydb/core/tablet_flat/flat_executor.cpp b/ydb/core/tablet_flat/flat_executor.cpp index 9f559a8374c..13cdb2d535e 100644 --- a/ydb/core/tablet_flat/flat_executor.cpp +++ b/ydb/core/tablet_flat/flat_executor.cpp @@ -97,6 +97,32 @@ TTableSnapshotContext::~TTableSnapshotContext() = default; using namespace NResourceBroker; +class TExecutor::TActiveTransactionZone { +public: + explicit TActiveTransactionZone(TExecutor* self) noexcept + : Self(self) + { + Y_DEBUG_ABORT_UNLESS(!Self->ActiveTransaction); + Self->ActiveTransaction = true; + Active = true; + } + + ~TActiveTransactionZone() noexcept { + Done(); + } + + void Done() noexcept { + if (Active) { + Self->ActiveTransaction = false; + Active = false; + } + } + +private: + TExecutor* Self; + bool Active = false; +}; + TExecutor::TExecutor( NFlatExecutorSetup::ITablet* owner, const TActorId& ownerActorId) @@ -876,6 +902,9 @@ void TExecutor::ApplyFollowerUpdate(THolder<TEvTablet::TFUpdateBody> update) { if (update->IsSnapshot) // do nothing over snapshot after initial one return; + // Protect against recursive transactions in callbacks + TActiveTransactionZone activeTransaction(this); + TString schemeUpdate; TString dataUpdate; TStackVec<TString> partSwitches; @@ -940,6 +969,11 @@ void TExecutor::ApplyFollowerUpdate(THolder<TEvTablet::TFUpdateBody> update) { if (schemeUpdate) { ReadResourceProfile(); ReflectSchemeSettings(); + Owner->OnFollowerSchemaUpdated(); + } + + if (dataUpdate) { + Owner->OnFollowerDataUpdated(); } } @@ -1006,8 +1040,10 @@ void TExecutor::ApplyFollowerAuxUpdate(const TString &auxBody) { const TString aux = NPageCollection::TSlicer::Lz4()->Decode(auxBody); TProtoBox<NKikimrExecutorFlat::TFollowerAux> proto(aux); - if (proto.HasUserAuxUpdate()) + if (proto.HasUserAuxUpdate()) { + TActiveTransactionZone activeTransaction(this); Owner->OnLeaderUserAuxUpdate(std::move(proto.GetUserAuxUpdate())); + } } void TExecutor::RequestFromSharedCache(TAutoPtr<NPageCollection::TFetch> fetch, @@ -1644,9 +1680,7 @@ void TExecutor::Enqueue(TAutoPtr<ITransaction> self, const TActorContext &ctx) { } void TExecutor::ExecuteTransaction(TAutoPtr<TSeat> seat, const TActorContext &ctx) { - Y_DEBUG_ABORT_UNLESS(!ActiveTransaction); - - ActiveTransaction = true; + TActiveTransactionZone activeTransaction(this); ++seat->Retries; THPTimer cpuTimer; @@ -1736,7 +1770,7 @@ void TExecutor::ExecuteTransaction(TAutoPtr<TSeat> seat, const TActorContext &ct } PrivatePageCache->ResetTouchesAndToLoad(false); - ActiveTransaction = false; + activeTransaction.Done(); PlanTransactionActivation(); } @@ -2813,7 +2847,7 @@ void TExecutor::Handle(TEvTablet::TEvCommitResult::TPtr &ev, const TActorContext Y_ABORT_UNLESS(msg->Generation == Generation()); const ui32 step = msg->Step; - ActiveTransaction = true; + TActiveTransactionZone activeTransaction(this); GcLogic->OnCommitLog(step, msg->ConfirmedOnSend, ctx); CommitManager->Confirm(step); @@ -2907,7 +2941,7 @@ void TExecutor::Handle(TEvTablet::TEvCommitResult::TPtr &ev, const TActorContext std::move(msg->GroupWrittenOps), ctx); - ActiveTransaction = false; + activeTransaction.Done(); PlanTransactionActivation(); MaybeRelaxRejectProbability(); @@ -3236,7 +3270,7 @@ void TExecutor::Handle(NOps::TEvResult *ops, TProdCompact *msg, bool cancelled) return Broken(); } - ActiveTransaction = true; + TActiveTransactionZone activeTransaction(this); const ui64 snapStamp = msg->Params->Edge.TxStamp ? msg->Params->Edge.TxStamp : MakeGenStepPair(Generation(), msg->Step); @@ -3469,7 +3503,7 @@ void TExecutor::Handle(NOps::TEvResult *ops, TProdCompact *msg, bool cancelled) Owner->CompactionComplete(tableId, OwnerCtx()); MaybeRelaxRejectProbability(); - ActiveTransaction = false; + activeTransaction.Done(); if (LogicSnap->MayFlush(false)) { MakeLogSnapshot(); diff --git a/ydb/core/tablet_flat/flat_executor.h b/ydb/core/tablet_flat/flat_executor.h index 594ccb0dbec..a8feffc9dda 100644 --- a/ydb/core/tablet_flat/flat_executor.h +++ b/ydb/core/tablet_flat/flat_executor.h @@ -413,6 +413,8 @@ class TExecutor THashMap<ui64, THolder<TScanSnapshot>> ScanSnapshots; ui64 ScanSnapshotId = 1; + class TActiveTransactionZone; + bool ActiveTransaction = false; bool BrokenTransaction = false; ui32 ActivateTransactionWaiting = 0; diff --git a/ydb/core/tablet_flat/flat_table.cpp b/ydb/core/tablet_flat/flat_table.cpp index e3feb7e3ea3..fb5f8d46ba7 100644 --- a/ydb/core/tablet_flat/flat_table.cpp +++ b/ydb/core/tablet_flat/flat_table.cpp @@ -830,6 +830,9 @@ void TTable::Update(ERowOp rop, TRawVals key, TOpsRef ops, TArrayRef<const TMemG } MemTable().Update(rop, key, ops, apart, rowVersion, CommittedTransactions); + if (TableObserver) { + TableObserver->OnUpdate(rop, key, ops, rowVersion); + } } void TTable::AddTxRef(ui64 txId) @@ -863,6 +866,10 @@ void TTable::UpdateTx(ERowOp rop, TRawVals key, TOpsRef ops, TArrayRef<const TMe } else { Y_DEBUG_ABORT_UNLESS(TxRefs[txId] > 0); } + + if (TableObserver) { + TableObserver->OnUpdateTx(rop, key, ops, txId); + } } void TTable::CommitTx(ui64 txId, TRowVersion rowVersion) @@ -1338,6 +1345,11 @@ TCompactionStats TTable::GetCompactionStats() const return stats; } +void TTable::SetTableObserver(TIntrusivePtr<ITableObserver> ptr) noexcept +{ + TableObserver = std::move(ptr); +} + void TPartStats::Add(const TPartView& partView) { PartsCount += 1; diff --git a/ydb/core/tablet_flat/flat_table.h b/ydb/core/tablet_flat/flat_table.h index 0b0cf796393..cd718f4be50 100644 --- a/ydb/core/tablet_flat/flat_table.h +++ b/ydb/core/tablet_flat/flat_table.h @@ -13,6 +13,7 @@ #include "flat_table_stats.h" #include "flat_table_subset.h" #include "flat_table_misc.h" +#include "flat_table_observer.h" #include "flat_sausage_solid.h" #include "util_basics.h" @@ -322,7 +323,7 @@ public: TCompactionStats GetCompactionStats() const; - void FillTxStatusCache(THashMap<TLogoBlobID, TSharedData>& cache) const noexcept; + void SetTableObserver(TIntrusivePtr<ITableObserver> ptr) noexcept; private: TMemTable& MemTable(); @@ -358,6 +359,7 @@ private: absl::flat_hash_set<ui64> CheckTransactions; TTransactionMap CommittedTransactions; TTransactionSet RemovedTransactions; + TIntrusivePtr<ITableObserver> TableObserver; private: struct TRollbackRemoveTxRef { diff --git a/ydb/core/tablet_flat/flat_table_observer.cpp b/ydb/core/tablet_flat/flat_table_observer.cpp new file mode 100644 index 00000000000..5adda17187d --- /dev/null +++ b/ydb/core/tablet_flat/flat_table_observer.cpp @@ -0,0 +1 @@ +#include "flat_table_observer.h" diff --git a/ydb/core/tablet_flat/flat_table_observer.h b/ydb/core/tablet_flat/flat_table_observer.h new file mode 100644 index 00000000000..374a0a0152d --- /dev/null +++ b/ydb/core/tablet_flat/flat_table_observer.h @@ -0,0 +1,31 @@ +#pragma once +#include "defs.h" +#include "flat_row_eggs.h" +#include "flat_update_op.h" + +#include <util/generic/ptr.h> + +namespace NKikimr::NTable { + + class ITableObserver : public TThrRefBase { + public: + /** + * Called when a new update is applied to the table + */ + virtual void OnUpdate( + ERowOp rop, + TArrayRef<const TRawTypeValue> key, + TArrayRef<const TUpdateOp> ops, + TRowVersion rowVersion) = 0; + + /** + * Called when an uncommitted update is applied to the table + */ + virtual void OnUpdateTx( + ERowOp rop, + TArrayRef<const TRawTypeValue> key, + TArrayRef<const TUpdateOp> ops, + ui64 txId) = 0; + }; + +} // namespace NKikimr::NTable diff --git a/ydb/core/tablet_flat/tablet_flat_executor.cpp b/ydb/core/tablet_flat/tablet_flat_executor.cpp index 6b1ccd28e4f..d075515a1e3 100644 --- a/ydb/core/tablet_flat/tablet_flat_executor.cpp +++ b/ydb/core/tablet_flat/tablet_flat_executor.cpp @@ -68,6 +68,14 @@ namespace NFlatExecutorSetup { void ITablet::OnFollowersCountChanged() { // nothing by default } + + void ITablet::OnFollowerSchemaUpdated() { + // nothing by default + } + + void ITablet::OnFollowerDataUpdated() { + // nothing by default + } } }} diff --git a/ydb/core/tablet_flat/tablet_flat_executor.h b/ydb/core/tablet_flat/tablet_flat_executor.h index 33ff77f951d..edd006ed266 100644 --- a/ydb/core/tablet_flat/tablet_flat_executor.h +++ b/ydb/core/tablet_flat/tablet_flat_executor.h @@ -505,6 +505,9 @@ namespace NFlatExecutorSetup { virtual void OnFollowersCountChanged(); + virtual void OnFollowerSchemaUpdated(); + virtual void OnFollowerDataUpdated(); + // create transaction? protected: ITablet(TTabletStorageInfo *info, const TActorId &tablet) diff --git a/ydb/core/tablet_flat/ya.make b/ydb/core/tablet_flat/ya.make index 8a6604935eb..7c97ebd5045 100644 --- a/ydb/core/tablet_flat/ya.make +++ b/ydb/core/tablet_flat/ya.make @@ -62,6 +62,8 @@ SRCS( flat_table_part.cpp flat_table_part.h flat_table_misc.cpp + flat_table_observer.cpp + flat_table_observer.h flat_update_op.h probes.cpp shared_handle.cpp |
