Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 6 additions & 5 deletions be/src/storage/data_dir.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
}

Expand Down Expand Up @@ -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(
Expand Down
25 changes: 24 additions & 1 deletion be/src/storage/tablet/tablet_manager.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}
Expand Down Expand Up @@ -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);
Expand All @@ -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();
Expand Down
6 changes: 6 additions & 0 deletions be/src/storage/tablet/tablet_meta.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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();
Expand Down Expand Up @@ -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());
Expand Down Expand Up @@ -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;
Expand Down
3 changes: 3 additions & 0 deletions be/src/storage/tablet/tablet_meta.h
Original file line number Diff line number Diff line change
Expand Up @@ -280,6 +280,8 @@ class TabletMeta : public MetadataAdder<TabletMeta> {
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;
Expand Down Expand Up @@ -395,6 +397,7 @@ class TabletMeta : public MetadataAdder<TabletMeta> {
// 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;
Expand Down
131 changes: 122 additions & 9 deletions be/test/olap/rowset/group_rowset_builder_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@
#include <stdlib.h>
#include <unistd.h>

#include <array>
#include <map>
#include <memory>
#include <string>
#include <vector>
Expand All @@ -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"

Expand All @@ -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);

Expand All @@ -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() {
Expand Down Expand Up @@ -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
Loading