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
10 changes: 10 additions & 0 deletions be/src/information_schema/schema_tso_status_scanner.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,11 @@ std::vector<SchemaScanner::ColumnDesc> SchemaTsoStatusScanner::_s_tso_status_col
{"CURRENT_TSO", TYPE_BIGINT, sizeof(int64_t), true},
{"CURRENT_TSO_PHYSICAL_TIME", TYPE_BIGINT, sizeof(int64_t), true},
{"CURRENT_TSO_LOGICAL_COUNTER", TYPE_BIGINT, sizeof(int64_t), true},
{.name = "COMMITTED_TSO", .type = TYPE_BIGINT, .size = sizeof(int64_t), .is_null = true},
{.name = "COMMITTED_TSO_PHYSICAL_TIME",
.type = TYPE_BIGINT,
.size = sizeof(int64_t),
.is_null = true},
};

SchemaTsoStatusScanner::SchemaTsoStatusScanner()
Expand All @@ -51,6 +56,11 @@ Status SchemaTsoStatusScanner::_get_tso_status_block_from_fe() {
TNetworkAddress master_addr = ExecEnv::GetInstance()->cluster_info()->master_fe_addr;

TSchemaTableRequestParams schema_table_request_params;
std::vector<std::string> columns;
for (const auto& column : _s_tso_status_columns) {
columns.emplace_back(column.name);
}
schema_table_request_params.__set_columns_name(columns);
TFetchSchemaTableDataRequest request;
request.__set_schema_table_name(TSchemaTableName::TSO_STATUS);
request.__set_schema_table_params(schema_table_request_params);
Expand Down
34 changes: 24 additions & 10 deletions be/test/exec/schema_scanner/schema_tso_status_scanner_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,7 @@
namespace doris {
namespace {

TRow create_tso_status_row(const std::array<int64_t, 4>& values) {
TRow create_tso_status_row(const std::array<int64_t, 6>& values) {
std::vector<TCell> cells;
cells.reserve(values.size());
for (int64_t value : values) {
Expand Down Expand Up @@ -65,7 +65,7 @@ std::unique_ptr<Block> create_output_block(SchemaTsoStatusScanner* scanner) {
}

void expect_tso_status_row(const Block& block, size_t row_idx,
const std::array<int64_t, 4>& expected) {
const std::array<int64_t, 6>& expected) {
ASSERT_EQ(expected.size(), block.columns());
for (size_t column_idx = 0; column_idx < expected.size(); ++column_idx) {
const auto& column = block.get_by_position(column_idx).column;
Expand All @@ -80,11 +80,13 @@ TEST(SchemaTsoStatusScannerTest, test_create_tso_status_scanner) {
auto scanner = SchemaScanner::create(TSchemaTableType::SCH_TSO_STATUS);
ASSERT_NE(nullptr, scanner);
EXPECT_EQ(TSchemaTableType::SCH_TSO_STATUS, scanner->type());
ASSERT_EQ(4, scanner->get_column_desc().size());
ASSERT_EQ(6, scanner->get_column_desc().size());
EXPECT_STREQ("WINDOW_END_PHYSICAL_TIME", scanner->get_column_desc()[0].name);
EXPECT_STREQ("CURRENT_TSO", scanner->get_column_desc()[1].name);
EXPECT_STREQ("CURRENT_TSO_PHYSICAL_TIME", scanner->get_column_desc()[2].name);
EXPECT_STREQ("CURRENT_TSO_LOGICAL_COUNTER", scanner->get_column_desc()[3].name);
EXPECT_STREQ("COMMITTED_TSO", scanner->get_column_desc()[4].name);
EXPECT_STREQ("COMMITTED_TSO_PHYSICAL_TIME", scanner->get_column_desc()[5].name);
for (const auto& column : scanner->get_column_desc()) {
EXPECT_EQ(TYPE_BIGINT, column.type);
EXPECT_TRUE(column.is_null);
Expand Down Expand Up @@ -140,16 +142,16 @@ TEST(SchemaTsoStatusScannerTest, test_process_tso_status_result_error) {
}

TEST(SchemaTsoStatusScannerTest, test_process_tso_status_result) {
const std::array<int64_t, 4> first_row = {1000, 2000, 3000, 4000};
const std::array<int64_t, 4> second_row = {1001, 2001, 3001, 4001};
const std::array<int64_t, 6> first_row = {1000, 2000, 3000, 4000};
const std::array<int64_t, 6> second_row = {1001, 2001, 3001, 4001};
auto result = create_tso_status_result(
{create_tso_status_row(first_row), create_tso_status_row(second_row)});

SchemaTsoStatusScanner scanner;
ASSERT_TRUE(scanner._process_tso_status_result(result).ok());

ASSERT_NE(nullptr, scanner._tso_status_block);
EXPECT_EQ(4, scanner._tso_status_block->columns());
EXPECT_EQ(6, scanner._tso_status_block->columns());
EXPECT_EQ(2, scanner._tso_status_block->rows());
EXPECT_EQ(2, scanner._total_rows);
expect_tso_status_row(*scanner._tso_status_block, 0, first_row);
Expand All @@ -158,7 +160,8 @@ TEST(SchemaTsoStatusScannerTest, test_process_tso_status_result) {

TEST(SchemaTsoStatusScannerTest, test_process_tso_status_result_schema_mismatch) {
TRow invalid_row = create_tso_status_row({1000, 2000, 3000, 4000});
invalid_row.column_value.pop_back();
invalid_row.column_value.resize(
4); // Response from an older FE cannot provide a committed prefix.
auto result = create_tso_status_result({invalid_row});

SchemaTsoStatusScanner scanner;
Expand All @@ -170,6 +173,17 @@ TEST(SchemaTsoStatusScannerTest, test_process_tso_status_result_schema_mismatch)
EXPECT_EQ(0, scanner._total_rows);
}

TEST(SchemaTsoStatusScannerTest, test_unknown_committed_tso_is_null) {
TRow row = create_tso_status_row({1000, 2000, 3000, 4000, 0, 0});
row.column_value[4].__set_isNull(true);
row.column_value[5].__set_isNull(true);
SchemaTsoStatusScanner scanner;
ASSERT_TRUE(scanner._process_tso_status_result(create_tso_status_result({row})).ok());
EXPECT_FALSE(scanner._tso_status_block->get_by_position(3).column->is_null_at(0));
EXPECT_TRUE(scanner._tso_status_block->get_by_position(4).column->is_null_at(0));
EXPECT_TRUE(scanner._tso_status_block->get_by_position(5).column->is_null_at(0));
}

TEST(SchemaTsoStatusScannerTest, test_get_next_block_empty_result) {
MockRuntimeState state;
SchemaScannerParam param;
Expand All @@ -188,9 +202,9 @@ TEST(SchemaTsoStatusScannerTest, test_get_next_block_empty_result) {
}

TEST(SchemaTsoStatusScannerTest, test_get_next_block_in_batches) {
const std::array<int64_t, 4> first_row = {1000, 2000, 3000, 4000};
const std::array<int64_t, 4> second_row = {1001, 2001, 3001, 4001};
const std::array<int64_t, 4> third_row = {1002, 2002, 3002, 4002};
const std::array<int64_t, 6> first_row = {1000, 2000, 3000, 4000};
const std::array<int64_t, 6> second_row = {1001, 2001, 3001, 4001};
const std::array<int64_t, 6> third_row = {1002, 2002, 3002, 4002};

MockRuntimeState state;
state._batch_size = 2;
Expand Down
8 changes: 8 additions & 0 deletions cloud/src/common/bvars.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,8 @@ BvarLatencyRecorderWithTag g_bvar_ms_create_meta_sync_point("ms", "create_meta_s
BvarLatencyRecorderWithTag g_bvar_ms_begin_sub_txn("ms", "begin_sub_txn");
BvarLatencyRecorderWithTag g_bvar_ms_abort_sub_txn("ms", "abort_sub_txn");
BvarLatencyRecorderWithTag g_bvar_ms_check_txn_conflict("ms", "check_txn_conflict");
BvarLatencyRecorderWithTag g_bvar_ms_get_tso_recovery_transactions("ms", "get_tso_recovery_transactions");
BvarLatencyRecorderWithTag g_bvar_ms_advance_tso_fence("ms", "advance_tso_fence");
BvarLatencyRecorderWithTag g_bvar_ms_abort_txn_with_coordinator("ms", "abort_txn_with_coordinator");
BvarLatencyRecorderWithTag g_bvar_ms_get_prepare_txn_by_coordinator("ms", "get_prepare_txn_by_coordinator");
BvarLatencyRecorderWithTag g_bvar_ms_clean_txn_label("ms", "clean_txn_label");
Expand Down Expand Up @@ -500,6 +502,9 @@ mBvarInt64Adder g_bvar_rpc_kv_abort_txn_with_coordinator_get_counter("rpc_kv_abo
mBvarInt64Adder g_bvar_rpc_kv_get_prepare_txn_by_coordinator_get_counter("rpc_kv_get_prepare_txn_by_coordinator_get_counter",{"instance_id"});
// check_txn_conflict
mBvarInt64Adder g_bvar_rpc_kv_check_txn_conflict_get_counter("rpc_kv_check_txn_conflict_get_counter",{"instance_id"});
mBvarInt64Adder g_bvar_rpc_kv_get_tso_recovery_transactions_get_counter("rpc_kv_get_tso_recovery_transactions_get_counter",{"instance_id"});
mBvarInt64Adder g_bvar_rpc_kv_advance_tso_fence_get_counter("rpc_kv_advance_tso_fence_get_counter",{"instance_id"});
mBvarInt64Adder g_bvar_rpc_kv_advance_tso_fence_put_counter("rpc_kv_advance_tso_fence_put_counter",{"instance_id"});
// clean_txn_label
mBvarInt64Adder g_bvar_rpc_kv_clean_txn_label_get_counter("rpc_kv_clean_txn_label_get_counter",{"instance_id"});
mBvarInt64Adder g_bvar_rpc_kv_clean_txn_label_put_counter("rpc_kv_clean_txn_label_put_counter",{"instance_id"});
Expand Down Expand Up @@ -710,6 +715,9 @@ mBvarInt64Adder g_bvar_rpc_kv_abort_txn_with_coordinator_get_bytes("rpc_kv_abort
mBvarInt64Adder g_bvar_rpc_kv_get_prepare_txn_by_coordinator_get_bytes("rpc_kv_get_prepare_txn_by_coordinator_get_bytes",{"instance_id"});
// check_txn_conflict
mBvarInt64Adder g_bvar_rpc_kv_check_txn_conflict_get_bytes("rpc_kv_check_txn_conflict_get_bytes",{"instance_id"});
mBvarInt64Adder g_bvar_rpc_kv_get_tso_recovery_transactions_get_bytes("rpc_kv_get_tso_recovery_transactions_get_bytes",{"instance_id"});
mBvarInt64Adder g_bvar_rpc_kv_advance_tso_fence_get_bytes("rpc_kv_advance_tso_fence_get_bytes",{"instance_id"});
mBvarInt64Adder g_bvar_rpc_kv_advance_tso_fence_put_bytes("rpc_kv_advance_tso_fence_put_bytes",{"instance_id"});
// clean_txn_label
mBvarInt64Adder g_bvar_rpc_kv_clean_txn_label_get_bytes("rpc_kv_clean_txn_label_get_bytes",{"instance_id"});
mBvarInt64Adder g_bvar_rpc_kv_clean_txn_label_put_bytes("rpc_kv_clean_txn_label_put_bytes",{"instance_id"});
Expand Down
8 changes: 8 additions & 0 deletions cloud/src/common/bvars.h
Original file line number Diff line number Diff line change
Expand Up @@ -553,6 +553,8 @@ extern BvarLatencyRecorderWithTag g_bvar_ms_get_txn;
extern BvarLatencyRecorderWithTag g_bvar_ms_get_current_max_txn_id;
extern BvarLatencyRecorderWithTag g_bvar_ms_create_meta_sync_point;
extern BvarLatencyRecorderWithTag g_bvar_ms_check_txn_conflict;
extern BvarLatencyRecorderWithTag g_bvar_ms_get_tso_recovery_transactions;
extern BvarLatencyRecorderWithTag g_bvar_ms_advance_tso_fence;
extern BvarLatencyRecorderWithTag g_bvar_ms_abort_txn_with_coordinator;
extern BvarLatencyRecorderWithTag g_bvar_ms_get_prepare_txn_by_coordinator;
extern BvarLatencyRecorderWithTag g_bvar_ms_begin_sub_txn;
Expand Down Expand Up @@ -907,6 +909,9 @@ extern mBvarInt64Adder g_bvar_rpc_kv_abort_sub_txn_put_counter;
extern mBvarInt64Adder g_bvar_rpc_kv_abort_txn_with_coordinator_get_counter;
extern mBvarInt64Adder g_bvar_rpc_kv_get_prepare_txn_by_coordinator_get_counter;
extern mBvarInt64Adder g_bvar_rpc_kv_check_txn_conflict_get_counter;
extern mBvarInt64Adder g_bvar_rpc_kv_get_tso_recovery_transactions_get_counter;
extern mBvarInt64Adder g_bvar_rpc_kv_advance_tso_fence_get_counter;
extern mBvarInt64Adder g_bvar_rpc_kv_advance_tso_fence_put_counter;
extern mBvarInt64Adder g_bvar_rpc_kv_clean_txn_label_get_counter;
extern mBvarInt64Adder g_bvar_rpc_kv_clean_txn_label_put_counter;
extern mBvarInt64Adder g_bvar_rpc_kv_clean_txn_label_del_counter;
Expand Down Expand Up @@ -1050,6 +1055,9 @@ extern mBvarInt64Adder g_bvar_rpc_kv_abort_sub_txn_put_bytes;
extern mBvarInt64Adder g_bvar_rpc_kv_abort_txn_with_coordinator_get_bytes;
extern mBvarInt64Adder g_bvar_rpc_kv_get_prepare_txn_by_coordinator_get_bytes;
extern mBvarInt64Adder g_bvar_rpc_kv_check_txn_conflict_get_bytes;
extern mBvarInt64Adder g_bvar_rpc_kv_get_tso_recovery_transactions_get_bytes;
extern mBvarInt64Adder g_bvar_rpc_kv_advance_tso_fence_get_bytes;
extern mBvarInt64Adder g_bvar_rpc_kv_advance_tso_fence_put_bytes;
extern mBvarInt64Adder g_bvar_rpc_kv_clean_txn_label_get_bytes;
extern mBvarInt64Adder g_bvar_rpc_kv_clean_txn_label_put_bytes;
extern mBvarInt64Adder g_bvar_rpc_kv_clean_txn_label_del_bytes;
Expand Down
23 changes: 23 additions & 0 deletions cloud/src/meta-service/meta_service.h
Original file line number Diff line number Diff line change
Expand Up @@ -145,6 +145,15 @@ class MetaServiceImpl : public cloud::MetaService {
CheckTxnConflictResponse* response,
::google::protobuf::Closure* done) override;

void get_tso_recovery_transactions(::google::protobuf::RpcController* controller,
const GetTsoRecoveryTransactionsRequest* request,
GetTsoRecoveryTransactionsResponse* response,
::google::protobuf::Closure* done) override;

void advance_tso_fence(::google::protobuf::RpcController* controller,
const AdvanceTsoFenceRequest* request, AdvanceTsoFenceResponse* response,
::google::protobuf::Closure* done) override;

void abort_txn_with_coordinator(::google::protobuf::RpcController* controller,
const AbortTxnWithCoordinatorRequest* request,
AbortTxnWithCoordinatorResponse* response,
Expand Down Expand Up @@ -619,6 +628,20 @@ class MetaServiceProxy final : public MetaService {
call_impl(&cloud::MetaService::check_txn_conflict, controller, request, response, done);
}

void get_tso_recovery_transactions(::google::protobuf::RpcController* controller,
const GetTsoRecoveryTransactionsRequest* request,
GetTsoRecoveryTransactionsResponse* response,
::google::protobuf::Closure* done) override {
call_impl(&cloud::MetaService::get_tso_recovery_transactions, controller, request, response,
done);
}

void advance_tso_fence(::google::protobuf::RpcController* controller,
const AdvanceTsoFenceRequest* request, AdvanceTsoFenceResponse* response,
::google::protobuf::Closure* done) override {
call_impl(&cloud::MetaService::advance_tso_fence, controller, request, response, done);
}

void abort_txn_with_coordinator(::google::protobuf::RpcController* controller,
const AbortTxnWithCoordinatorRequest* request,
AbortTxnWithCoordinatorResponse* response,
Expand Down
5 changes: 5 additions & 0 deletions cloud/src/meta-service/meta_service_helper.h
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,11 @@ inline std::pair<MetaServiceCode, std::string> resolve_response_code_and_msg(Met
"[TXN_ALREADY_COMMITED will be converted to code=UNDEFINED_ERR for old version "
"clients]";
return {MetaServiceCode::UNDEFINED_ERR, std::move(msg)};
case MetaServiceCode::TXN_COMMIT_TSO_FENCED:
msg += std::string((msg.empty() ? "" : ", ")) +
"[TXN_COMMIT_TSO_FENCED will be converted to code=UNDEFINED_ERR for old "
"version clients]";
return {MetaServiceCode::UNDEFINED_ERR, std::move(msg)};
default:
return {code, std::move(msg)};
}
Expand Down
Loading
Loading