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