summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorAleksei Borzenkov <[email protected]>2024-01-26 08:57:21 +0300
committerGitHub <[email protected]>2024-01-26 08:57:21 +0300
commita8d82c5c4e82ecf2055f2f0f07e87aeeff5e0306 (patch)
treee765e513dc64b4d29d0826e30df9753a6d656984
parent79d70ca8feff2f57b30bb609a0e5c2e5846d4870 (diff)
Support table changes observer KIKIMR-20853 (#1310)
-rw-r--r--ydb/core/tablet_flat/CMakeLists.darwin-arm64.txt1
-rw-r--r--ydb/core/tablet_flat/CMakeLists.darwin-x86_64.txt1
-rw-r--r--ydb/core/tablet_flat/CMakeLists.linux-aarch64.txt1
-rw-r--r--ydb/core/tablet_flat/CMakeLists.linux-x86_64.txt1
-rw-r--r--ydb/core/tablet_flat/CMakeLists.windows-x86_64.txt1
-rw-r--r--ydb/core/tablet_flat/flat_database.cpp5
-rw-r--r--ydb/core/tablet_flat/flat_database.h3
-rw-r--r--ydb/core/tablet_flat/flat_executor.cpp52
-rw-r--r--ydb/core/tablet_flat/flat_executor.h2
-rw-r--r--ydb/core/tablet_flat/flat_table.cpp12
-rw-r--r--ydb/core/tablet_flat/flat_table.h4
-rw-r--r--ydb/core/tablet_flat/flat_table_observer.cpp1
-rw-r--r--ydb/core/tablet_flat/flat_table_observer.h31
-rw-r--r--ydb/core/tablet_flat/tablet_flat_executor.cpp8
-rw-r--r--ydb/core/tablet_flat/tablet_flat_executor.h3
-rw-r--r--ydb/core/tablet_flat/ya.make2
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