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

gavinchou pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git


The following commit(s) were added to refs/heads/master by this push:
     new 798eb24b176 [fix](binlog) Fix row binlog recovery for multi-tablet 
transactions (#67508)
798eb24b176 is described below

commit 798eb24b176c6862bdc66c94bfa875b8e861341e
Author: Luwei <[email protected]>
AuthorDate: Tue Sep 8 15:00:54 2026 +0800

    [fix](binlog) Fix row binlog recovery for multi-tablet transactions (#67508)
    
    ### What problem does this PR solve?
    
    Issue Number: close #67091
    
    Related PR: None
    
    Problem Summary: A BE restart indexed row-binlog rowsets only by
    transaction ID, so multiple tablet pairs in one transaction could attach
    the wrong companion rowset. Persist each base tablet's companion ID and
    recover using both transaction and tablet IDs.
    
    ### Release note
    
    Fix incorrect row-binlog companion recovery after a BE restart for
    multi-tablet transactions.
    
    ### Check List (For Author)
    
    - Test: Unit Test
    - Added and ran
    GroupRowsetBuilderTest.recoverMultipleRowBinlogPairsInOneTxn and
    GroupRowsetBuilderTest.*
    - Behavior changed: Yes. BE restart recovery now attaches each base
    rowset to its persisted companion tablet.
    - Does this need documentation: No
---
 be/src/storage/data_dir.cpp                       |  11 +-
 be/src/storage/tablet/tablet_manager.cpp          |  25 ++++-
 be/src/storage/tablet/tablet_meta.cpp             |   6 +
 be/src/storage/tablet/tablet_meta.h               |   3 +
 be/test/olap/rowset/group_rowset_builder_test.cpp | 131 ++++++++++++++++++++--
 5 files changed, 161 insertions(+), 15 deletions(-)

diff --git a/be/src/storage/data_dir.cpp b/be/src/storage/data_dir.cpp
index 6b580b48618..6419351700a 100644
--- a/be/src/storage/data_dir.cpp
+++ b/be/src/storage/data_dir.cpp
@@ -520,11 +520,11 @@ Status DataDir::load() {
     }
 
     // Row binlog rowset is now a normal rowset under its own binlog tablet, 
loaded above.
-    // Index them by txn id so each base rowset can re-attach its paired 
binlog rowset on recovery.
-    std::map<int64_t, RowsetMetaSharedPtr> txn_id_to_row_binlog_meta;
+    // Index them by txn and tablet id so each base rowset can re-attach its 
paired binlog rowset.
+    std::map<std::pair<int64_t, int64_t>, RowsetMetaSharedPtr> 
row_binlog_metas;
     for (auto&& rowset_meta : dir_rowset_metas) {
         if (rowset_meta->is_row_binlog()) {
-            txn_id_to_row_binlog_meta[rowset_meta->txn_id()] = rowset_meta;
+            row_binlog_metas[{rowset_meta->txn_id(), 
rowset_meta->tablet_id()}] = rowset_meta;
         }
     }
 
@@ -558,8 +558,9 @@ Status DataDir::load() {
         }
 
         RowBinlogTxnInfo attach_row_binlog;
-        if (auto it = txn_id_to_row_binlog_meta.find(rowset_meta->txn_id());
-            it != txn_id_to_row_binlog_meta.end()) {
+        if (auto it = row_binlog_metas.find(
+                    {rowset_meta->txn_id(), 
tablet->tablet_meta()->binlog_tablet_id()});
+            it != row_binlog_metas.end()) {
             const RowsetMetaSharedPtr& attach_row_binlog_rowset_meta = 
it->second;
             DCHECK_EQ(attach_row_binlog_rowset_meta->rowset_state(), 
rowset_meta->rowset_state());
             TabletSharedPtr binlog_tablet = 
_engine.tablet_manager()->get_tablet(
diff --git a/be/src/storage/tablet/tablet_manager.cpp 
b/be/src/storage/tablet/tablet_manager.cpp
index dea24b37de6..5db2159c1e4 100644
--- a/be/src/storage/tablet/tablet_manager.cpp
+++ b/be/src/storage/tablet/tablet_manager.cpp
@@ -288,9 +288,11 @@ Status TabletManager::create_tablet(const 
TCreateTabletReq& request, std::vector
     // same) already exist, then just return true(an duplicate request). But if
     // tablet_id exist but with different schema_hash, return an error(report 
task will
     // eventually trigger its deletion).
+    bool tablet_exists = false;
     {
         SCOPED_TIMER(ADD_TIMER(profile, "GetTabletUnlocked"));
-        if (_get_tablet_unlocked(tablet_id) != nullptr) {
+        tablet_exists = _get_tablet_unlocked(tablet_id) != nullptr;
+        if (tablet_exists && !is_colocated_row_binlog) {
             LOG(INFO) << "success to create tablet. tablet already exist. 
tablet_id=" << tablet_id;
             return Status::OK();
         }
@@ -331,6 +333,24 @@ Status TabletManager::create_tablet(const 
TCreateTabletReq& request, std::vector
         }
     }
 
+    auto persist_row_binlog_pair = [&]() {
+        CHECK(is_colocated_row_binlog);
+        std::lock_guard base_tablet_wlock(base_tablet->get_header_lock());
+        CHECK(base_tablet->tablet_meta()->binlog_tablet_id() == 0 ||
+              base_tablet->tablet_meta()->binlog_tablet_id() == tablet_id)
+                << "base tablet " << base_tablet->tablet_id()
+                << " is already paired with row-binlog tablet "
+                << base_tablet->tablet_meta()->binlog_tablet_id() << ", new 
row-binlog tablet "
+                << tablet_id;
+        base_tablet->tablet_meta()->set_binlog_tablet_id(tablet_id);
+        base_tablet->save_meta();
+    };
+    if (tablet_exists) {
+        persist_row_binlog_pair();
+        LOG(INFO) << "success to create tablet. tablet already exist. 
tablet_id=" << tablet_id;
+        return Status::OK();
+    }
+
     TabletSharedPtr tablet = _internal_create_tablet_unlocked(
             request, is_schema_change_or_atomic_restore, 
is_colocated_row_binlog, base_tablet.get(),
             stores, profile);
@@ -339,6 +359,9 @@ Status TabletManager::create_tablet(const TCreateTabletReq& 
request, std::vector
         return Status::Error<CE_CMD_PARAMS_ERROR>("fail to create tablet. 
tablet_id={}",
                                                   request.tablet_id);
     }
+    if (is_colocated_row_binlog) {
+        persist_row_binlog_pair();
+    }
 
     LOG(INFO) << "success to create tablet. tablet_id=" << tablet_id
               << ", tablet_path=" << tablet->tablet_path();
diff --git a/be/src/storage/tablet/tablet_meta.cpp 
b/be/src/storage/tablet/tablet_meta.cpp
index 4d57e98491d..db8783c334b 100644
--- a/be/src/storage/tablet/tablet_meta.cpp
+++ b/be/src/storage/tablet/tablet_meta.cpp
@@ -269,6 +269,7 @@ TabletMeta::TabletMeta(const TabletMeta& b)
           _delete_bitmap(b._delete_bitmap),
           _binlog_config(b._binlog_config),
           _tablet_role(b._tablet_role),
+          _binlog_tablet_id(b._binlog_tablet_id),
           _compaction_policy(b._compaction_policy),
           
_time_series_compaction_goal_size_mbytes(b._time_series_compaction_goal_size_mbytes),
           _time_series_compaction_file_count_threshold(
@@ -905,6 +906,7 @@ void TabletMeta::init_from_pb(const TabletMetaPB& 
tablet_meta_pb) {
         _binlog_config = tablet_meta_pb.binlog_config();
     }
     _tablet_role = tablet_meta_pb.tablet_role();
+    _binlog_tablet_id = tablet_meta_pb.binlog_tablet_id();
     _compaction_policy = tablet_meta_pb.compaction_policy();
     _time_series_compaction_goal_size_mbytes =
             tablet_meta_pb.time_series_compaction_goal_size_mbytes();
@@ -1008,6 +1010,9 @@ void TabletMeta::to_meta_pb(TabletMetaPB* tablet_meta_pb, 
bool cloud_get_rowset_
     }
     _binlog_config.to_pb(tablet_meta_pb->mutable_binlog_config());
     tablet_meta_pb->set_tablet_role(_tablet_role);
+    if (_binlog_tablet_id > 0) {
+        tablet_meta_pb->set_binlog_tablet_id(_binlog_tablet_id);
+    }
     tablet_meta_pb->set_compaction_policy(compaction_policy());
     tablet_meta_pb->set_time_series_compaction_goal_size_mbytes(
             time_series_compaction_goal_size_mbytes());
@@ -1252,6 +1257,7 @@ bool operator==(const TabletMeta& a, const TabletMeta& b) 
{
     if (a._in_restore_mode != b._in_restore_mode) return false;
     if (a._preferred_rowset_type != b._preferred_rowset_type) return false;
     if (a._storage_policy_id != b._storage_policy_id) return false;
+    if (a._binlog_tablet_id != b._binlog_tablet_id) return false;
     if (a._compaction_policy != b._compaction_policy) return false;
     if (a._time_series_compaction_goal_size_mbytes != 
b._time_series_compaction_goal_size_mbytes)
         return false;
diff --git a/be/src/storage/tablet/tablet_meta.h 
b/be/src/storage/tablet/tablet_meta.h
index 0efce3d3f2e..3b01f94df8a 100644
--- a/be/src/storage/tablet/tablet_meta.h
+++ b/be/src/storage/tablet/tablet_meta.h
@@ -280,6 +280,8 @@ public:
         return _tablet_role == TabletRolePB::TABLET_ROLE_ROW_BINLOG;
     }
     void set_tablet_role(TabletRolePB tablet_role) { _tablet_role = 
tablet_role; }
+    int64_t binlog_tablet_id() const { return _binlog_tablet_id; }
+    void set_binlog_tablet_id(int64_t binlog_tablet_id) { _binlog_tablet_id = 
binlog_tablet_id; }
 
     void set_compaction_policy(std::string compaction_policy) {
         _compaction_policy = compaction_policy;
@@ -395,6 +397,7 @@ private:
     // binlog config
     BinlogConfig _binlog_config {};
     TabletRolePB _tablet_role = TabletRolePB::TABLET_ROLE_DATA;
+    int64_t _binlog_tablet_id = 0;
 
     // meta for compaction
     std::string _compaction_policy;
diff --git a/be/test/olap/rowset/group_rowset_builder_test.cpp 
b/be/test/olap/rowset/group_rowset_builder_test.cpp
index 325604f712a..a4b65e95d6a 100644
--- a/be/test/olap/rowset/group_rowset_builder_test.cpp
+++ b/be/test/olap/rowset/group_rowset_builder_test.cpp
@@ -22,6 +22,8 @@
 #include <stdlib.h>
 #include <unistd.h>
 
+#include <array>
+#include <map>
 #include <memory>
 #include <string>
 #include <vector>
@@ -39,6 +41,7 @@
 #include "storage/storage_engine.h"
 #include "storage/tablet/tablet.h"
 #include "storage/tablet/tablet_manager.h"
+#include "storage/tablet/tablet_meta_manager.h"
 #include "storage/tablet_info.h"
 #include "testutil/creators.h"
 
@@ -47,14 +50,7 @@ namespace doris {
 static const uint32_t MAX_PATH_LEN = 1024;
 static StorageEngine* engine_ref = nullptr;
 
-static void set_up() {
-    char buffer[MAX_PATH_LEN];
-    EXPECT_NE(getcwd(buffer, MAX_PATH_LEN), nullptr);
-    config::storage_root_path = std::string(buffer) + "/data_test";
-    auto st = 
io::global_local_filesystem()->delete_directory(config::storage_root_path);
-    ASSERT_TRUE(st.ok()) << st;
-    st = 
io::global_local_filesystem()->create_directory(config::storage_root_path);
-    ASSERT_TRUE(st.ok()) << st;
+static void open_engine() {
     std::vector<StorePath> paths;
     paths.emplace_back(config::storage_root_path, -1);
 
@@ -64,10 +60,26 @@ static void set_up() {
     engine_ref = engine.get();
     Status s = engine->open();
     ASSERT_TRUE(s.ok()) << s;
+    ExecEnv::GetInstance()->set_storage_engine(std::move(engine));
+}
 
+static void set_up() {
+    char buffer[MAX_PATH_LEN];
+    EXPECT_NE(getcwd(buffer, MAX_PATH_LEN), nullptr);
+    config::storage_root_path = std::string(buffer) + "/data_test";
+    auto st = 
io::global_local_filesystem()->delete_directory(config::storage_root_path);
+    ASSERT_TRUE(st.ok()) << st;
+    st = 
io::global_local_filesystem()->create_directory(config::storage_root_path);
+    ASSERT_TRUE(st.ok()) << st;
     ExecEnv* exec_env = doris::ExecEnv::GetInstance();
     exec_env->set_memtable_memory_limiter(new MemTableMemoryLimiter());
-    exec_env->set_storage_engine(std::move(engine));
+    open_engine();
+}
+
+static void restart_engine() {
+    engine_ref = nullptr;
+    ExecEnv::GetInstance()->set_storage_engine(nullptr);
+    open_engine();
 }
 
 static void tear_down() {
@@ -167,4 +179,105 @@ TEST_F(GroupRowsetBuilderTest, buildWithRowBinlogMeta) {
     ASSERT_TRUE(res.ok());
 }
 
+TEST_F(GroupRowsetBuilderTest, recoverMultipleRowBinlogPairsInOneTxn) {
+    constexpr int64_t partition_id = 10100;
+    constexpr int64_t txn_id = 20100;
+    constexpr int64_t index_id = 30100;
+    constexpr int64_t row_binlog_index_id = 30101;
+    constexpr int32_t schema_hash = 40100;
+    constexpr int32_t row_binlog_schema_hash = 40101;
+    constexpr std::array<std::pair<int64_t, int64_t>, 2> tablet_pairs = 
{std::pair {10100, 10101},
+                                                                         
std::pair {10200, 10201}};
+
+    auto base_request = testutil::create_tablet_request(
+            0, schema_hash, partition_id, 1, TKeysType::UNIQUE_KEYS,
+            {{"k1", TPrimitiveType::INT, true}, {"v1", TPrimitiveType::INT, 
false}});
+    base_request.__set_enable_unique_key_merge_on_write(true);
+    testutil::enable_row_binlog(&base_request);
+    auto row_binlog_schema = 
testutil::create_row_binlog_tablet_schema(base_request.tablet_schema,
+                                                                       
row_binlog_schema_hash);
+
+    RuntimeProfile profile("CreateTablet");
+    for (const auto& [base_tablet_id, row_binlog_tablet_id] : tablet_pairs) {
+        base_request.tablet_id = base_tablet_id;
+        ASSERT_TRUE(engine_ref->create_tablet(base_request, &profile).ok());
+
+        auto row_binlog_request = base_request;
+        row_binlog_request.tablet_id = row_binlog_tablet_id;
+        row_binlog_request.tablet_schema = row_binlog_schema;
+        row_binlog_request.__set_base_tablet_id(base_tablet_id);
+        
row_binlog_request.__set_tablet_role(TTabletRole::TABLET_ROLE_ROW_BINLOG);
+        ASSERT_TRUE(engine_ref->create_tablet(row_binlog_request, 
&profile).ok());
+
+        auto base_tablet = 
engine_ref->tablet_manager()->get_tablet(base_tablet_id);
+        ASSERT_NE(base_tablet, nullptr);
+        TabletMetaPB in_memory_meta_pb;
+        base_tablet->tablet_meta()->to_meta_pb(&in_memory_meta_pb, false);
+        EXPECT_EQ(in_memory_meta_pb.binlog_tablet_id(), row_binlog_tablet_id);
+
+        TabletMetaSharedPtr persisted_meta = std::make_shared<TabletMeta>();
+        ASSERT_TRUE(TabletMetaManager::get_meta(base_tablet->data_dir(), 
base_tablet_id,
+                                                schema_hash, persisted_meta)
+                            .ok());
+        TabletMetaPB persisted_meta_pb;
+        persisted_meta->to_meta_pb(&persisted_meta_pb, false);
+        EXPECT_EQ(persisted_meta_pb.binlog_tablet_id(), row_binlog_tablet_id);
+    }
+
+    TDescriptorTable tdesc_tbl =
+            testutil::create_descriptor_table({{TYPE_INT, "k1", false}, 
{TYPE_INT, "v1", false}});
+    auto schema_param = testutil::create_table_schema_param(
+            tdesc_tbl, index_id, schema_hash, 
base_request.tablet_schema.columns,
+            row_binlog_index_id, row_binlog_schema_hash, 
&row_binlog_schema.columns);
+    ASSERT_NE(schema_param, nullptr);
+
+    PUniqueId load_id;
+    load_id.set_hi(0);
+    load_id.set_lo(1);
+    for (const auto& [base_tablet_id, row_binlog_tablet_id] : tablet_pairs) {
+        WriteRequest data_req;
+        data_req.tablet_id = base_tablet_id;
+        data_req.schema_hash = schema_hash;
+        data_req.txn_id = txn_id;
+        data_req.partition_id = partition_id;
+        data_req.index_id = index_id;
+        data_req.load_id = load_id;
+        data_req.table_schema_param = schema_param;
+        data_req.write_req_type = WriteRequestType::DATA;
+
+        WriteRequest row_binlog_req = data_req;
+        row_binlog_req.tablet_id = row_binlog_tablet_id;
+        row_binlog_req.index_id = row_binlog_index_id;
+        row_binlog_req.schema_hash = row_binlog_schema_hash;
+        row_binlog_req.write_req_type = WriteRequestType::ROW_BINLOG;
+
+        WriteRequest group_req = data_req;
+        group_req.write_req_type = WriteRequestType::GROUP;
+
+        GroupRowsetBuilder builder(*engine_ref, group_req, data_req, 
row_binlog_req, &profile);
+        ASSERT_TRUE(builder.init().ok());
+        ASSERT_TRUE(builder.rowset_writer()->flush().ok());
+        ASSERT_TRUE(builder.build_rowset().ok());
+        ASSERT_TRUE(builder.commit_txn().ok());
+    }
+
+    restart_engine();
+
+    std::map<TabletInfo, RowsetSharedPtr> rowsets;
+    std::map<TabletInfo, std::shared_ptr<TabletTxnInfo>> txn_infos;
+    engine_ref->txn_manager()->get_txn_related_tablets(txn_id, partition_id, 
&rowsets, &txn_infos);
+    ASSERT_EQ(txn_infos.size(), tablet_pairs.size());
+    for (const auto& [base_tablet_id, row_binlog_tablet_id] : tablet_pairs) {
+        auto base_tablet = 
engine_ref->tablet_manager()->get_tablet(base_tablet_id);
+        ASSERT_NE(base_tablet, nullptr);
+        auto txn_info = txn_infos.find(base_tablet->get_tablet_info());
+        ASSERT_NE(txn_info, txn_infos.end());
+        ASSERT_NE(txn_info->second->attach_row_binlog.tablet, nullptr);
+        ASSERT_NE(txn_info->second->attach_row_binlog.rowset, nullptr);
+        EXPECT_EQ(txn_info->second->attach_row_binlog.tablet->tablet_id(), 
row_binlog_tablet_id);
+        
EXPECT_EQ(txn_info->second->attach_row_binlog.rowset->rowset_meta()->tablet_id(),
+                  row_binlog_tablet_id);
+    }
+}
+
 } // namespace doris


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

Reply via email to