This is an automated email from the ASF dual-hosted git repository.

yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/branch-4.1 by this push:
     new 30535aa6f03 branch-4.1: [fix](cloud) Repair tablet indexes missing 
db_id during transaction commit (#67855) (#68127)
30535aa6f03 is described below

commit 30535aa6f037b7d7d69a0d8114e6e69c147027c9
Author: Yixuan Wang <[email protected]>
AuthorDate: Thu Sep 17 21:58:09 2026 +0800

    branch-4.1: [fix](cloud) Repair tablet indexes missing db_id during 
transaction commit (#67855) (#68127)
    
    pick: https://github.com/apache/doris/pull/67855
---
 cloud/src/common/bvars.cpp                  |   1 +
 cloud/src/common/bvars.h                    |   1 +
 cloud/src/meta-service/meta_service_txn.cpp |  86 ++++++++--
 cloud/test/txn_lazy_commit_test.cpp         | 257 +++++++++++++++++++++++++++-
 4 files changed, 321 insertions(+), 24 deletions(-)

diff --git a/cloud/src/common/bvars.cpp b/cloud/src/common/bvars.cpp
index 004dd1a8cc2..ff1c832f1d6 100644
--- a/cloud/src/common/bvars.cpp
+++ b/cloud/src/common/bvars.cpp
@@ -119,6 +119,7 @@ bvar::Adder<int64_t> 
g_bvar_ms_rate_limit_trigger_ms_resource(
         "ms", "rate_limit_trigger_ms_resource");
 bvar::Adder<int64_t> g_bvar_ms_rate_limit_trigger_test_injection(
         "ms", "rate_limit_trigger_test_injection");
+bvar::Adder<int64_t> g_bvar_ms_repair_tablet_index("ms", 
"repair_tablet_index");
 bvar::Status<int64_t> 
g_bvar_ms_cpu_usage_percent("ms_process_cpu_usage_percent", -1);
 bvar::Status<int64_t> 
g_bvar_ms_memory_usage_percent("ms_process_memory_usage_percent", -1);
 bvar::Adder<int64_t> g_bvar_update_delete_bitmap_fail_counter;
diff --git a/cloud/src/common/bvars.h b/cloud/src/common/bvars.h
index b488283670e..7a02be09ae0 100644
--- a/cloud/src/common/bvars.h
+++ b/cloud/src/common/bvars.h
@@ -627,6 +627,7 @@ extern bvar::Adder<int64_t> 
g_bvar_ms_rate_limit_trigger_fdb_cluster;
 extern bvar::Adder<int64_t> g_bvar_ms_rate_limit_trigger_fdb_client_thread;
 extern bvar::Adder<int64_t> g_bvar_ms_rate_limit_trigger_ms_resource;
 extern bvar::Adder<int64_t> g_bvar_ms_rate_limit_trigger_test_injection;
+extern bvar::Adder<int64_t> g_bvar_ms_repair_tablet_index;
 extern bvar::Status<int64_t> g_bvar_ms_cpu_usage_percent;
 extern bvar::Status<int64_t> g_bvar_ms_memory_usage_percent;
 extern bvar::Adder<int64_t> g_bvar_update_delete_bitmap_fail_counter;
diff --git a/cloud/src/meta-service/meta_service_txn.cpp 
b/cloud/src/meta-service/meta_service_txn.cpp
index 33df289d70a..d08a1854b0f 100644
--- a/cloud/src/meta-service/meta_service_txn.cpp
+++ b/cloud/src/meta-service/meta_service_txn.cpp
@@ -18,12 +18,14 @@
 #include <gen_cpp/cloud.pb.h>
 #include <gen_cpp/olap_file.pb.h>
 
+#include <algorithm>
 #include <chrono>
 #include <cstdint>
 #include <limits>
 #include <ranges>
 #include <tuple>
 
+#include "common/bvars.h"
 #include "common/config.h"
 #include "common/logging.h"
 #include "common/stats.h"
@@ -48,6 +50,10 @@ namespace doris::cloud {
 
 static constexpr std::string_view kMetaSyncPointDummyKey = 
"__meta_service_sync_point_dummy_key__";
 
+void repair_tablet_index(std::shared_ptr<TxnKv>& txn_kv, MetaServiceCode& 
code, std::string& msg,
+                         const std::string& instance_id, int64_t db_id, 
int64_t txn_id,
+                         const std::vector<int64_t>& tablet_ids, bool 
is_versioned_write);
+
 struct TableStats {
     int64_t updated_row_count = 0;
 
@@ -1647,6 +1653,29 @@ void MetaServiceImpl::commit_txn_immediately(
             }
         }
 
+        std::vector<int64_t> repair_tablet_ids;
+        for (const auto& [tablet_id, tablet_idx] : tablet_ids) {
+            if (!tablet_idx.has_db_id()) {
+                repair_tablet_ids.push_back(tablet_id);
+            }
+        }
+        bool need_repair_tablet_idx = !repair_tablet_ids.empty();
+
+        
TEST_SYNC_POINT_CALLBACK("commit_txn_immediately::need_repair_tablet_idx",
+                                 &need_repair_tablet_idx);
+        if (need_repair_tablet_idx) {
+            stats.get_bytes += txn->get_bytes();
+            stats.get_counter += txn->num_get_keys();
+            txn.reset();
+            repair_tablet_index(txn_kv_, code, msg, instance_id, db_id, 
txn_id, repair_tablet_ids,
+                                is_versioned_write);
+            if (code != MetaServiceCode::OK) {
+                LOG(WARNING) << "repair_tablet_index failed, txn_id=" << 
txn_id << " code=" << code;
+                return;
+            }
+            continue;
+        }
+
         std::unordered_map<int64_t, std::tuple<int64_t, int64_t>> 
partition_indexes;
         for (auto& [_, i] : tmp_rowsets_meta) {
             int64_t tablet_id = i.tablet_id();
@@ -2083,18 +2112,17 @@ void MetaServiceImpl::commit_txn_immediately(
 
 // rewrite TabletIndexPB for fill db_id, in case of historical reasons
 // TabletIndexPB missing db_id
-void repair_tablet_index(
-        std::shared_ptr<TxnKv>& txn_kv, MetaServiceCode& code, std::string& 
msg,
-        const std::string& instance_id, int64_t db_id, int64_t txn_id,
-        const std::vector<std::pair<std::string, doris::RowsetMetaCloudPB>>& 
tmp_rowsets_meta,
-        bool is_versioned_write) {
+void repair_tablet_index(std::shared_ptr<TxnKv>& txn_kv, MetaServiceCode& 
code, std::string& msg,
+                         const std::string& instance_id, int64_t db_id, 
int64_t txn_id,
+                         const std::vector<int64_t>& tablet_ids, bool 
is_versioned_write) {
     std::stringstream ss;
     std::vector<std::string> tablet_idx_keys;
-    for (auto& [_, i] : tmp_rowsets_meta) {
-        tablet_idx_keys.push_back(meta_tablet_idx_key({instance_id, 
i.tablet_id()}));
+    for (int64_t tablet_id : tablet_ids) {
+        tablet_idx_keys.push_back(meta_tablet_idx_key({instance_id, 
tablet_id}));
     }
 
     for (size_t i = 0; i < tablet_idx_keys.size(); i += 
config::max_tablet_index_num_per_batch) {
+        int64_t repaired_tablet_index_num = 0;
         size_t end = (i + config::max_tablet_index_num_per_batch) > 
tablet_idx_keys.size()
                              ? tablet_idx_keys.size()
                              : i + config::max_tablet_index_num_per_batch;
@@ -2156,6 +2184,7 @@ void repair_tablet_index(
                     return;
                 }
                 txn->put(sub_tablet_idx_keys[j], idx_val);
+                ++repaired_tablet_index_num;
                 LOG(INFO) << " repair tablet index txn_id=" << txn_id
                           << " tablet_idx_pb:" << 
tablet_idx_pb.ShortDebugString()
                           << " key=" << hex(sub_tablet_idx_keys[j]);
@@ -2185,6 +2214,7 @@ void repair_tablet_index(
             LOG(WARNING) << msg;
             return;
         }
+        g_bvar_ms_repair_tablet_index << repaired_tablet_index_num;
     }
     code = MetaServiceCode::OK;
 }
@@ -2257,9 +2287,13 @@ void MetaServiceImpl::commit_txn_eventually(
             }
         }
 
-        bool need_repair_tablet_idx =
-                std::any_of(tablet_ids.begin(), tablet_ids.end(),
-                            [](const auto& pair) { return 
!pair.second.has_db_id(); });
+        std::vector<int64_t> repair_tablet_ids;
+        for (const auto& [tablet_id, tablet_idx] : tablet_ids) {
+            if (!tablet_idx.has_db_id()) {
+                repair_tablet_ids.push_back(tablet_id);
+            }
+        }
+        bool need_repair_tablet_idx = !repair_tablet_ids.empty();
 
         
TEST_SYNC_POINT_CALLBACK("commit_txn_eventually::need_repair_tablet_idx",
                                  &need_repair_tablet_idx);
@@ -2267,7 +2301,7 @@ void MetaServiceImpl::commit_txn_eventually(
             stats.get_bytes += txn->get_bytes();
             stats.get_counter += txn->num_get_keys();
             txn.reset();
-            repair_tablet_index(txn_kv_, code, msg, instance_id, db_id, 
txn_id, tmp_rowsets_meta,
+            repair_tablet_index(txn_kv_, code, msg, instance_id, db_id, 
txn_id, repair_tablet_ids,
                                 is_versioned_write);
             if (code != MetaServiceCode::OK) {
                 LOG(WARNING) << "repair_tablet_index failed, txn_id=" << 
txn_id << " code=" << code;
@@ -2825,10 +2859,11 @@ void MetaServiceImpl::commit_txn_with_sub_txn(const 
CommitTxnRequest* request,
         std::unordered_map<int64_t, TabletIndexPB> tablet_ids;
         std::vector<int64_t> acquired_tablet_ids;
         for (const auto& [_, tmp_rowsets_meta] : sub_txn_to_tmp_rowsets_meta) {
-            for (const auto& [_, i] : tmp_rowsets_meta) {
-                acquired_tablet_ids.push_back(i.tablet_id());
+            for (const auto& [_, rowset_meta] : tmp_rowsets_meta) {
+                acquired_tablet_ids.push_back(rowset_meta.tablet_id());
             }
         }
+
         if (!is_versioned_read) {
             // Read tablet indexes in batch.
             std::tie(code, msg) =
@@ -2837,8 +2872,7 @@ void MetaServiceImpl::commit_txn_with_sub_txn(const 
CommitTxnRequest* request,
                 return;
             }
         } else {
-            TxnErrorCode err =
-                    meta_reader.get_tablet_indexes(txn.get(), 
acquired_tablet_ids, &tablet_ids);
+            err = meta_reader.get_tablet_indexes(txn.get(), 
acquired_tablet_ids, &tablet_ids);
             if (err != TxnErrorCode::TXN_OK) {
                 code = cast_as<ErrCategory::READ>(err);
                 msg = fmt::format("failed to get tablet indexes, err={}", err);
@@ -2847,6 +2881,28 @@ void MetaServiceImpl::commit_txn_with_sub_txn(const 
CommitTxnRequest* request,
             }
         }
 
+        std::vector<int64_t> repair_tablet_ids;
+        for (const auto& [tablet_id, tablet_idx] : tablet_ids) {
+            if (!tablet_idx.has_db_id()) {
+                repair_tablet_ids.push_back(tablet_id);
+            }
+        }
+        bool need_repair_tablet_idx = !repair_tablet_ids.empty();
+        
TEST_SYNC_POINT_CALLBACK("commit_txn_with_sub_txn::need_repair_tablet_idx",
+                                 &need_repair_tablet_idx);
+        if (need_repair_tablet_idx) {
+            stats.get_bytes += txn->get_bytes();
+            stats.get_counter += txn->num_get_keys();
+            txn.reset();
+            repair_tablet_index(txn_kv_, code, msg, instance_id, db_id, 
txn_id, repair_tablet_ids,
+                                is_versioned_write);
+            if (code != MetaServiceCode::OK) {
+                LOG(WARNING) << "repair_tablet_index failed, txn_id=" << 
txn_id << " code=" << code;
+                return;
+            }
+            continue;
+        }
+
         // {table/partition} -> version
         std::unordered_map<int64_t, int64_t> new_versions;
         std::unordered_map<int64_t, std::tuple<int64_t, int64_t>> 
partition_indexes;
diff --git a/cloud/test/txn_lazy_commit_test.cpp 
b/cloud/test/txn_lazy_commit_test.cpp
index 41374e67a75..7d73e72ac25 100644
--- a/cloud/test/txn_lazy_commit_test.cpp
+++ b/cloud/test/txn_lazy_commit_test.cpp
@@ -51,11 +51,9 @@
 using namespace doris::cloud;
 
 namespace doris::cloud {
-void repair_tablet_index(
-        std::shared_ptr<TxnKv>& txn_kv, MetaServiceCode& code, std::string& 
msg,
-        const std::string& instance_id, int64_t db_id, int64_t txn_id,
-        const std::vector<std::pair<std::string, doris::RowsetMetaCloudPB>>& 
tmp_rowsets_meta,
-        bool is_versioned_write);
+void repair_tablet_index(std::shared_ptr<TxnKv>& txn_kv, MetaServiceCode& 
code, std::string& msg,
+                         const std::string& instance_id, int64_t db_id, 
int64_t txn_id,
+                         const std::vector<int64_t>& tablet_ids, bool 
is_versioned_write);
 };
 
 static std::shared_ptr<TxnKv> txn_kv;
@@ -538,7 +536,11 @@ TEST(TxnLazyCommitTest, RepairTabletIndexTest) {
 
     MetaServiceCode code = MetaServiceCode::UNDEFINED_ERR;
     std::string msg;
-    repair_tablet_index(txn_kv, code, msg, mock_instance, db_id, txn_id, 
tmp_rowsets_meta, false);
+    std::vector<int64_t> tablet_ids;
+    for (int i = 0; i < 2001; ++i) {
+        tablet_ids.push_back(tablet_id_base + i);
+    }
+    repair_tablet_index(txn_kv, code, msg, mock_instance, db_id, txn_id, 
tablet_ids, false);
     ASSERT_EQ(code, MetaServiceCode::OK);
 
     {
@@ -1128,7 +1130,7 @@ TEST(TxnLazyCommitVersionedReadTest, 
DISABLED_CommitTxnEventuallyWithoutDbIdTest
     }
 }
 
-TEST(TxnLazyCommitTest, CommitTxnImmediatelyTest) {
+TEST(TxnLazyCommitTest, CommitTxnImmediatelyRepairTabletIndexTest) {
     auto txn_kv = get_mem_txn_kv();
 
     int64_t db_id = 983153141;
@@ -1136,8 +1138,18 @@ TEST(TxnLazyCommitTest, CommitTxnImmediatelyTest) {
     int64_t index_id = 80124;
     int64_t partition_id = 8989313;
     bool commit_txn_immediatelly_hit = false;
+    int repair_tablet_idx_count = 0;
 
     auto sp = SyncPoint::get_instance();
+    sp->set_call_back("commit_txn_immediately::need_repair_tablet_idx", 
[&](auto&& args) {
+        bool need_repair_tablet_idx = *try_any_cast<bool*>(args[0]);
+        if (repair_tablet_idx_count == 0) {
+            ASSERT_TRUE(need_repair_tablet_idx);
+        } else {
+            ASSERT_FALSE(need_repair_tablet_idx);
+        }
+        repair_tablet_idx_count++;
+    });
     sp->set_call_back("commit_txn_immediately::finish", [&](auto&& args) {
         MetaServiceCode code = *try_any_cast<MetaServiceCode*>(args[0]);
         ASSERT_EQ(code, MetaServiceCode::OK);
@@ -1192,6 +1204,7 @@ TEST(TxnLazyCommitTest, CommitTxnImmediatelyTest) {
                                  &res, nullptr);
         ASSERT_EQ(res.status().code(), MetaServiceCode::OK);
         ASSERT_TRUE(commit_txn_immediatelly_hit);
+        ASSERT_EQ(repair_tablet_idx_count, 2);
         ASSERT_TRUE(res.has_is_lazy_commit());
         ASSERT_FALSE(res.is_lazy_commit());
         ASSERT_FALSE(res.has_is_lazy_commit_incomplete());
@@ -1201,16 +1214,242 @@ TEST(TxnLazyCommitTest, CommitTxnImmediatelyTest) {
     {
         std::unique_ptr<Transaction> txn;
         ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
-        std::string mock_instance = "test_instance";
         for (int i = 0; i < config::txn_lazy_commit_rowsets_thresold; ++i) {
             int64_t tablet_id = tablet_id_base + i;
-            check_tablet_idx_without_db_id(txn, tablet_id);
+            check_tablet_idx_db_id(txn, db_id, tablet_id);
             check_tmp_rowset_not_exist(txn, tablet_id, txn_id);
             check_rowset_meta_exist(txn, tablet_id, 2);
         }
     }
 }
 
+TEST(TxnLazyCommitTest, CommitTxnImmediatelyRepairOnlyMissingDbIdTest) {
+    auto txn_kv = get_mem_txn_kv();
+    auto meta_service = get_meta_service(txn_kv, true);
+    const int64_t db_id = 983153143;
+    const int64_t table_id = 71419095;
+    const int64_t index_id = 80126;
+    const int64_t partition_id = 8989316;
+    const int64_t tablet_id_base = 31311420;
+    const int64_t bad_tablet_id = tablet_id_base + 1;
+    const std::string instance_id = "test_instance";
+
+    brpc::Controller cntl;
+    BeginTxnRequest begin_req;
+    begin_req.set_cloud_unique_id("test_cloud_unique_id");
+    auto* txn_info = begin_req.mutable_txn_info();
+    txn_info->set_db_id(db_id);
+    txn_info->set_label("test_commit_txn_repair_only_missing_db_id");
+    txn_info->add_table_ids(table_id);
+    txn_info->set_timeout_ms(36000);
+    BeginTxnResponse begin_res;
+    meta_service->begin_txn(&cntl, &begin_req, &begin_res, nullptr);
+    ASSERT_EQ(begin_res.status().code(), MetaServiceCode::OK);
+    const int64_t txn_id = begin_res.txn_id();
+
+    for (int i = 0; i < 3; ++i) {
+        const int64_t tablet_id = tablet_id_base + i;
+        if (tablet_id == bad_tablet_id) {
+            create_tablet_without_db_id(meta_service.get(), table_id, 
index_id, partition_id,
+                                        tablet_id);
+        } else {
+            create_tablet_with_db_id(meta_service.get(), db_id, table_id, 
index_id, partition_id,
+                                     tablet_id);
+        }
+        auto rowset = create_rowset(txn_id, tablet_id, index_id, partition_id);
+        CreateRowsetResponse res;
+        prepare_rowset(meta_service.get(), rowset, res);
+        ASSERT_EQ(res.status().code(), MetaServiceCode::OK);
+        commit_rowset(meta_service.get(), rowset, res);
+        ASSERT_EQ(res.status().code(), MetaServiceCode::OK);
+    }
+
+    std::vector<std::string> original_index_values(3);
+    {
+        std::unique_ptr<Transaction> txn;
+        ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+        check_tablet_idx_without_db_id(txn, bad_tablet_id);
+        for (int i = 0; i < 3; ++i) {
+            auto key = meta_tablet_idx_key({instance_id, tablet_id_base + i});
+            ASSERT_EQ(txn->get(key, &original_index_values[i]), 
TxnErrorCode::TXN_OK);
+        }
+    }
+
+    std::vector<bool> repair_required;
+    auto sp = SyncPoint::get_instance();
+    sp->set_call_back("commit_txn_immediately::need_repair_tablet_idx", 
[&](auto&& args) {
+        repair_required.push_back(*try_any_cast<bool*>(args[0]));
+    });
+    sp->enable_processing();
+    DORIS_CLOUD_DEFER {
+        sp->clear_all_call_backs();
+        sp->clear_trace();
+        sp->disable_processing();
+    };
+
+    const auto repair_count_before = g_bvar_ms_repair_tablet_index.get_value();
+    CommitTxnRequest commit_req;
+    commit_req.set_cloud_unique_id("test_cloud_unique_id");
+    commit_req.set_db_id(db_id);
+    commit_req.set_txn_id(txn_id);
+    commit_req.set_is_2pc(false);
+    commit_req.set_enable_txn_lazy_commit(false);
+    CommitTxnResponse commit_res;
+    meta_service->commit_txn(&cntl, &commit_req, &commit_res, nullptr);
+    ASSERT_EQ(commit_res.status().code(), MetaServiceCode::OK);
+    ASSERT_FALSE(commit_res.is_lazy_commit());
+    ASSERT_EQ(repair_required, (std::vector<bool> {true, false}));
+    ASSERT_EQ(g_bvar_ms_repair_tablet_index.get_value() - repair_count_before, 
1);
+
+    std::unique_ptr<Transaction> txn;
+    ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+    for (int i = 0; i < 3; ++i) {
+        const int64_t tablet_id = tablet_id_base + i;
+        check_tablet_idx_db_id(txn, db_id, tablet_id);
+        std::string value;
+        ASSERT_EQ(txn->get(meta_tablet_idx_key({instance_id, tablet_id}), 
&value),
+                  TxnErrorCode::TXN_OK);
+        if (tablet_id == bad_tablet_id) {
+            TabletIndexPB expected;
+            ASSERT_TRUE(expected.ParseFromString(original_index_values[i]));
+            expected.set_db_id(db_id);
+            ASSERT_EQ(value, expected.SerializeAsString());
+        } else {
+            ASSERT_EQ(value, original_index_values[i]);
+        }
+        check_tmp_rowset_not_exist(txn, tablet_id, txn_id);
+        check_rowset_meta_exist(txn, tablet_id, 2);
+    }
+}
+
+TEST(TxnLazyCommitTest, CommitTxnWithSubTxnRepairTabletIndexTest) {
+    auto txn_kv = get_mem_txn_kv();
+
+    int64_t db_id = 983153142;
+    int64_t table_id = 71419094;
+    int64_t index_id = 80125;
+    int64_t partition_id1 = 8989314;
+    int64_t partition_id2 = 8989315;
+    int64_t tablet_id1 = 31311415;
+    int64_t tablet_id2 = 31311416;
+    int repair_tablet_idx_count = 0;
+    bool commit_txn_with_sub_txn_hit = false;
+
+    auto sp = SyncPoint::get_instance();
+    sp->set_call_back("commit_txn_with_sub_txn::need_repair_tablet_idx", 
[&](auto&& args) {
+        bool need_repair_tablet_idx = *try_any_cast<bool*>(args[0]);
+        if (repair_tablet_idx_count == 0) {
+            ASSERT_TRUE(need_repair_tablet_idx);
+        } else {
+            ASSERT_FALSE(need_repair_tablet_idx);
+        }
+        repair_tablet_idx_count++;
+    });
+    sp->set_call_back("commit_txn_with_sub_txn::finish", [&](auto&& args) {
+        MetaServiceCode code = *try_any_cast<MetaServiceCode*>(args[0]);
+        ASSERT_EQ(code, MetaServiceCode::OK);
+        commit_txn_with_sub_txn_hit = true;
+    });
+    sp->enable_processing();
+    DORIS_CLOUD_DEFER {
+        sp->clear_all_call_backs();
+        sp->clear_trace();
+        sp->disable_processing();
+    };
+
+    auto meta_service = get_meta_service(txn_kv, true);
+    int64_t txn_id = 0;
+    {
+        brpc::Controller cntl;
+        BeginTxnRequest req;
+        req.set_cloud_unique_id("test_cloud_unique_id");
+        TxnInfoPB txn_info_pb;
+        txn_info_pb.set_db_id(db_id);
+        txn_info_pb.set_label("test_commit_txn_with_sub_txn_repair");
+        txn_info_pb.add_table_ids(table_id);
+        txn_info_pb.set_timeout_ms(36000);
+        req.mutable_txn_info()->CopyFrom(txn_info_pb);
+        BeginTxnResponse res;
+        
meta_service->begin_txn(reinterpret_cast<::google::protobuf::RpcController*>(&cntl),
 &req,
+                                &res, nullptr);
+        ASSERT_EQ(res.status().code(), MetaServiceCode::OK);
+        txn_id = res.txn_id();
+    }
+
+    create_tablet_without_db_id(meta_service.get(), table_id, index_id, 
partition_id1, tablet_id1);
+    auto tmp_rowset1 = create_rowset(txn_id, tablet_id1, index_id, 
partition_id1);
+    CreateRowsetResponse rowset_res;
+    prepare_rowset(meta_service.get(), tmp_rowset1, rowset_res);
+    ASSERT_EQ(rowset_res.status().code(), MetaServiceCode::OK);
+    commit_rowset(meta_service.get(), tmp_rowset1, rowset_res);
+    ASSERT_EQ(rowset_res.status().code(), MetaServiceCode::OK);
+
+    int64_t sub_txn_id = 0;
+    {
+        brpc::Controller cntl;
+        BeginSubTxnRequest req;
+        req.set_cloud_unique_id("test_cloud_unique_id");
+        req.set_txn_id(txn_id);
+        req.set_sub_txn_num(0);
+        req.set_db_id(db_id);
+        req.set_label("test_commit_txn_with_sub_txn_repair_sub");
+        req.mutable_table_ids()->Add(table_id);
+        BeginSubTxnResponse res;
+        
meta_service->begin_sub_txn(reinterpret_cast<::google::protobuf::RpcController*>(&cntl),
+                                    &req, &res, nullptr);
+        ASSERT_EQ(res.status().code(), MetaServiceCode::OK);
+        ASSERT_TRUE(res.has_sub_txn_id());
+        sub_txn_id = res.sub_txn_id();
+    }
+
+    create_tablet_without_db_id(meta_service.get(), table_id, index_id, 
partition_id2, tablet_id2);
+    auto tmp_rowset2 = create_rowset(sub_txn_id, tablet_id2, index_id, 
partition_id2);
+    rowset_res.Clear();
+    prepare_rowset(meta_service.get(), tmp_rowset2, rowset_res);
+    ASSERT_EQ(rowset_res.status().code(), MetaServiceCode::OK);
+    commit_rowset(meta_service.get(), tmp_rowset2, rowset_res);
+    ASSERT_EQ(rowset_res.status().code(), MetaServiceCode::OK);
+
+    {
+        brpc::Controller cntl;
+        CommitTxnRequest req;
+        req.set_cloud_unique_id("test_cloud_unique_id");
+        req.set_db_id(db_id);
+        req.set_txn_id(txn_id);
+        req.set_is_txn_load(true);
+
+        SubTxnInfo sub_txn_info1;
+        sub_txn_info1.set_sub_txn_id(txn_id);
+        sub_txn_info1.set_table_id(table_id);
+        sub_txn_info1.mutable_base_tablet_ids()->Add(tablet_id1);
+
+        SubTxnInfo sub_txn_info2;
+        sub_txn_info2.set_sub_txn_id(sub_txn_id);
+        sub_txn_info2.set_table_id(table_id);
+        sub_txn_info2.mutable_base_tablet_ids()->Add(tablet_id2);
+
+        req.mutable_sub_txn_infos()->Add(std::move(sub_txn_info1));
+        req.mutable_sub_txn_infos()->Add(std::move(sub_txn_info2));
+        CommitTxnResponse res;
+        
meta_service->commit_txn(reinterpret_cast<::google::protobuf::RpcController*>(&cntl),
 &req,
+                                 &res, nullptr);
+        ASSERT_EQ(res.status().code(), MetaServiceCode::OK);
+        ASSERT_EQ(repair_tablet_idx_count, 2);
+        ASSERT_TRUE(commit_txn_with_sub_txn_hit);
+    }
+
+    {
+        std::unique_ptr<Transaction> txn;
+        ASSERT_EQ(txn_kv->create_txn(&txn), TxnErrorCode::TXN_OK);
+        check_tablet_idx_db_id(txn, db_id, tablet_id1);
+        check_tablet_idx_db_id(txn, db_id, tablet_id2);
+        check_tmp_rowset_not_exist(txn, tablet_id1, txn_id);
+        check_tmp_rowset_not_exist(txn, tablet_id2, sub_txn_id);
+        check_rowset_meta_exist(txn, tablet_id1, 2);
+        check_rowset_meta_exist(txn, tablet_id2, 2);
+    }
+}
+
 TEST(TxnLazyCommitTest, CommitTxnEventuallyWithFailedLazyCommitTaskTest) {
     auto txn_kv = get_mem_txn_kv();
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to