summaryrefslogtreecommitdiffstats
diff options
context:
space:
mode:
authorAlek5andr-Kotov <[email protected]>2024-10-02 13:51:05 +0300
committerGitHub <[email protected]>2024-10-02 13:51:05 +0300
commitbb6b434bfbe8312ea6b418cf867738c5b1c8d133 (patch)
treedb3f18279af6a4a16db1016cf90336c1d7d735cf
parent785e5b728b8abf9a340c24d83f6123b1335982f5 (diff)
use `TTxControl` to commit transactions (#9890)
-rw-r--r--ydb/public/sdk/cpp/client/ydb_table/impl/table_client.h34
-rw-r--r--ydb/public/sdk/cpp/client/ydb_table/impl/transaction.cpp13
-rw-r--r--ydb/public/sdk/cpp/client/ydb_table/impl/transaction.h1
-rw-r--r--ydb/public/sdk/cpp/client/ydb_table/table.cpp7
-rw-r--r--ydb/public/sdk/cpp/client/ydb_table/table.h54
-rw-r--r--ydb/public/sdk/cpp/client/ydb_topic/ut/topic_to_table_ut.cpp72
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);
+}
+
}
}