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
156 changes: 92 additions & 64 deletions cpp/src/reader/chunk_reader.cc
Original file line number Diff line number Diff line change
Expand Up @@ -315,7 +315,7 @@ int ChunkReader::skip_cur_page() {
}

int ChunkReader::decode_cur_page_data(TsBlock*& ret_tsblock, Filter* filter,
PageArena& pa) {
PageArena& pa, int* row_offset) {
int ret = E_OK;

// Step 1: make sure we load the whole page data in @in_stream_
Expand Down Expand Up @@ -396,8 +396,8 @@ int ChunkReader::decode_cur_page_data(TsBlock*& ret_tsblock, Filter* filter,
// ret = decode_tv_buf_into_tsblock(time_buf, value_buf, time_buf_size,
// value_buf_size, ret_tsblock,
// filter);
ret = decode_tv_buf_into_tsblock_by_datatype(time_in_, value_in_,
ret_tsblock, filter, &pa);
ret = decode_tv_buf_into_tsblock_by_datatype(
time_in_, value_in_, ret_tsblock, filter, &pa, row_offset);
// if we return during @decode_tv_buf_into_tsblock, we should keep
// @uncompressed_buf_ valid until all TV pairs are decoded.
if (ret != E_OVERFLOW) {
Expand All @@ -414,29 +414,32 @@ int ChunkReader::decode_cur_page_data(TsBlock*& ret_tsblock, Filter* filter,
return ret;
}

#define DECODE_TYPED_TV_INTO_TSBLOCK(CppType, ReadType, time_in, value_in, \
row_appender) \
do { \
int64_t time = 0; \
CppType value; \
while (time_decoder_->has_remaining(time_in)) { \
ASSERT(value_decoder_->has_remaining(value_in)); \
if (UNLIKELY(!row_appender.add_row())) { \
ret = E_OVERFLOW; \
break; \
} else if (RET_FAIL(time_decoder_->read_int64(time, time_in))) { \
} else if (RET_FAIL(value_decoder_->read_##ReadType(value, \
value_in))) { \
} else if (filter != nullptr && !filter->satisfy(time, value)) { \
row_appender.backoff_add_row(); \
continue; \
} else { \
/*std::cout << "decoder: time=" << time << ", value=" << value \
* << std::endl;*/ \
row_appender.append(0, (char*)&time, sizeof(time)); \
row_appender.append(1, (char*)&value, sizeof(value)); \
} \
} \
#define DECODE_TYPED_TV_INTO_TSBLOCK(CppType, ReadType, time_in, value_in, \
row_appender) \
do { \
int64_t time = 0; \
CppType value; \
while (time_decoder_->has_remaining(time_in)) { \
ASSERT(value_decoder_->has_remaining(value_in)); \
if (UNLIKELY(!row_appender.add_row())) { \
ret = E_OVERFLOW; \
break; \
} else if (RET_FAIL(time_decoder_->read_int64(time, time_in))) { \
} else if (RET_FAIL(value_decoder_->read_##ReadType(value, \
value_in))) { \
} else if (filter != nullptr && !filter->satisfy(time, value)) { \
row_appender.backoff_add_row(); \
continue; \
} else { \
if (row_offset > 0) { \
--row_offset; \
row_appender.backoff_add_row(); \
continue; \
} \
row_appender.append(0, (char*)&time, sizeof(time)); \
row_appender.append(1, (char*)&value, sizeof(value)); \
} \
} \
} while (false)

