From cf913701376f878e54fda737593343f536ccbf6b Mon Sep 17 00:00:00 2001 From: Luwei <814383175@qq.com> Date: Thu, 3 Sep 2026 22:28:09 +0800 Subject: [PATCH] [fix](be) Fix row binlog recovery for multi-tablet transactions ### 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 + .../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 6b580b4861802b..6419351700a7e4 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 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, 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 dea24b37de66b1..5db2159c1e418c 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("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 1f6ee1012da583..7b1908f20c3c5e 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( @@ -902,6 +903,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(); @@ -1005,6 +1007,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()); @@ -1249,6 +1254,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 0efce3d3f2eef5..3b01f94df8a3c3 100644 --- a/be/src/storage/tablet/tablet_meta.h +++ b/be/src/storage/tablet/tablet_meta.h @@ -280,6 +280,8 @@ class TabletMeta : public MetadataAdder { 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 @@ class TabletMeta : public MetadataAdder { // 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 325604f712a909..a4b65e95d6ad84 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 #include +#include +#include #include #include #include @@ -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 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, 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(); + 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 rowsets; + std::map> 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