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
6 changes: 3 additions & 3 deletions cmake/arrow.cmake
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ SET(ARROW_ACERO_LIB "${ARROW_INSTALL_DIR}/lib/libarrow_acero.a" CACHE FILEPATH "
SET(ARROW_BUNDLED_DEP_LIB "${ARROW_INSTALL_DIR}/lib/libarrow_bundled_dependencies.a" CACHE FILEPATH "arrow dependencies." FORCE)

FILE(WRITE ${ARROW_SOURCES_DIR}/src/build.sh
"cd cpp && cmake -DCMAKE_BUILD_TYPE=release -DARROW_JEMALLOC=OFF -DARROW_BUILD_SHARED=OFF -DARROW_PARQUET=ON -DARROW_WITH_ZLIB=ON -DARROW_WITH_ZSTD=ON -DARROW_WITH_LZ4=ON -DARROW_WITH_SNAPPY=ON -DARROW_COMPUTE=ON -DARROW_ACERO=ON -DARROW_FILESYSTEM=ON -DARROW_JSON=ON -DARROW_PARQUET=ON -DARROW_BUILD_TESTS=OFF -DARROW_BUILD_STATIC=ON -DARROW_WITH_RE2=ON -DARROW_RE2_VENDORED=OFF -DRE2_HOME=${RE2_INSTALL_DIR} -DCMAKE_PREFIX_PATH=${RE2_INSTALL_DIR} && make -j4"
"#!/bin/bash\nset -ex\ncd cpp\ncmake -DCMAKE_BUILD_TYPE=release -DARROW_JEMALLOC=OFF -DARROW_BUILD_SHARED=OFF -DARROW_PARQUET=ON -DARROW_WITH_ZLIB=ON -DARROW_WITH_ZSTD=ON -DARROW_WITH_LZ4=ON -DARROW_WITH_SNAPPY=ON -DARROW_COMPUTE=ON -DARROW_ACERO=ON -DARROW_FILESYSTEM=ON -DARROW_JSON=ON -DARROW_PARQUET=ON -DARROW_BUILD_TESTS=OFF -DARROW_BUILD_STATIC=ON -DARROW_WITH_RE2=ON -DARROW_RE2_VENDORED=OFF -DRE2_HOME=${RE2_INSTALL_DIR} -DARROW_GFLAGS=OFF -DCMAKE_PREFIX_PATH=${RE2_INSTALL_DIR} -DCMAKE_CXX_FLAGS=\"-I${GFLAGS_INSTALL_DIR}/include\" && make -j4\n"
)

ExternalProject_Add(
Expand All @@ -48,7 +48,7 @@ ExternalProject_Add(
COMMAND cp -r ${ARROW_SOURCES_DIR}/src/extern_arrow/cpp/src/parquet ${ARROW_INCLUDE_DIR}/
)

ADD_DEPENDENCIES(extern_arrow zlib snappy zstd lz4 re2 protobuf rapidjson)
ADD_DEPENDENCIES(extern_arrow zlib snappy zstd lz4 re2 protobuf rapidjson gflags)
ADD_LIBRARY(arrow STATIC IMPORTED GLOBAL)
SET_PROPERTY(TARGET arrow PROPERTY IMPORTED_LOCATION ${ARROW_LIBRARIES})
ADD_LIBRARY(parquet STATIC IMPORTED GLOBAL)
Expand All @@ -57,4 +57,4 @@ ADD_LIBRARY(acero STATIC IMPORTED GLOBAL)
SET_PROPERTY(TARGET acero PROPERTY IMPORTED_LOCATION ${ARROW_ACERO_LIB})
ADD_LIBRARY(arrow_deps STATIC IMPORTED GLOBAL)
SET_PROPERTY(TARGET arrow_deps PROPERTY IMPORTED_LOCATION ${ARROW_BUNDLED_DEP_LIB})
ADD_DEPENDENCIES(arrow parquet acero arrow_deps extern_arrow)
ADD_DEPENDENCIES(arrow parquet acero arrow_deps extern_arrow )
2 changes: 1 addition & 1 deletion cmake/boost.cmake
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ SET(Boost_VERSION "106300")
SET(Boost_LIB_VERSION "1_63_0")
SET(BOOST_VER "1.63.0")
SET(BOOST_TAR "boost_1_63_0" CACHE STRING "" FORCE)
SET(BOOST_URL "https://jaist.dl.sourceforge.net/project/boost/boost/1.63.0/${BOOST_TAR}.tar.gz" CACHE STRING "" FORCE)
SET(BOOST_URL "https://sourceforge.net/projects/boost/files/boost/1.63.0/boost_1_63_0.tar.gz" CACHE STRING "" FORCE)

MESSAGE(STATUS "BOOST_TAR: ${BOOST_TAR}, BOOST_URL: ${BOOST_URL}")

Expand Down
8 changes: 4 additions & 4 deletions include/column/column_record.h
Original file line number Diff line number Diff line change
Expand Up @@ -77,11 +77,11 @@ class ColumnRecord {
const std::vector<FieldInfo>& fields, int64_t& userid) {
return encode_row_key(record_batch, record_batch->num_rows() - 1, fields, userid);
}

static std::shared_ptr<arrow::Field> make_schema(const std::string& name, arrow::Type::type type);
static ExprValue get_vectorized_value(const std::shared_ptr<arrow::Array>& array, int row_idx);
static std::shared_ptr<arrow::Array> make_array_from_exprvalue(
const pb::PrimitiveType type, const ExprValue& expr_value, const int length);

static std::shared_ptr<ColumnSchemaInfo> make_column_schema(int64_t tableid,
SmartTable table_info, SmartIndex pri_info,
const std::unordered_map<int32_t, FieldInfo*>& field_id2info_map);

int init();

Expand Down
28 changes: 17 additions & 11 deletions include/column/file_manager.h
Original file line number Diff line number Diff line change
Expand Up @@ -187,7 +187,7 @@ class ParquetFile {

// 获取每个parquet文件中符合条件的rowgroup和rowranges
int get_qualified_rowgroup_and_rowranges(
const std::vector<pb::PossibleIndex::Range>& key_ranges,
const pb::PossibleIndex& possible_index,
std::vector<int>& rowgroup_indices,
std::vector<std::vector<std::pair<int64_t, int64_t>>>& rowranges);

Expand Down Expand Up @@ -229,7 +229,8 @@ class ParquetFile {

::arrow::Status GetRecordBatchReader(std::unique_ptr<::arrow::RecordBatchReader>* out);
static bool check_interval_overlapped(
const pb::PossibleIndex::Range& index_range, const std::string& file_start_key, const std::string& file_end_key);
const pb::PossibleIndex::Range& index_range, bool is_eq, bool is_left_open, bool is_right_open,
const std::string& file_start_key, const std::string& file_end_key);

private:
std::shared_ptr<ColumnFileInfo> _file_info;
Expand All @@ -243,29 +244,29 @@ class ParquetFile {
};

struct ParquetFileReaderOptions {
bool need_order_info = true;
int64_t raftindex = 0;
std::shared_ptr<ColumnSchemaInfo> schema_info = nullptr;
std::shared_ptr<ColumnFileInfo> file_info = nullptr;
pb::PossibleIndex* pos_index = nullptr;
std::map<std::string, FieldInfo> lower_short_name_fields;
std::shared_ptr<arrow::Schema> schema = nullptr;
};

class ParquetFileReader : public ::arrow::RecordBatchReader {
public:
ParquetFileReader(ParquetFileReaderOptions& options) : _options(options) {
_parquet_file = std::make_shared<ParquetFile>(_options.file_info);
}
ParquetFileReader(ParquetFileReaderOptions& options, std::shared_ptr<ParquetFile> file) : _options(options), _parquet_file(file) { }
virtual ~ParquetFileReader() { }

int init();

std::shared_ptr<arrow::Schema> schema() const override { return nullptr; }
std::shared_ptr<arrow::Schema> schema() const override { return _options.schema; }

virtual ::arrow::Status ReadNext(std::shared_ptr<::arrow::RecordBatch>* batch) override;
private:
ParquetFileReaderOptions _options;
std::shared_ptr<ParquetFile> _parquet_file;
std::unique_ptr<::arrow::RecordBatchReader> _reader;
bool _init = false;
int64_t _raftindex = 0;
std::shared_ptr<ReadContents> _read_contents = nullptr;

};

class ParquetFileManager {
Expand All @@ -275,6 +276,11 @@ class ParquetFileManager {
return &instance;
}

void close() {
std::unique_lock<std::mutex> l(_mutex);
_lru_cache.clear();
}

bool link_file(const std::string& old_path, const std::string& new_path) {
return ::link(old_path.c_str(), new_path.c_str()) == 0;
}
Expand Down Expand Up @@ -335,7 +341,7 @@ class ColumnFileManager {
int load_snapshot(bool restart);
int pick_minor_compact_file(int64_t applied_index, int64_t& start_version);
int pick_major_compact_file(std::vector<std::shared_ptr<ColumnFileInfo>>& file_infos);
int pick_base_compact_file(std::vector<std::shared_ptr<ColumnFileInfo>>& file_infos);
int pick_base_compact_file(std::vector<std::shared_ptr<ColumnFileInfo>>& file_infos, bool only_read_base);
int finish_minor_compact(const std::shared_ptr<ColumnFileInfo>& new_file, int64_t last_max_version);
int finish_major_compact(const std::vector<std::shared_ptr<ColumnFileInfo>>& old_files,
const std::vector<std::shared_ptr<ColumnFileInfo>>& new_files, bool is_base);
Expand Down
3 changes: 2 additions & 1 deletion include/column/parquet_cache.h
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
#include "file_system.h"
#include "column_record.h"
#include "rocksdb_filesystem.h"
#include "arrow_io_excutor.h"
namespace baikaldb {
DECLARE_int64(parquet_cache_size_mb);

Expand Down Expand Up @@ -238,7 +239,7 @@ class ParquetArrowReader {

class ParquetArrowReadableFile : public ::arrow::io::RandomAccessFile {
public:
ParquetArrowReadableFile(const std::shared_ptr<ParquetArrowReader>& reader, int64_t size, ::arrow::MemoryPool* pool = ::arrow::default_memory_pool()) :
ParquetArrowReadableFile(const std::shared_ptr<ParquetArrowReader>& reader, int64_t size, ::arrow::MemoryPool* pool = GetMemoryPoolForRead()) :
_file_reader(reader), _file_size(size), _pool(pool) { }

~ParquetArrowReadableFile() override {
Expand Down
118 changes: 111 additions & 7 deletions include/column/row2column.h
Original file line number Diff line number Diff line change
Expand Up @@ -77,39 +77,143 @@ class RocksdbBaseReader : public Row2ColumnReader {
bool _init = false;
};

struct RaftLogCacheIter {
TimeCost begin_time;
int64_t commit_index = -1;
std::map<int64_t, pb::StoreReq> log_index_req_map;
};

struct RaftLogCache {
bool empty() {
return txn_id_raft_log_map.empty();
}

int64_t time() {
int64_t t = 0;
for (auto iter : txn_id_raft_log_map) {
if (t < iter.second->begin_time.get_time()) {
t = iter.second->begin_time.get_time();
}
}

return t;
}

std::map<int64_t, std::shared_ptr<RaftLogCacheIter>> txn_id_raft_log_map;
};

class RaftLogMgr {
public:
~RaftLogMgr() {}

static RaftLogMgr* get_instance() {
static RaftLogMgr _instance;
return &_instance;
}

std::shared_ptr<RaftLogCache> get_raft_log_cache(int64_t region_id) {
std::unique_lock<bthread::Mutex> l(_lock);
auto iter = _region_id_raft_log.find(region_id);
if (iter == _region_id_raft_log.end()) {
auto cache = std::make_shared<RaftLogCache>();
return cache;
} else {
auto cache = iter->second;
_region_id_raft_log.erase(iter);
return cache;
}
}

void release_raft_log_cache(int64_t region_id, std::shared_ptr<RaftLogCache> cache) {
if (!cache->empty()) {
std::unique_lock<bthread::Mutex> l(_lock);
if (_region_id_raft_log.count(region_id) > 0) {
DB_COLUMN_FATAL("region_id: %ld, cache not empty", region_id);
}
_region_id_raft_log[region_id] = cache;
}
}


private:
bthread::Mutex _lock;
std::map<int64_t, std::shared_ptr<RaftLogCache>> _region_id_raft_log;

private:
RaftLogMgr() {}
DISALLOW_COPY_AND_ASSIGN(RaftLogMgr);
};

class RaftLogReader : public Row2ColumnReader {
public:
RaftLogReader(const Row2ColOptions& options) : Row2ColumnReader(options) {
RaftLogReader(const Row2ColOptions& options) : Row2ColumnReader(options), _region_id(options.region_id) {
_column_record = std::make_shared<ColumnRecord>(_schema_info->schema_with_order_info, _options.read_batch_size);
_raft_log_cache = RaftLogMgr::get_instance()->get_raft_log_cache(_region_id);
}
virtual ~RaftLogReader() {}

virtual ~RaftLogReader() {
int64_t time_cost = 0;
if (!_raft_log_cache->txn_id_raft_log_map.empty()) {
time_cost = _raft_log_cache->time();
if (time_cost > 15 * 60 * 1000 * 1000ULL) {
DB_COLUMN_FATAL("region_id: %ld, time_cost: %ld", _region_id, time_cost);
}
}

RaftLogMgr::get_instance()->release_raft_log_cache(_region_id, _raft_log_cache);
}

int init();

int get_raft_log(int64_t start_index, int64_t end_index, uint64_t txn_id, std::map<int64_t, pb::StoreReq>& pre_reqs);

arrow::Status ReadNext(std::shared_ptr<arrow::RecordBatch>* out) override;

std::shared_ptr<::arrow::Schema> schema() const override {
return _schema_info->schema_with_order_info;
}

int delete_column_txn_log_index() {
return MetaWriter::get_instance()->delete_column_txn_log_index(_region_id, _txn_ids);
}

int64_t get_last_raft_index() const {
return _last_index;
}

int64_t row_count() const {
return _total_row_nums;
}

int64_t put_count() const {
return _put_count;
}

int64_t delete_count() const {
return _delete_count;
}

void commit(int64_t txn_id, int64_t raft_index);

void rollback(int64_t txn_id, int64_t raft_index);

void insert(int64_t txn_id, int64_t raft_index, pb::StoreReq& request);

int get(int64_t txn_id, std::map<int64_t, pb::StoreReq>& log_index_req_map);

private:
int64_t _first_index = -1;
int64_t _last_index = -1;
int64_t _skip_count = 0;
int64_t _put_count = 0;
int64_t _first_index = -1;
int64_t _last_index = -1;
int64_t _skip_count = 0;
int64_t _put_count = 0;
int64_t _delete_count = 0;
int64_t _merge_count = 0;
int64_t _merge_count = 0;
bool _init = false;
std::vector<std::shared_ptr<arrow::RecordBatch>> _batchs;
int64_t _region_id = 0;
int64_t _commited_txn_id = -1;
std::shared_ptr<RaftLogCache> _raft_log_cache;
std::vector<uint64_t> _txn_ids;
};

} // baikaldb
1 change: 1 addition & 0 deletions include/common/baikal_heartbeat.h
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ struct HeartBeatTableName {
std::string namespace_name;
std::string database;
std::string table_name;
std::set<int64_t> partition_ids;
};

struct SubTableNames {
Expand Down
Loading