Skip to content
Merged
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
4 changes: 2 additions & 2 deletions include/column/sort_merge.h
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,7 @@ struct UnorderRow {
}
}

return true;
return false; // equal is not less-than
}
};

Expand Down Expand Up @@ -306,7 +306,7 @@ struct SingleRow : public SlotIndex {
}
}

return true;
return false; // equal is not less-than
}
private:
CacheRow* _row = nullptr;
Expand Down
2 changes: 1 addition & 1 deletion include/common/schema_factory.h
Original file line number Diff line number Diff line change
Expand Up @@ -517,7 +517,7 @@ class RangePartition : public Partition {
bool operator() (const Range& range1, const Range& range2) {
int64_t res = range1.right_value.compare(range2.right_value);
if (res == 0) {
return range1.partition_type <= range2.partition_type;
return range1.partition_type < range2.partition_type;
} else {
return res < 0;
}
Expand Down
4 changes: 2 additions & 2 deletions include/meta_server/meta_util.h
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,7 @@ struct RangeComparator {
}
const pb::Expr& pre_right_value = left.range().right_value();
const pb::Expr& next_right_value = right.range().right_value();
return compare(pre_right_value, next_right_value) <= 0;
return compare(pre_right_value, next_right_value) < 0;
}
};

Expand All @@ -156,7 +156,7 @@ struct PointerRangeComparator {
const pb::Expr& next_right_value = right->range().right_value();
int64_t res = compare(pre_right_value, next_right_value);
if (res == 0) {
return left->type() <= right->type();
return left->type() < right->type();
} else {
return res < 0;
}
Expand Down
5 changes: 4 additions & 1 deletion src/exec/dml_node.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -493,7 +493,10 @@ int DMLNode::insert_row(RuntimeState* state, SmartRecord record, bool is_update)
}
ret = vector_index_map[info.id]->insert_vector(_txn, vectors, pk_str, record);
if (ret < 0) {
DB_WARNING_STATE(state, "vector_index fail insert, index_id: %ld", info.id);
if (ret != -2) {
// -2是向量维度对不上, 减少报警
DB_WARNING_STATE(state, "vector_index fail insert, index_id: %ld", info.id);
}
return ret;
}
continue;
Expand Down
12 changes: 11 additions & 1 deletion src/exec/dual_scan_node.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -116,7 +116,17 @@ int DualScanNode::open(RuntimeState* state) {
ExecNode* join = _sub_query_node->get_node(pb::JOIN_NODE);
_sub_query_runtime_state->is_simple_select = (join == nullptr);
_sub_query_node->set_delay_fetcher_store(_delay_fetcher_store);

// 当 delay_fetcher_store=true(UNION 向量化并发执行)时,子查询内的 JoinNode 也需要感知 delay,
// 以便 index join 路径在打开 outer node 前正确设置 delay,使其进入 EXEC_ARROW_ACERO
if (_delay_fetcher_store) {
std::vector<ExecNode*> join_nodes;
_sub_query_node->get_node(pb::JOIN_NODE, join_nodes);
for (ExecNode* node : join_nodes) {
if (node != nullptr) {
node->set_delay_fetcher_store(true);
}
}
}
ret = _sub_query_node->open(_sub_query_runtime_state);
if (ret < 0) {
DB_WARNING("Fail to open _sub_query_node");
Expand Down
9 changes: 6 additions & 3 deletions src/exec/fetcher_store.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -316,7 +316,6 @@ void OnRPCDone::select_resource_insulate_read_addr(pb::RegionInfo& info,
bool& select_without_leader,
bool& resource_insulate_read) {
// offline强制隔离
// TODO:全局索引update也会有SELECT,查询不在事务会不会有问题
_state->txn_id = 0;
std::vector<std::string> valid_addrs;
if (info.learners_size() > 0) {
Expand Down Expand Up @@ -389,8 +388,12 @@ void OnRPCDone::select_addr(pb::RegionInfo& info,
&& !_state->client_conn()->user_info->resource_tag.empty()) {
insulate_resource_tag = _state->client_conn()->user_info->resource_tag;
}
if (!insulate_resource_tag.empty() && _op_type == pb::OP_SELECT) {
return select_resource_insulate_read_addr(info, addr,
if (!insulate_resource_tag.empty() && _op_type == pb::OP_SELECT
&& _state->client_conn() != nullptr
&& _state->client_conn()->query_ctx != nullptr
&& (_state->client_conn()->query_ctx->stmt_type == parser::NT_SELECT
|| _state->client_conn()->query_ctx->stmt_type == parser::NT_UNION)) {
return select_resource_insulate_read_addr(info, addr,
insulate_resource_tag, select_without_leader, resource_insulate_read);
}

Expand Down
6 changes: 6 additions & 0 deletions src/exec/join_node.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -938,6 +938,9 @@ int JoinNode::hash_join(RuntimeState* state) {
&& (!_use_index_join || state->sign_exec_type == SignExecType::SIGN_EXEC_ARROW_FORCE_NO_INDEX_JOIN)) {
return no_index_hash_join(state);
}
if (is_delay_fetcher_store()) {
_outer_node->set_delay_fetcher_store(true);
}
int ret = _outer_node->open(state);
if (ret < 0) {
DB_WARNING("ExecNode::outer table open fail");
Expand Down Expand Up @@ -1051,6 +1054,9 @@ int JoinNode::nested_loop_join(RuntimeState* state) {
&& (!_use_index_join || state->sign_exec_type == SignExecType::SIGN_EXEC_ARROW_FORCE_NO_INDEX_JOIN)) {
return no_index_hash_join(state);
}
if (is_delay_fetcher_store()) {
_outer_node->set_delay_fetcher_store(true);
}
int ret = _outer_node->open(state);
if (ret < 0) {
DB_WARNING("ExecNode:: left table open fail");
Expand Down
10 changes: 8 additions & 2 deletions src/exec/rocksdb_scan_node.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -665,7 +665,10 @@ int RocksdbScanNode::open(RuntimeState* state) {
_scan_conjuncts,
_pre_filter_conjuncts);
if (ret < 0) {
DB_FATAL("vector_index fail search, index:%ld, table:%ld", _index_info->id, _table_info->id);
if (ret != -2) {
// -2是向量维度对不上, 减少报警
DB_FATAL("vector_index fail search, index:%ld, table:%ld", _index_info->id, _table_info->id);
}
return -1;
}
if (!_vector_filter_conjuncts.empty()) {
Expand Down Expand Up @@ -1270,7 +1273,10 @@ int RocksdbScanNode::index_ddl_work(RuntimeState* state, MemRow* row) {
}
ret = vector_index_map[index_id]->insert_vector(txn, word, pk_key.data(), record);
if (ret < 0) {
DB_WARNING_STATE(state, "vector_index fail insert, index_id: %ld", index_id);
if (ret != -2) {
// -2是向量维度对不上, 减少报警
DB_WARNING_STATE(state, "vector_index fail insert, index_id: %ld", index_id);
}
return ret;
}
return 0;
Expand Down
2 changes: 1 addition & 1 deletion src/runtime/arrow_io_excutor.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -284,7 +284,7 @@ void GlobalArrowExecutor::execute(RuntimeState* state, arrow::Result<std::shared
});
{
std::unique_lock<decltype(mu)> lock(mu);
if (!done) {
while (!done) {
cond.wait(lock);
}
}
Expand Down
4 changes: 2 additions & 2 deletions src/vector_index/vector_index.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -739,7 +739,7 @@ int VectorIndex::add_to_faiss(SmartFaissIndex faiss_index,
std::vector<int64_t> cache_idxs;
if ((int)cache_vectors.size() != _dimension) {
DB_FATAL("cache_vectors.size:%lu != _dimension:%d", cache_vectors.size(), _dimension);
return -1;
return -2;
}
cache_idxs.emplace_back(cache_idx++);
std::unordered_set<int32_t> need_not_cache_fields;
Expand Down Expand Up @@ -925,7 +925,7 @@ int VectorIndex::search(
from_chars_to_float_vec(search_data, search_vector);
if ((int)search_vector.size() != _dimension) {
DB_FATAL("search_vector.size:%lu != _dimension:%d", search_vector.size(), _dimension);
return -1;
return -2;
}
int64_t k = topk;
std::vector<int64_t> flat_idxs;
Expand Down
8 changes: 6 additions & 2 deletions test/test_dms.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,11 @@ TEST(test_dms_utils, case_sql_rewrite) {

std::vector<std::string> select_fields {"k1", "k2"};
EXPECT_EQ(SqlRewrite::rewrite_to_select(stmt, select_fields, sql, rewrite_sql), 0);
EXPECT_EQ(rewrite_sql, "SELECT k1,k2 FROM tbl WHERE k3 = v3 and k4 = v4");
EXPECT_EQ(rewrite_sql, "SELECT `k1`,`k2` FROM tbl WHERE k3 = v3 and k4 = v4");

select_fields = {"k1", "1"};
EXPECT_EQ(SqlRewrite::rewrite_to_select(stmt, select_fields, sql, rewrite_sql), 0);
EXPECT_EQ(rewrite_sql, "SELECT `k1`,1 FROM tbl WHERE k3 = v3 and k4 = v4");
}
{
parser::SqlParser parser;
Expand All @@ -49,7 +53,7 @@ TEST(test_dms_utils, case_sql_rewrite) {

std::vector<std::string> select_fields {"k1", "k2"};
EXPECT_EQ(SqlRewrite::rewrite_to_select(stmt, select_fields, sql, rewrite_sql), 0);
EXPECT_EQ(rewrite_sql, "SELECT k1,k2 FROM tbl WHERE k3 = v3 and k4 = v4");
EXPECT_EQ(rewrite_sql, "SELECT `k1`,`k2` FROM tbl WHERE k3 = v3 and k4 = v4");
}
{
parser::SqlParser parser;
Expand Down
17 changes: 17 additions & 0 deletions test/test_fetcher_store.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
#include "gtest/gtest.h"
#include "schema_factory.h"
#include "fetcher_store.h"
#include "query_context.h"
#include "proto/meta.interface.pb.h"
#include "proto/store.interface.pb.h"
#include <gflags/gflags.h>
Expand Down Expand Up @@ -85,8 +86,17 @@ class FetcherStoreTest : public testing::Test {
}
}
}
void set_select_client_conn(RuntimeState& state) {
auto ctx = std::make_shared<QueryContext>();
ctx->stmt_type = parser::NT_SELECT;
auto conn = std::make_shared<NetworkSocket>();
conn->query_ctx = ctx;
_select_conns.push_back(conn);
state.set_client_conn(conn.get());
}
std::map<std::string, std::string> _instance_info;
std::map<std::string, std::string> _instance_logical_map;
std::vector<std::shared_ptr<NetworkSocket>> _select_conns;
};

// 非select一定选leader
Expand Down Expand Up @@ -245,6 +255,7 @@ TEST_F(FetcherStoreTest, test_resource_insulate_read_without_learner) {
// 选watt peer
FetcherStore fetcher_store;
RuntimeState state;
set_select_client_conn(state);
ExecNode store_request;
OnSingleRPCDone done(&fetcher_store, &state, &store_request, &region, 1, 1, 0, 0, pb::OP_SELECT, false);
for (int i = 0; i < 5; ++i) {
Expand All @@ -257,6 +268,7 @@ TEST_F(FetcherStoreTest, test_resource_insulate_read_without_learner) {
// watt peer返回not_leader读失败,选其他peer
FetcherStore fetcher_store;
RuntimeState state;
set_select_client_conn(state);
ExecNode store_request;
OnSingleRPCDone done(&fetcher_store, &state, &store_request, &region, 1, 1, 0, 0, pb::OP_SELECT, false);
done.select_addr();
Expand All @@ -279,6 +291,7 @@ TEST_F(FetcherStoreTest, test_resource_insulate_read_without_learner) {
// watt peer can not access,选其他peer
FetcherStore fetcher_store;
RuntimeState state;
set_select_client_conn(state);
ExecNode store_request;
OnSingleRPCDone done(&fetcher_store, &state, &store_request, &region, 1, 1, 0, 0, pb::OP_SELECT, false);
done.select_addr();
Expand Down Expand Up @@ -321,6 +334,7 @@ TEST_F(FetcherStoreTest, test_resource_insulate_read_with_learner) {
FLAGS_fetcher_learner_read = false;
FetcherStore fetcher_store;
RuntimeState state;
set_select_client_conn(state);
ExecNode store_request;
OnSingleRPCDone done(&fetcher_store, &state, &store_request, &region, 1, 1, 0, 0, pb::OP_SELECT, false);
for (int i = 0; i < 5; ++i) {
Expand All @@ -335,6 +349,7 @@ TEST_F(FetcherStoreTest, test_resource_insulate_read_with_learner) {
// learner方式失败选peer
FetcherStore fetcher_store;
RuntimeState state;
set_select_client_conn(state);
ExecNode store_request;
OnSingleRPCDone done(&fetcher_store, &state, &store_request, &region, 1, 1, 0, 0, pb::OP_SELECT, false);
done.select_addr();
Expand All @@ -358,6 +373,7 @@ TEST_F(FetcherStoreTest, test_resource_insulate_read_with_learner) {
// learner not access选peer, peer失败选其他peer
FetcherStore fetcher_store;
RuntimeState state;
set_select_client_conn(state);
ExecNode store_request;
OnSingleRPCDone done(&fetcher_store, &state, &store_request, &region, 1, 1, 0, 0, pb::OP_SELECT, false);
done.select_addr();
Expand Down Expand Up @@ -534,6 +550,7 @@ TEST_F(FetcherStoreTest, test_retry_later_choose_other_peer_read2) {
FLAGS_fetcher_resource_tag = "";
FetcherStore fetcher_store;
RuntimeState state;
set_select_client_conn(state);
ExecNode store_request;
OnSingleRPCDone done(&fetcher_store, &state, &store_request, &region, 1, 1, 0, 0, pb::OP_SELECT, false);
done.select_addr();
Expand Down
Loading