diff options
| author | Alek5andr-Kotov <[email protected]> | 2024-10-02 13:51:05 +0300 |
|---|---|---|
| committer | GitHub <[email protected]> | 2024-10-02 13:51:05 +0300 |
| commit | bb6b434bfbe8312ea6b418cf867738c5b1c8d133 (patch) | |
| tree | db3f18279af6a4a16db1016cf90336c1d7d735cf | |
| parent | 785e5b728b8abf9a340c24d83f6123b1335982f5 (diff) | |
use `TTxControl` to commit transactions (#9890)
6 files changed, 149 insertions, 32 deletions
diff --git a/ydb/public/sdk/cpp/client/ydb_table/impl/table_client.h b/ydb/public/sdk/cpp/client/ydb_table/impl/table_client.h index 2d20737210e..da681f0d959 100644 --- a/ydb/public/sdk/cpp/client/ydb_table/impl/table_client.h +++ b/ydb/public/sdk/cpp/client/ydb_table/impl/table_client.h @@ -79,7 +79,7 @@ public: CacheMissCounter.Inc(); return ::NYdb::NSessionPool::InjectSessionStatusInterception(session.SessionImpl_, - ExecuteDataQueryInternal(session, query, txControl, params, settings, false), + ExecuteDataQueryImpl(session, query, txControl, params, settings, false), true, GetMinTimeToTouch(Settings_.SessionPoolSettings_)); } @@ -96,7 +96,7 @@ public: return ::NYdb::NSessionPool::InjectSessionStatusInterception<TDataQueryResult>( session.SessionImpl_, - session.Client_->ExecuteDataQueryInternal(session, dataQuery, txControl, params, settings, fromCache), + session.Client_->ExecuteDataQueryImpl(session, dataQuery, txControl, params, settings, fromCache), true, GetMinTimeToTouch(session.Client_->Settings_.SessionPoolSettings_), cb); @@ -172,6 +172,32 @@ private: static void CollectQuerySize(const TDataQuery&, NSdkStats::TAtomicHistogram<::NMonitoring::THistogram>&); template <typename TQueryType, typename TParamsType> + TAsyncDataQueryResult ExecuteDataQueryImpl(const TSession& session, const TQueryType& query, + const TTxControl& txControl, TParamsType params, + const TExecDataQuerySettings& settings, bool fromCache + ) { + if (!txControl.Tx_.Defined() || !txControl.CommitTx_) { + return ExecuteDataQueryInternal(session, query, txControl, params, settings, fromCache); + } + + auto onPrecommitCompleted = [this, session, query, txControl, params, settings, fromCache](const NThreading::TFuture<TStatus>& f) { + TStatus status = f.GetValueSync(); + if (!status.IsSuccess()) { + return NThreading::MakeFuture(TDataQueryResult(std::move(status), + {}, + txControl.Tx_, + Nothing(), + false, + Nothing())); + } + + return ExecuteDataQueryInternal(session, query, txControl, params, settings, fromCache); + }; + + return txControl.Tx_->Precommit().Apply(onPrecommitCompleted); + } + + template <typename TQueryType, typename TParamsType> TAsyncDataQueryResult ExecuteDataQueryInternal(const TSession& session, const TQueryType& query, const TTxControl& txControl, TParamsType params, const TExecDataQuerySettings& settings, bool fromCache @@ -180,8 +206,8 @@ private: request.set_session_id(session.GetId()); auto txControlProto = request.mutable_tx_control(); txControlProto->set_commit_tx(txControl.CommitTx_); - if (txControl.TxId_) { - txControlProto->set_tx_id(*txControl.TxId_); + if (txControl.Tx_.Defined()) { + txControlProto->set_tx_id(txControl.Tx_->GetId()); } else { SetTxSettings(txControl.BeginTx_, txControlProto->mutable_begin_tx()); } diff --git a/ydb/public/sdk/cpp/client/ydb_table/impl/transaction.cpp b/ydb/public/sdk/cpp/client/ydb_table/impl/transaction.cpp index 2c00e9ae9d4..f6d690b0eb1 100644 --- a/ydb/public/sdk/cpp/client/ydb_table/impl/transaction.cpp +++ b/ydb/public/sdk/cpp/client/ydb_table/impl/transaction.cpp @@ -9,10 +9,8 @@ TTransaction::TImpl::TImpl(const TSession& session, const TString& txId) { } -TAsyncCommitTransactionResult TTransaction::TImpl::Commit(const TCommitTxSettings& settings) +TAsyncStatus TTransaction::TImpl::Precommit() const { - ChangesAreAccepted = false; - auto result = NThreading::MakeFuture(TStatus(EStatus::SUCCESS, {})); for (auto& callback : PrecommitCallbacks) { @@ -27,6 +25,15 @@ TAsyncCommitTransactionResult TTransaction::TImpl::Commit(const TCommitTxSetting result = result.Apply(action); } + return result; +} + +TAsyncCommitTransactionResult TTransaction::TImpl::Commit(const TCommitTxSettings& settings) +{ + ChangesAreAccepted = false; + + auto result = Precommit(); + auto precommitsCompleted = [this, settings](const TAsyncStatus& result) mutable { if (const TStatus& status = result.GetValue(); !status.IsSuccess()) { return NThreading::MakeFuture(TCommitTransactionResult(TStatus(status), Nothing())); diff --git a/ydb/public/sdk/cpp/client/ydb_table/impl/transaction.h b/ydb/public/sdk/cpp/client/ydb_table/impl/transaction.h index 12df59f8065..e73ac616bb1 100644 --- a/ydb/public/sdk/cpp/client/ydb_table/impl/transaction.h +++ b/ydb/public/sdk/cpp/client/ydb_table/impl/transaction.h @@ -16,6 +16,7 @@ public: return !TxId_.empty(); } + TAsyncStatus Precommit() const; TAsyncCommitTransactionResult Commit(const TCommitTxSettings& settings = TCommitTxSettings()); TAsyncStatus Rollback(const TRollbackTxSettings& settings = TRollbackTxSettings()); diff --git a/ydb/public/sdk/cpp/client/ydb_table/table.cpp b/ydb/public/sdk/cpp/client/ydb_table/table.cpp index 767964fc4d7..6ef08362cca 100644 --- a/ydb/public/sdk/cpp/client/ydb_table/table.cpp +++ b/ydb/public/sdk/cpp/client/ydb_table/table.cpp @@ -1988,7 +1988,7 @@ const TString& TSession::GetId() const { //////////////////////////////////////////////////////////////////////////////// TTxControl::TTxControl(const TTransaction& tx) - : TxId_(tx.GetId()) + : Tx_(tx) {} TTxControl::TTxControl(const TTxSettings& begin) @@ -2011,6 +2011,11 @@ bool TTransaction::IsActive() const return TransactionImpl_->IsActive(); } +TAsyncStatus TTransaction::Precommit() const +{ + return TransactionImpl_->Precommit(); +} + TAsyncCommitTransactionResult TTransaction::Commit(const TCommitTxSettings& settings) { return TransactionImpl_->Commit(settings); } diff --git a/ydb/public/sdk/cpp/client/ydb_table/table.h b/ydb/public/sdk/cpp/client/ydb_table/table.h index 3b50f810f7e..2d1b681431f 100644 --- a/ydb/public/sdk/cpp/client/ydb_table/table.h +++ b/ydb/public/sdk/cpp/client/ydb_table/table.h @@ -1274,30 +1274,7 @@ private: ETransactionMode Mode_; }; -class TTxControl { - friend class TTableClient; - -public: - using TSelf = TTxControl; - - static TTxControl Tx(const TTransaction& tx) { - return TTxControl(tx); - } - - static TTxControl BeginTx(const TTxSettings& settings = TTxSettings()) { - return TTxControl(settings); - } - - FLUENT_SETTING_FLAG(CommitTx); - -private: - TTxControl(const TTransaction& tx); - TTxControl(const TTxSettings& begin); - -private: - TMaybe<TString> TxId_; - TTxSettings BeginTx_; -}; +class TTxControl; enum class EAutoPartitioningPolicy { Disabled = 1, @@ -1844,6 +1821,8 @@ public: private: TTransaction(const TSession& session, const TString& txId); + TAsyncStatus Precommit() const; + class TImpl; std::shared_ptr<TImpl> TransactionImpl_; @@ -1851,6 +1830,33 @@ private: //////////////////////////////////////////////////////////////////////////////// +class TTxControl { + friend class TTableClient; + +public: + using TSelf = TTxControl; + + static TTxControl Tx(const TTransaction& tx) { + return TTxControl(tx); + } + + static TTxControl BeginTx(const TTxSettings& settings = TTxSettings()) { + return TTxControl(settings); + } + + FLUENT_SETTING_FLAG(CommitTx); + +private: + TTxControl(const TTransaction& tx); + TTxControl(const TTxSettings& begin); + +private: + TMaybe<TTransaction> Tx_; + TTxSettings BeginTx_; +}; + +//////////////////////////////////////////////////////////////////////////////// + //! Represents query identificator (e.g. used for prepared query) class TDataQuery { friend class TTableClient; diff --git a/ydb/public/sdk/cpp/client/ydb_topic/ut/topic_to_table_ut.cpp b/ydb/public/sdk/cpp/client/ydb_topic/ut/topic_to_table_ut.cpp index 4b529bfaaad..52f67d3f6f6 100644 --- a/ydb/public/sdk/cpp/client/ydb_topic/ut/topic_to_table_ut.cpp +++ b/ydb/public/sdk/cpp/client/ydb_topic/ut/topic_to_table_ut.cpp @@ -172,6 +172,8 @@ protected: void CheckTabletKeys(const TString& topicName); void DumpPQTabletKeys(const TString& topicName); + NTable::TDataQueryResult ExecuteDataQuery(NTable::TSession session, const TString& query, const NTable::TTxControl& control); + private: template<class E> E ReadEvent(TTopicReadSessionPtr reader, NTable::TTransaction& tx); @@ -1535,6 +1537,13 @@ void TFixture::TestTheCompletionOfATransaction(const TTransactionCompletionTestD } } +NTable::TDataQueryResult TFixture::ExecuteDataQuery(NTable::TSession session, const TString& query, const NTable::TTxControl& control) +{ + auto status = session.ExecuteDataQuery(query, control).GetValueSync(); + UNIT_ASSERT_C(status.IsSuccess(), status.GetIssues().ToString()); + return status; +} + Y_UNIT_TEST_F(WriteToTopic_Demo_11, TFixture) { for (auto endOfTransaction : {Commit, Rollback, CloseTableSession}) { @@ -2148,6 +2157,69 @@ Y_UNIT_TEST_F(WriteToTopic_Demo_41, TFixture) CommitTx(tx, EStatus::SESSION_EXPIRED); } +Y_UNIT_TEST_F(WriteToTopic_Demo_42, TFixture) +{ + CreateTopic("topic_A", TEST_CONSUMER); + + NTable::TSession tableSession = CreateTableSession(); + NTable::TTransaction tx = BeginTx(tableSession); + + for (size_t k = 0; k < 100; ++k) { + WriteToTopic("topic_A", TEST_MESSAGE_GROUP_ID, TString(1'000'000, 'a'), &tx); + } + + CloseTopicWriteSession("topic_A", TEST_MESSAGE_GROUP_ID); // gracefully close + + CommitTx(tx, EStatus::SUCCESS); + + auto messages = ReadFromTopic("topic_A", TEST_CONSUMER, TDuration::Seconds(2)); + UNIT_ASSERT_VALUES_EQUAL(messages.size(), 100); +} + +Y_UNIT_TEST_F(WriteToTopic_Demo_43, TFixture) +{ + // The recording stream will run into a quota. Before the commit, the client will receive confirmations + // for some of the messages. The `ExecuteDataQuery` call will wait for the rest. + CreateTopic("topic_A", TEST_CONSUMER); + + NTable::TSession tableSession = CreateTableSession(); + NTable::TTransaction tx = BeginTx(tableSession); + + for (size_t k = 0; k < 100; ++k) { + WriteToTopic("topic_A", TEST_MESSAGE_GROUP_ID, TString(1'000'000, 'a'), &tx); + } + + ExecuteDataQuery(tableSession, "SELECT 1", NTable::TTxControl::Tx(tx).CommitTx(true)); + + auto messages = ReadFromTopic("topic_A", TEST_CONSUMER, TDuration::Seconds(60)); + UNIT_ASSERT_VALUES_EQUAL(messages.size(), 100); +} + +Y_UNIT_TEST_F(WriteToTopic_Demo_44, TFixture) +{ + CreateTopic("topic_A", TEST_CONSUMER); + + NTable::TSession tableSession = CreateTableSession(); + + auto result = ExecuteDataQuery(tableSession, "SELECT 1", NTable::TTxControl::BeginTx()); + + NTable::TTransaction tx = *result.GetTransaction(); + + for (size_t k = 0; k < 100; ++k) { + WriteToTopic("topic_A", TEST_MESSAGE_GROUP_ID, TString(1'000'000, 'a'), &tx); + } + + WaitForAcks("topic_A", TEST_MESSAGE_GROUP_ID); + + auto messages = ReadFromTopic("topic_A", TEST_CONSUMER, TDuration::Seconds(60)); + UNIT_ASSERT_VALUES_EQUAL(messages.size(), 0); + + ExecuteDataQuery(tableSession, "SELECT 2", NTable::TTxControl::Tx(tx).CommitTx(true)); + + messages = ReadFromTopic("topic_A", TEST_CONSUMER, TDuration::Seconds(60)); + UNIT_ASSERT_VALUES_EQUAL(messages.size(), 100); +} + } } |