int ChunkReader::i32_DECODE_TYPED_TV_INTO_TSBLOCK(ByteStream& time_in,
Expand Down Expand Up @@ -467,19 +470,19 @@ int ChunkReader::i32_DECODE_TYPED_TV_INTO_TSBLOCK(ByteStream& time_in,
}

int ChunkReader::i32_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in,
RowAppender& row_appender,
Filter* filter) {
RowAppender& row_appender, Filter* filter,
int& row_offset) {
int ret = E_OK;
const int BATCH = 129;
int64_t times[BATCH];
int32_t values[BATCH];

while (time_decoder_->has_remaining(time_in)) {
// Cap each pass to what the appender can still hold; the old
// "remaining < BATCH → OVERFLOW" check made progress impossible on
// TsBlocks with capacity below BATCH.
int eff_batch =
std::min(BATCH, static_cast<int>(row_appender.remaining()));
// Skipped rows consume no output capacity. Bound the batch so every
// accepted row fits after the remaining offset has been consumed.
int eff_batch = std::min(
BATCH,
std::max(static_cast<int>(row_appender.remaining()), row_offset));
if (eff_batch <= 0) {
ret = E_OVERFLOW;
break;
Expand Down Expand Up @@ -543,6 +546,10 @@ int ChunkReader::i32_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in,
!filter->satisfy(times[i], (int64_t)values[i])) {
continue;
}
if (row_offset > 0) {
--row_offset;
continue;
}
if (UNLIKELY(!row_appender.add_row())) {
ret = E_OVERFLOW;
break;
Expand All @@ -556,16 +563,17 @@ int ChunkReader::i32_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in,
}

int ChunkReader::i64_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in,
RowAppender& row_appender,
Filter* filter) {
RowAppender& row_appender, Filter* filter,
int& row_offset) {
int ret = E_OK;
const int BATCH = 129;
int64_t times[BATCH];
int64_t values[BATCH];

while (time_decoder_->has_remaining(time_in)) {
int eff_batch =
std::min(BATCH, static_cast<int>(row_appender.remaining()));
int eff_batch = std::min(
BATCH,
std::max(static_cast<int>(row_appender.remaining()), row_offset));
if (eff_batch <= 0) {
ret = E_OVERFLOW;
break;
Expand Down Expand Up @@ -629,6 +637,10 @@ int ChunkReader::i64_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in,
!filter->satisfy(times[i], values[i])) {
continue;
}
if (row_offset > 0) {
--row_offset;
continue;
}
if (UNLIKELY(!row_appender.add_row())) {
ret = E_OVERFLOW;
break;
Expand All @@ -644,15 +656,16 @@ int ChunkReader::i64_DECODE_TV_BATCH(ByteStream& time_in, ByteStream& value_in,
int ChunkReader::float_DECODE_TV_BATCH(ByteStream& time_in,
ByteStream& value_in,
RowAppender& row_appender,
Filter* filter) {
Filter* filter, int& row_offset) {
int ret = E_OK;
const int BATCH = 129;
int64_t times[BATCH];
float values[BATCH];

while (time_decoder_->has_remaining(time_in)) {
int eff_batch =
std::min(BATCH, static_cast<int>(row_appender.remaining()));
int eff_batch = std::min(
BATCH,
std::max(static_cast<int>(row_appender.remaining()), row_offset));
if (eff_batch <= 0) {
ret = E_OVERFLOW;
break;
Expand Down Expand Up @@ -712,6 +725,10 @@ int ChunkReader::float_DECODE_TV_BATCH(ByteStream& time_in,
if (filter != nullptr && !block_all_pass && !time_mask[i]) {
continue;
}
if (row_offset > 0) {
--row_offset;
continue;
}
if (UNLIKELY(!row_appender.add_row())) {
ret = E_OVERFLOW;
break;
Expand All @@ -727,15 +744,16 @@ int ChunkReader::float_DECODE_TV_BATCH(ByteStream& time_in,
int ChunkReader::double_DECODE_TV_BATCH(ByteStream& time_in,
ByteStream& value_in,
RowAppender& row_appender,
Filter* filter) {
Filter* filter, int& row_offset) {
int ret = E_OK;
const int BATCH = 129;
int64_t times[BATCH];
double values[BATCH];

while (time_decoder_->has_remaining(time_in)) {
int eff_batch =
std::min(BATCH, static_cast<int>(row_appender.remaining()));
int eff_batch = std::min(
BATCH,
std::max(static_cast<int>(row_appender.remaining()), row_offset));
if (eff_batch <= 0) {
ret = E_OVERFLOW;
break;
Expand Down Expand Up @@ -795,6 +813,10 @@ int ChunkReader::double_DECODE_TV_BATCH(ByteStream& time_in,
if (filter != nullptr && !block_all_pass && !time_mask[i]) {
continue;
}
if (row_offset > 0) {
--row_offset;
continue;
}
if (UNLIKELY(!row_appender.add_row())) {
ret = E_OVERFLOW;
break;
Expand All @@ -807,11 +829,9 @@ int ChunkReader::double_DECODE_TV_BATCH(ByteStream& time_in,
return ret;
}

int ChunkReader::STRING_DECODE_TYPED_TV_INTO_TSBLOCK(ByteStream& time_in,
ByteStream& value_in,
RowAppender& row_appender,
PageArena& pa,
Filter* filter) {
int ChunkReader::STRING_DECODE_TYPED_TV_INTO_TSBLOCK(
ByteStream& time_in, ByteStream& value_in, RowAppender& row_appender,
PageArena& pa, Filter* filter, int& row_offset) {
int ret = E_OK;
int64_t time = 0;
common::String value;
Expand All @@ -825,6 +845,10 @@ int ChunkReader::STRING_DECODE_TYPED_TV_INTO_TSBLOCK(ByteStream& time_in,
} else if (filter != nullptr && !filter->satisfy(time, value)) {
row_appender.backoff_add_row();
continue;
} else if (row_offset > 0) {
--row_offset;
row_appender.backoff_add_row();
continue;
} else {
row_appender.append(0, (char*)&time, sizeof(time));
row_appender.append(1, value.buf_, value.len_);
Expand All @@ -833,12 +857,13 @@ int ChunkReader::STRING_DECODE_TYPED_TV_INTO_TSBLOCK(ByteStream& time_in,
return ret;
}

int ChunkReader::decode_tv_buf_into_tsblock_by_datatype(ByteStream& time_in,
ByteStream& value_in,
TsBlock* ret_tsblock,
Filter* filter,
common::PageArena* pa) {
int ChunkReader::decode_tv_buf_into_tsblock_by_datatype(
ByteStream& time_in, ByteStream& value_in, TsBlock* ret_tsblock,
Filter* filter, common::PageArena* pa, int* remaining_offset) {
int ret = E_OK;
int unused_offset = 0;
int& row_offset =
remaining_offset == nullptr ? unused_offset : *remaining_offset;
RowAppender row_appender(ret_tsblock);
switch (chunk_header_.data_type_) {
case common::BOOLEAN:
Expand All @@ -847,33 +872,36 @@ int ChunkReader::decode_tv_buf_into_tsblock_by_datatype(ByteStream& time_in,
break;
case common::DATE:
case common::INT32:
ret =
i32_DECODE_TV_BATCH(time_in_, value_in_, row_appender, filter);
ret = i32_DECODE_TV_BATCH(time_in_, value_in_, row_appender, filter,
row_offset);
break;
case TIMESTAMP:
case common::INT64:
ret =
i64_DECODE_TV_BATCH(time_in_, value_in_, row_appender, filter);
ret = i64_DECODE_TV_BATCH(time_in_, value_in_, row_appender, filter,
row_offset);
break;
case common::FLOAT:
ret = float_DECODE_TV_BATCH(time_in_, value_in_, row_appender,
filter);
filter, row_offset);
break;
case common::DOUBLE:
ret = double_DECODE_TV_BATCH(time_in_, value_in_, row_appender,
filter);
filter, row_offset);
break;
case common::TEXT:
case common::BLOB:
case common::STRING:
ret = STRING_DECODE_TYPED_TV_INTO_TSBLOCK(
time_in, value_in, row_appender, *pa, filter);
time_in, value_in, row_appender, *pa, filter, row_offset);
break;
default:
ret = E_NOT_SUPPORT;
ASSERT(false);
}
if (ret_tsblock->get_row_count() == 0 && ret == E_OK) {
// The offset-aware iterator must keep scanning after a page whose
// matching rows were all consumed by the offset (or time filter).
if (remaining_offset == nullptr && ret_tsblock->get_row_count() == 0 &&
ret == E_OK) {
ret = E_NO_MORE_DATA;
}
return ret;
Expand Down Expand Up @@ -917,8 +945,8 @@ int ChunkReader::get_next_page(TsBlock* ret_tsblock, Filter* oneshoot_filter,
}

if (prev_page_not_finish()) {
ret = decode_tv_buf_into_tsblock_by_datatype(time_in_, value_in_,
ret_tsblock, filter, &pa);
ret = decode_tv_buf_into_tsblock_by_datatype(
time_in_, value_in_, ret_tsblock, filter, &pa, &row_offset);
if (ret == E_OVERFLOW) {
ret = E_OK;
} else {
Expand Down Expand Up @@ -953,7 +981,7 @@ int ChunkReader::get_next_page(TsBlock* ret_tsblock, Filter* oneshoot_filter,
}

if (IS_SUCC(ret)) {
ret = decode_cur_page_data(ret_tsblock, filter, pa);
ret = decode_cur_page_data(ret_tsblock, filter, pa, &row_offset);
}
return ret;
}
Expand Down
19 changes: 11 additions & 8 deletions cpp/src/reader/chunk_reader.h
Original file line number Diff line number Diff line change
Expand Up @@ -91,7 +91,7 @@ class ChunkReader : public IChunkReader {
bool cur_page_fully_satisfies_filter(Filter* filter);
int skip_cur_page();
int decode_cur_page_data(common::TsBlock*& ret_tsblock, Filter* filter,
common::PageArena& pa);
common::PageArena& pa, int* row_offset = nullptr);
bool prev_page_not_finish() const {
return (time_decoder_ && time_decoder_->has_remaining(time_in_)) ||
time_in_.has_remaining();
Expand All @@ -101,30 +101,33 @@ class ChunkReader : public IChunkReader {
common::ByteStream& value_in,
common::TsBlock* ret_tsblock,
Filter* filter,
common::PageArena* pa = nullptr);
common::PageArena* pa = nullptr,
int* remaining_offset = nullptr);
int i32_DECODE_TYPED_TV_INTO_TSBLOCK(common::ByteStream& time_in,
common::ByteStream& value_in,
common::RowAppender& row_appender,
Filter* filter);
int i32_DECODE_TV_BATCH(common::ByteStream& time_in,
common::ByteStream& value_in,
common::RowAppender& row_appender, Filter* filter);
common::RowAppender& row_appender, Filter* filter,
int& row_offset);
int i64_DECODE_TV_BATCH(common::ByteStream& time_in,
common::ByteStream& value_in,
common::RowAppender& row_appender, Filter* filter);
common::RowAppender& row_appender, Filter* filter,
int& row_offset);
int float_DECODE_TV_BATCH(common::ByteStream& time_in,
common::ByteStream& value_in,
common::RowAppender& row_appender,
Filter* filter);
common::RowAppender& row_appender, Filter* filter,
int& row_offset);
int double_DECODE_TV_BATCH(common::ByteStream& time_in,
common::ByteStream& value_in,
common::RowAppender& row_appender,
Filter* filter);
Filter* filter, int& row_offset);
int STRING_DECODE_TYPED_TV_INTO_TSBLOCK(common::ByteStream& time_in,
common::ByteStream& value_in,
common::RowAppender& row_appender,
common::PageArena& pa,
Filter* filter);
Filter* filter, int& row_offset);

private:
RandomAccessReadFile* read_file_;
Expand Down
Loading
Loading