This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/master by this push:
new 573c93c63b8 [improvement](segment): Batch read array data by row IDs
(#67817)
573c93c63b8 is described below
commit 573c93c63b88fb68ac1d318dcd3e1672dd2361e6
Author: foxtail463 <[email protected]>
AuthorDate: Thu Sep 24 20:21:29 2026 +0800
[improvement](segment): Batch read array data by row IDs (#67817)
### Problem Summary:
ARRAY row-ID reads seek and materialize each parent row separately,
repeating metadata processing and column setup and limiting sparse-read
throughput.
### Solution:
Batch-read parent null maps and ordered offset endpoints, then coalesce
physically contiguous item ranges for bulk
element reads. Preserve append, metadata-only, and lazy-materialization
semantics, and track temporary buffers through Doris allocators.
---------
Co-authored-by: yangtao555 <[email protected]>
---
be/src/storage/segment/column_reader.cpp | 190 ++++-
be/test/storage/segment/column_reader_test.cpp | 1062 +++++++++++++++++++++++-
2 files changed, 1244 insertions(+), 8 deletions(-)
diff --git a/be/src/storage/segment/column_reader.cpp
b/be/src/storage/segment/column_reader.cpp
index 8bd5d0f2073..79ac783f2bc 100644
--- a/be/src/storage/segment/column_reader.cpp
+++ b/be/src/storage/segment/column_reader.cpp
@@ -41,6 +41,7 @@
#include "core/column/column_nullable.h"
#include "core/column/column_struct.h"
#include "core/column/column_vector.h"
+#include "core/custom_allocator.h"
#include "core/data_type/data_type_agg_state.h"
#include "core/data_type/data_type_factory.hpp"
#include "core/data_type/data_type_nullable.h"
@@ -2438,6 +2439,11 @@ void ArrayFileColumnIterator::collect_prefetchers(
}
}
+// Materialize selected parent rows without repeating the ARRAY
seek/next_batch setup per row.
+// Requires rowids in nondecreasing segment-local order. Batch-read ordered
metadata, then coalesce
+// adjacent source item spans into fewer item reads.
+// Normally append complete rows to dst; in LAZY mode, fill missing children
without duplicating
+// parent offsets/null-map that were already materialized and filtered in the
predicate phase.
Status ArrayFileColumnIterator::read_by_rowids(const rowid_t* rowids, const
size_t count,
MutableColumnPtr& dst) {
if (!need_to_read()) {
@@ -2448,11 +2454,189 @@ Status ArrayFileColumnIterator::read_by_rowids(const
rowid_t* rowids, const size
_recovery_from_place_holder_column(dst);
+ if (count == 0) {
+ return Status::OK();
+ }
+
+ DCHECK(std::is_sorted(rowids, rowids + count));
+
+ // A null-only consumer needs one nested row per null marker, but no
lengths or item data.
+ if (read_null_map_only()) {
+ DORIS_CHECK(is_column_nullable(*dst));
+ auto& nullable_column = assert_cast<ColumnNullable&>(*dst);
+ if (_null_iterator) {
+ auto null_map_ptr = nullable_column.get_null_map_column_ptr();
+ MutableColumnPtr null_map_column = std::move(null_map_ptr);
+ RETURN_IF_ERROR(_null_iterator->read_by_rowids(rowids, count,
null_map_column));
+ } else {
+ // A nullable schema can read an old non-nullable segment, which
has no null stream.
+ nullable_column.get_null_map_column_ptr()->insert_many_vals(0,
count);
+ }
+ auto& column_array = assert_cast<ColumnArray&,
TypeCheckOnRelease::DISABLE>(
+ nullable_column.get_nested_column());
+ column_array.insert_many_defaults(count);
+ return Status::OK();
+ }
+
+ auto& column_array = assert_cast<ColumnArray&,
TypeCheckOnRelease::DISABLE>(
+ is_column_nullable(*dst) ?
static_cast<ColumnNullable&>(*dst).get_nested_column()
+ : *dst);
+ // The parent reader can stay active solely for lazy children. This flag
controls writing
+ // parent metadata to dst, not reading source offsets to locate those
children on disk.
+ const bool read_meta_columns = need_to_read_meta_columns();
+
+ if (_array_reader->is_nullable()) {
+ if (UNLIKELY(!is_column_nullable(*dst))) {
+ return Status::InternalError(
+ "unexpected non-nullable destination column for nullable
array reader");
+ }
+ auto& nullable_column = static_cast<ColumnNullable&>(*dst);
+ if (read_meta_columns) {
+ MutableColumnPtr null_map_column =
nullable_column.get_null_map_column_ptr();
+ RETURN_IF_ERROR(_null_iterator->read_by_rowids(rowids, count,
null_map_column));
+ } else {
+ DORIS_CHECK(nullable_column.get_null_map_column().size() == count);
+ }
+ } else if (read_meta_columns && is_column_nullable(*dst)) {
+
static_cast<ColumnNullable&>(*dst).get_null_map_column_ptr()->insert_many_vals(0,
count);
+ }
+
+ // Array row r spans [offset[r], offset[r + 1]) in the source item stream.
offset_rowids
+ // identifies entries in the offset stream, not item ordinals. Read both
endpoints in one
+ // ordered pass to avoid revisiting pages for the ends; adjacent rows
share an endpoint:
+ // rowids [1, 2, 8] need offset entries [1, 2, 3, 8, 9].
+ DorisVector<rowid_t> offset_rowids;
+ offset_rowids.reserve(count * 2);
for (size_t i = 0; i < count; ++i) {
- // TODO(cambyszju): now read array one by one, need optimize later
- RETURN_IF_ERROR(seek_to_ordinal(rowids[i]));
+ offset_rowids.push_back(rowids[i]);
+ const auto next_rowid = static_cast<uint64_t>(rowids[i]) + 1;
+ if (next_rowid < _array_reader->num_rows() &&
+ (i + 1 == count || next_rowid != rowids[i + 1])) {
+ offset_rowids.push_back(static_cast<rowid_t>(next_rowid));
+ }
+ }
+ MutableColumnPtr source_offsets_column = ColumnOffset64::create();
+ source_offsets_column->reserve(offset_rowids.size() + 1);
+ RETURN_IF_ERROR(_offset_iterator->read_by_rowids(offset_rowids.data(),
offset_rowids.size(),
+ source_offsets_column));
+ // source_offsets contains element ordinals in the segment, not file byte
positions.
+ // Destination offsets instead delimit the compact item stream after row
selection.
+ auto& source_offsets =
assert_cast<ColumnOffset64&>(*source_offsets_column).get_data();
+ DORIS_CHECK(source_offsets.size() == offset_rowids.size());
+ if (static_cast<uint64_t>(rowids[count - 1]) + 1 ==
_array_reader->num_rows()) {
+ // The last array row has no rowid + 1. Consume its start offset, then
obtain the end
+ // offset from the page-tail sentinel written by OffsetColumnWriter.
+ RETURN_IF_ERROR(_offset_iterator->seek_to_ordinal(rowids[count - 1]));
size_t num_read = 1;
- RETURN_IF_ERROR(next_batch(&num_read, dst));
+ bool has_null = false;
+ MutableColumnPtr last_start = ColumnOffset64::create();
+ RETURN_IF_ERROR(_offset_iterator->next_batch(&num_read, last_start,
&has_null));
+ if (UNLIKELY(num_read != 1)) {
+ return Status::Corruption("failed to read the last array offset");
+ }
+ ordinal_t next_start = 0;
+ RETURN_IF_ERROR(_offset_iterator->_peek_one_offset(&next_start));
+ source_offsets.push_back(next_start);
+ }
+
+ if (!read_meta_columns) {
+ DORIS_CHECK(column_array.size() == count);
+ }
+
+ // Obtain writable child owners once for the whole batch, and restore them
on every exit.
+ // Moving the owners avoids introducing extra sharing just to read into
their columns.
+ MutableColumnPtr output_offsets_ptr;
+ ColumnArray::ColumnOffsets* output_offsets = nullptr;
+ if (read_meta_columns) {
+ output_offsets_ptr =
IColumn::mutate(std::move(column_array.get_offsets_ptr()));
+ output_offsets = assert_cast<ColumnArray::ColumnOffsets*,
TypeCheckOnRelease::DISABLE>(
+ output_offsets_ptr.get());
+ }
+ Defer defer_offsets {[&] {
+ if (read_meta_columns) {
+ auto typed_offsets_ptr =
ColumnArray::ColumnOffsets::cast_to_column_mutptr(
+ assert_cast<ColumnArray::ColumnOffsets*,
TypeCheckOnRelease::DISABLE>(
+ output_offsets_ptr.get()));
+ output_offsets_ptr = nullptr;
+ column_array.get_offsets_ptr() = std::move(typed_offsets_ptr);
+ }
+ }};
+
+ auto items_ptr = IColumn::mutate(std::move(column_array.get_data_ptr()));
+ Defer defer_items {[&] { column_array.get_data_ptr() =
std::move(items_ptr); }};
+ // Append exactly one selected source span; short reads must not yield an
incomplete ARRAY.
+ auto read_item_range = [&](ordinal_t start, size_t item_count) -> Status {
+ if (item_count == 0) {
+ return Status::OK();
+ }
+ size_t num_read = item_count;
+ bool has_null = false;
+ RETURN_IF_ERROR(_item_iterator->seek_to_ordinal(start));
+ RETURN_IF_ERROR(_item_iterator->next_batch(&num_read, items_ptr,
&has_null));
+ if (UNLIKELY(num_read != item_count)) {
+ return Status::Corruption("array item reader returned {} items,
expected {}", num_read,
+ item_count);
+ }
+ return Status::OK();
+ };
+
+ // Cumulative end in dst's item stream, not a source ordinal. Continue the
previous end
+ // so this batch can append to a non-empty destination.
+ uint64_t output_offset =
+ read_meta_columns && !output_offsets->empty() ?
output_offsets->get_data().back() : 0;
+ size_t total_item_count = 0; // Number of placeholder item slots needed by
OFFSET_ONLY.
+ // Delay reading this pending source span so adjacent item ranges share
one read, even
+ // when their parent rowids are separated by unselected rows with no
physical items.
+ ordinal_t range_start = 0;
+ size_t range_size = 0;
+ if (read_meta_columns) {
+ output_offsets->get_data().reserve(output_offsets->size() + count);
+ }
+ // Cursor in the compact endpoint buffer, not a parent rowid or a source
item ordinal.
+ // Each iteration leaves it on that row's end; only adjacent parent rows
reuse it as a start.
+ size_t offset_index = 0;
+ for (size_t i = 0; i < count; ++i) {
+ if (i > 0 && static_cast<uint64_t>(rowids[i - 1]) + 1 != rowids[i]) {
+ ++offset_index;
+ }
+ const ordinal_t item_start = source_offsets[offset_index++];
+ const ordinal_t item_end = source_offsets[offset_index];
+ if (UNLIKELY(item_end < item_start)) {
+ return Status::Corruption("invalid array element offsets: start
{}, end {}", item_start,
+ item_end);
+ }
+ const size_t item_count = static_cast<size_t>(item_end - item_start);
+ // A nullable parent can legally retain nested payload for a null row.
Preserve the raw
+ // offset span so the nested column stays aligned with the parent
offsets, especially when
+ // lazy materialization fills only the item subtree in a later phase.
+ total_item_count += item_count;
+ if (read_meta_columns) {
+ output_offset += item_count;
+ output_offsets->get_data().push_back(output_offset);
+ } else {
+ DCHECK_EQ(column_array.size_at(i), item_count);
+ }
+ if (read_offset_only() || item_count == 0) {
+ continue;
+ }
+
+ if (range_size == 0) {
+ range_start = item_start;
+ range_size = item_count;
+ } else if (range_start + range_size == item_start) {
+ range_size += item_count;
+ } else {
+ RETURN_IF_ERROR(read_item_range(range_start, range_size));
+ range_start = item_start;
+ range_size = item_count;
+ }
+ }
+
+ DCHECK_EQ(offset_index + 1, source_offsets.size());
+ if (read_offset_only()) {
+ items_ptr->insert_many_defaults(total_item_count);
+ } else {
+ RETURN_IF_ERROR(read_item_range(range_start, range_size));
}
return Status::OK();
}
diff --git a/be/test/storage/segment/column_reader_test.cpp
b/be/test/storage/segment/column_reader_test.cpp
index f2b4eec5a3a..51b843c43bb 100644
--- a/be/test/storage/segment/column_reader_test.cpp
+++ b/be/test/storage/segment/column_reader_test.cpp
@@ -23,6 +23,7 @@
#include <gmock/gmock.h>
#include <gtest/gtest.h>
+#include <array>
#include <chrono>
#include <cstdint>
#include <iterator>
@@ -36,6 +37,11 @@
#include "common/config.h"
#include "core/column/column_map.h"
#include "core/column/column_struct.h"
+#include "core/data_type/data_type_array.h"
+#include "core/data_type/data_type_map.h"
+#include "core/data_type/data_type_nullable.h"
+#include "core/data_type/data_type_number.h"
+#include "core/data_type/data_type_struct.h"
#include "io/fs/file_reader.h"
#include "io/fs/file_system.h"
#include "io/fs/file_writer.h"
@@ -202,6 +208,10 @@ public:
routed_predicate_access_paths.clear();
}
+ void convert_to_place_holder_column(MutableColumnPtr& dst, size_t count) {
+ _convert_to_place_holder_column(dst, count);
+ }
+
std::vector<ordinal_t> seek_ordinals;
std::vector<size_t> next_batch_sizes;
std::vector<std::vector<rowid_t>> read_by_rowids_batches;
@@ -273,17 +283,27 @@ private:
class RowidOffsetFileColumnIterator final : public FileColumnIterator {
public:
- RowidOffsetFileColumnIterator() :
FileColumnIterator(create_test_reader(false, 10)) {}
+ RowidOffsetFileColumnIterator()
+ : RowidOffsetFileColumnIterator(
+ std::vector<ordinal_t> {0, 1, 2, 3, 4, 5, 6, 7, 8, 9,
10}) {}
+
+ explicit RowidOffsetFileColumnIterator(std::vector<ordinal_t> offsets)
+ : FileColumnIterator(create_test_reader(false, offsets.size() -
1)),
+ _offsets(std::move(offsets)) {
+ get_current_page()->next_array_item_ordinal = _offsets.back();
+ }
Status seek_to_ordinal(ordinal_t ord) override {
+ seek_ordinals.emplace_back(ord);
_current_ordinal = ord;
return Status::OK();
}
Status next_batch(size_t* n, MutableColumnPtr& dst, bool* has_null)
override {
+ next_batch_sizes.emplace_back(*n);
auto& offsets = assert_cast<ColumnOffset64&,
TypeCheckOnRelease::DISABLE>(*dst);
for (size_t i = 0; i < *n; ++i) {
- offsets.insert_value(_current_ordinal + i);
+ offsets.insert_value(_offsets[_current_ordinal + i]);
}
_current_ordinal += *n;
if (has_null != nullptr) {
@@ -294,16 +314,22 @@ public:
Status read_by_rowids(const rowid_t* rowids, const size_t count,
MutableColumnPtr& dst) override {
+ read_by_rowids_batches.emplace_back(rowids, rowids + count);
auto& offsets = assert_cast<ColumnOffset64&,
TypeCheckOnRelease::DISABLE>(*dst);
for (size_t i = 0; i < count; ++i) {
- offsets.insert_value(rowids[i]);
+ offsets.insert_value(_offsets[rowids[i]]);
}
return Status::OK();
}
ordinal_t get_current_ordinal() const override { return _current_ordinal; }
+ std::vector<ordinal_t> seek_ordinals;
+ std::vector<size_t> next_batch_sizes;
+ std::vector<std::vector<rowid_t>> read_by_rowids_batches;
+
private:
+ std::vector<ordinal_t> _offsets;
ordinal_t _current_ordinal = 0;
};
@@ -344,6 +370,129 @@ struct TrackingOffsetIterator {
TrackingFileColumnIterator* tracker = nullptr;
};
+// item_data is the nested writer's raw input and differs by element type. It
is not the
+// outer array's own data layout; write_array_column only adds the outer array
metadata.
+void write_array_column(const std::string& file_name, ColumnMetaPB* meta,
+ const TabletColumn& tablet_column, size_t num_rows,
+ const std::vector<uint64_t>& item_data,
+ const std::vector<uint64_t>& outer_offsets,
+ const std::vector<uint8_t>& item_null_map,
+ const std::vector<uint8_t>* outer_null_map) {
+ auto fs = io::global_local_filesystem();
+ io::FileWriterPtr file_writer;
+ auto st = fs->create_file(file_name, &file_writer);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+
+ ColumnWriterOptions writer_options;
+ writer_options.meta = meta;
+ std::unique_ptr<ColumnWriter> writer;
+ st = ColumnWriter::create(writer_options, &tablet_column,
file_writer.get(), &writer);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ ASSERT_TRUE(writer->init().ok());
+
+ std::vector<uint64_t> outer_data
{static_cast<uint64_t>(item_null_map.size()),
+
reinterpret_cast<uint64_t>(outer_offsets.data()),
+
reinterpret_cast<uint64_t>(item_data.data()),
+
reinterpret_cast<uint64_t>(item_null_map.data())};
+ ASSERT_TRUE(writer->append(outer_null_map ? outer_null_map->data() :
nullptr, outer_data.data(),
+ num_rows)
+ .ok());
+ ASSERT_TRUE(writer->finish().ok());
+ ASSERT_TRUE(writer->write_data().ok());
+ ASSERT_TRUE(writer->write_ordinal_index().ok());
+ ASSERT_TRUE(file_writer->close().ok());
+}
+
+void read_array_baseline(const std::shared_ptr<ColumnReader>& reader,
+ const TabletColumn& tablet_column, const DataTypePtr&
column_type,
+ io::FileReader* file_reader, MutableColumnPtr*
baseline) {
+ ColumnIteratorUPtr iterator;
+ OlapReaderStatistics stats;
+ ASSERT_TRUE(reader->new_iterator(&iterator, &tablet_column).ok());
+ ColumnIteratorOptions options;
+ options.stats = &stats;
+ options.file_reader = file_reader;
+ ASSERT_TRUE(iterator->init(options).ok());
+ ASSERT_TRUE(iterator->seek_to_ordinal(0).ok());
+
+ *baseline = column_type->create_column();
+ size_t rows_to_read = reader->num_rows();
+ bool has_null = false;
+ ASSERT_TRUE(iterator->next_batch(&rows_to_read, *baseline,
&has_null).ok());
+ ASSERT_EQ(reader->num_rows(), rows_to_read);
+}
+
+void check_lazy_array_read_matches_baseline(const
std::shared_ptr<ColumnReader>& reader,
+ const TabletColumn& tablet_column,
+ const DataTypePtr& column_type,
+ const std::vector<rowid_t>& rowids,
+ const IColumn& baseline,
io::FileReader* file_reader) {
+ ColumnIteratorUPtr iterator;
+ OlapReaderStatistics stats;
+ ASSERT_TRUE(reader->new_iterator(&iterator, &tablet_column).ok());
+ ColumnIteratorOptions options;
+ options.stats = &stats;
+ options.file_reader = file_reader;
+ ASSERT_TRUE(iterator->init(options).ok());
+
+ iterator->set_column_name(tablet_column.name());
+ TColumnAccessPaths all_paths
{create_data_access_path({tablet_column.name()})};
+ TColumnAccessPaths predicate_paths {
+ create_meta_access_path({tablet_column.name(),
ColumnIterator::ACCESS_OFFSET})};
+ ASSERT_TRUE(iterator->set_access_paths(all_paths, predicate_paths).ok());
+
+ iterator->set_read_phase(ColumnIterator::ReadPhase::PREDICATE);
+ auto lazy = column_type->create_column();
+ ASSERT_TRUE(iterator->seek_to_ordinal(0).ok());
+ size_t rows_to_read = baseline.size();
+ bool has_null = false;
+ ASSERT_TRUE(iterator->next_batch(&rows_to_read, lazy, &has_null).ok());
+ ASSERT_EQ(baseline.size(), rows_to_read);
+
+ IColumn::Filter filter;
+ filter.resize_fill(baseline.size(), 0);
+ for (rowid_t rowid : rowids) {
+ filter[rowid] = 1;
+ }
+ lazy = IColumn::mutate(lazy->filter(filter, rowids.size()));
+
+ iterator->set_read_phase(ColumnIterator::ReadPhase::LAZY);
+ ASSERT_TRUE(iterator->need_to_read());
+ ASSERT_FALSE(iterator->need_to_read_meta_columns());
+ ASSERT_TRUE(iterator->read_by_rowids(rowids.data(), rowids.size(),
lazy).ok());
+ iterator->finalize_lazy_phase(lazy);
+
+ ASSERT_EQ(rowids.size(), lazy->size());
+ for (size_t i = 0; i < rowids.size(); ++i) {
+ EXPECT_EQ(0, lazy->compare_at(i, rowids[i], baseline, 1));
+ }
+}
+
+void check_array_read_by_rowids_matches_baseline(const
std::shared_ptr<ColumnReader>& reader,
+ const TabletColumn&
tablet_column,
+ const DataTypePtr&
column_type,
+ const std::vector<rowid_t>&
rowids,
+ const IColumn& baseline,
+ io::FileReader* file_reader) {
+ ColumnIteratorUPtr iterator;
+ OlapReaderStatistics stats;
+ ASSERT_TRUE(reader->new_iterator(&iterator, &tablet_column).ok());
+ ColumnIteratorOptions options;
+ options.stats = &stats;
+ options.file_reader = file_reader;
+ ASSERT_TRUE(iterator->init(options).ok());
+
+ auto actual = column_type->create_column();
+ ASSERT_TRUE(iterator->read_by_rowids(rowids.data(), rowids.size(),
actual).ok());
+ ASSERT_EQ(rowids.size(), actual->size());
+ for (size_t i = 0; i < rowids.size(); ++i) {
+ EXPECT_EQ(0, actual->compare_at(i, rowids[i], baseline, 1));
+ }
+
+ check_lazy_array_read_matches_baseline(reader, tablet_column, column_type,
rowids, baseline,
+ file_reader);
+}
+
TrackingOffsetIterator create_tracking_offset_iterator() {
auto file_iterator =
std::make_unique<TrackingFileColumnIterator>(create_test_reader());
auto* tracker = file_iterator.get();
@@ -451,6 +600,751 @@ TEST_F(ColumnReaderTest,
NullMapOnlyReadBySparseRowidsAcrossPages) {
EXPECT_EQ(2, nullable_col.get_nested_column().size());
}
+TEST_F(ColumnReaderTest, ArrayReadByRowidsMatchesSequentialReadAcrossPages) {
+ // The generated offsets and non-null INT items both exceed the default 64
KiB data-page
+ // size, so the selected rowids exercise page transitions in both child
streams.
+ constexpr size_t num_rows = 12000;
+ constexpr size_t max_items_per_row = 3;
+ ColumnMetaPB meta;
+ TabletColumn
array_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_ARRAY);
+ TabletColumn
item_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_INT, true);
+ array_column.add_sub_column(item_column);
+
+ std::vector<int32_t> item_values;
+ std::vector<uint8_t> item_null_map;
+ std::vector<uint64_t> array_offsets(num_rows + 1, 0);
+ std::vector<uint8_t> array_null_map(num_rows, 0);
+ item_values.reserve(num_rows * max_items_per_row);
+ item_null_map.reserve(num_rows * max_items_per_row);
+ for (size_t row = 0; row < num_rows; ++row) {
+ const bool is_null = row % 7 == 1;
+ array_null_map[row] = is_null;
+ const size_t item_count = row % 11 == 0 ? 0 : (is_null ? 3 : row % 3 +
1);
+ for (size_t item = 0; item < item_count; ++item) {
+ item_values.push_back(static_cast<int32_t>(row * 10 + item));
+ item_null_map.push_back((row + item) % 5 == 0);
+ }
+ array_offsets[row + 1] = item_values.size();
+ }
+
+ const std::string file_name = COLUMN_READER_FILE_TEST_DIR +
"/array_read_by_rowids";
+ auto fs = io::global_local_filesystem();
+ {
+ io::FileWriterPtr file_writer;
+ auto st = fs->create_file(file_name, &file_writer);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+
+ ColumnWriterOptions writer_options;
+ writer_options.meta = &meta;
+ writer_options.meta->set_column_id(0);
+ writer_options.meta->set_unique_id(0);
+
writer_options.meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_ARRAY));
+ writer_options.meta->set_length(0);
+ writer_options.meta->set_encoding(DEFAULT_ENCODING);
+ writer_options.meta->set_compression(CompressionTypePB::LZ4F);
+ writer_options.meta->set_is_nullable(true);
+
+ auto* child_meta = meta.add_children_columns();
+ child_meta->set_column_id(1);
+ child_meta->set_unique_id(1);
+
child_meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_INT));
+ child_meta->set_length(0);
+ child_meta->set_encoding(BIT_SHUFFLE);
+ child_meta->set_compression(CompressionTypePB::LZ4F);
+ child_meta->set_is_nullable(true);
+
+ std::unique_ptr<ColumnWriter> writer;
+ st = ColumnWriter::create(writer_options, &array_column,
file_writer.get(), &writer);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ ASSERT_TRUE(writer->init().ok());
+ // Use the vectorized append path so even null parent rows keep their
physical item span.
+ const std::array<uint64_t, 4> array_data
{static_cast<uint64_t>(item_values.size()),
+
reinterpret_cast<uint64_t>(array_offsets.data()),
+
reinterpret_cast<uint64_t>(item_values.data()),
+
reinterpret_cast<uint64_t>(item_null_map.data())};
+ ASSERT_TRUE(writer->append(array_null_map.data(), array_data.data(),
num_rows).ok());
+ ASSERT_TRUE(writer->finish().ok());
+ ASSERT_TRUE(writer->write_data().ok());
+ ASSERT_TRUE(writer->write_ordinal_index().ok());
+ ASSERT_TRUE(file_writer->close().ok());
+ }
+
+ io::FileReaderSPtr file_reader;
+ auto st = fs->open_file(file_name, &file_reader);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ ColumnReaderOptions reader_options;
+ std::shared_ptr<ColumnReader> reader;
+ st = ColumnReader::create(reader_options, meta, num_rows, file_reader,
&reader);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+
+ DataTypePtr array_type = std::make_shared<DataTypeArray>(
+
std::make_shared<DataTypeNullable>(std::make_shared<DataTypeInt32>()));
+ DataTypePtr column_type = std::make_shared<DataTypeNullable>(array_type);
+ auto create_iterator = [&](OlapReaderStatistics* stats,
+ ColumnIteratorUPtr* iterator) -> Status {
+ RETURN_IF_ERROR(reader->new_iterator(iterator, &array_column));
+ ColumnIteratorOptions options;
+ options.stats = stats;
+ options.file_reader = file_reader.get();
+ return (*iterator)->init(options);
+ };
+
+ MutableColumnPtr baseline = column_type->create_column();
+ {
+ ColumnIteratorUPtr iterator;
+ OlapReaderStatistics stats;
+ st = create_iterator(&stats, &iterator);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ ASSERT_TRUE(iterator->seek_to_ordinal(0).ok());
+ size_t rows_to_read = num_rows;
+ bool has_null = false;
+ st = iterator->next_batch(&rows_to_read, baseline, &has_null);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ ASSERT_EQ(num_rows, rows_to_read);
+ }
+
+ std::vector<rowid_t> rowids {0, 1, 2, 7, 8, 11, 4095, 4096, 8191, 8192};
+ for (rowid_t rowid = 10000; rowid <= 11000; ++rowid) {
+ rowids.push_back(rowid);
+ }
+ rowids.push_back(11998);
+ rowids.push_back(11999);
+ MutableColumnPtr actual = column_type->create_column();
+ {
+ ColumnIteratorUPtr iterator;
+ OlapReaderStatistics stats;
+ st = create_iterator(&stats, &iterator);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ st = iterator->read_by_rowids(rowids.data(), rowids.size(), actual);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ }
+
+ ASSERT_EQ(1, array_null_map[1]);
+ ASSERT_GT(array_offsets[2] - array_offsets[1], 0);
+ const auto& baseline_nullable = assert_cast<const
ColumnNullable&>(*baseline);
+ const auto& baseline_array = assert_cast<const ColumnArray&,
TypeCheckOnRelease::DISABLE>(
+ baseline_nullable.get_nested_column());
+ ASSERT_EQ(1, baseline_nullable.get_null_map_data()[1]);
+ ASSERT_GT(baseline_array.size_at(1), 0);
+
+ // Compare physical item spans as well as logical rows, including payload
under NULL parents.
+ auto check_selected_rows = [&](const std::vector<rowid_t>& selected_rowids,
+ const IColumn& result) {
+ ASSERT_EQ(selected_rowids.size(), result.size());
+ const auto& actual_nullable = assert_cast<const
ColumnNullable&>(result);
+ const auto& actual_array = assert_cast<const ColumnArray&,
TypeCheckOnRelease::DISABLE>(
+ actual_nullable.get_nested_column());
+ size_t expected_item_count = 0;
+ size_t actual_item_index = 0;
+ for (size_t i = 0; i < selected_rowids.size(); ++i) {
+ const auto rowid = selected_rowids[i];
+ EXPECT_EQ(baseline_nullable.get_null_map_data()[rowid],
+ actual_nullable.get_null_map_data()[i]);
+
+ const size_t row_item_count = baseline_array.size_at(rowid);
+ expected_item_count += row_item_count;
+ EXPECT_EQ(expected_item_count, actual_array.get_offsets()[i]);
+ for (size_t item = 0; item < row_item_count; ++item) {
+ EXPECT_EQ(0, actual_array.get_data().compare_at(
+ actual_item_index,
baseline_array.offset_at(rowid) + item,
+ baseline_array.get_data(), 1));
+ ++actual_item_index;
+ }
+ }
+ EXPECT_EQ(expected_item_count, actual_array.get_data().size());
+ };
+ check_selected_rows(rowids, *actual);
+
+ {
+ SCOPED_TRACE("predicate metadata -> filter -> lazy items");
+ ColumnIteratorUPtr iterator;
+ OlapReaderStatistics stats;
+ ASSERT_TRUE(create_iterator(&stats, &iterator).ok());
+ iterator->set_column_name("a");
+ // The predicate needs the parent lengths, while the output needs
complete arrays.
+ // Let the same access-path API as production assign parent/child read
requirements.
+ TColumnAccessPaths all_paths {create_data_access_path({"a"})};
+ TColumnAccessPaths predicate_paths {
+ create_meta_access_path({"a", ColumnIterator::ACCESS_OFFSET})};
+ ASSERT_TRUE(iterator->set_access_paths(all_paths,
predicate_paths).ok());
+ iterator->set_read_phase(ColumnIterator::ReadPhase::PREDICATE);
+ auto lazy = column_type->create_column();
+ ASSERT_TRUE(iterator->seek_to_ordinal(0).ok());
+ size_t rows_to_read = num_rows;
+ bool has_null = false;
+ ASSERT_TRUE(iterator->next_batch(&rows_to_read, lazy, &has_null).ok());
+ ASSERT_EQ(num_rows, rows_to_read);
+ {
+ const auto& array = assert_cast<const ColumnArray&>(
+ assert_cast<const
ColumnNullable&>(*lazy).get_nested_column());
+ ASSERT_EQ(item_values.size(), array.get_data().size());
+ // Nullable item defaults are NULL: the predicate phase must not
read real items.
+ EXPECT_TRUE(array.get_data().only_null());
+ }
+
+ IColumn::Filter filter;
+ filter.resize_fill(num_rows, 0);
+ for (rowid_t rowid : rowids) {
+ filter[rowid] = 1;
+ }
+ lazy = IColumn::mutate(lazy->filter(filter, rowids.size()));
+ iterator->set_read_phase(ColumnIterator::ReadPhase::LAZY);
+ ASSERT_TRUE(iterator->need_to_read());
+ ASSERT_FALSE(iterator->need_to_read_meta_columns());
+ ASSERT_TRUE(iterator->read_by_rowids(rowids.data(), rowids.size(),
lazy).ok());
+ iterator->finalize_lazy_phase(lazy);
+ check_selected_rows(rowids, *lazy);
+ }
+}
+
+TEST_F(ColumnReaderTest, ArrayReadByRowidsNestedArrayMatchesSequentialRead) {
+ constexpr size_t num_rows = 12000;
+ ColumnMetaPB meta;
+ TabletColumn
array_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_ARRAY);
+ TabletColumn
nested_array_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_ARRAY);
+ TabletColumn
int_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_INT, true);
+ nested_array_column.add_sub_column(int_column);
+ array_column.add_sub_column(nested_array_column);
+ array_column.set_is_nullable(true);
+ nested_array_column.set_is_nullable(true);
+ array_column.set_name("a");
+ nested_array_column.set_name("item");
+ int_column.set_name("item");
+
+ meta.set_column_id(0);
+ meta.set_unique_id(0);
+ meta.set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_ARRAY));
+ meta.set_length(0);
+ meta.set_encoding(DEFAULT_ENCODING);
+ meta.set_compression(CompressionTypePB::LZ4F);
+ meta.set_is_nullable(true);
+ auto* nested_meta = meta.add_children_columns();
+ nested_meta->set_column_id(1);
+ nested_meta->set_unique_id(1);
+
nested_meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_ARRAY));
+ nested_meta->set_length(0);
+ nested_meta->set_encoding(DEFAULT_ENCODING);
+ nested_meta->set_compression(CompressionTypePB::LZ4F);
+ nested_meta->set_is_nullable(true);
+ auto* int_meta = nested_meta->add_children_columns();
+ int_meta->set_column_id(2);
+ int_meta->set_unique_id(2);
+ int_meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_INT));
+ int_meta->set_length(0);
+ int_meta->set_encoding(BIT_SHUFFLE);
+ int_meta->set_compression(CompressionTypePB::LZ4F);
+ int_meta->set_is_nullable(true);
+
+ std::vector<int32_t> int_values;
+ std::vector<uint8_t> int_null_map;
+ std::vector<uint64_t> nested_offsets {0};
+ std::vector<uint8_t> nested_item_null_map;
+ std::vector<uint64_t> outer_offsets(num_rows + 1, 0);
+ std::vector<uint8_t> outer_null_map(num_rows, 0);
+ for (size_t row = 0; row < num_rows; ++row) {
+ const bool outer_is_null = row % 7 == 1;
+ outer_null_map[row] = outer_is_null;
+ const size_t item_count = row % 11 == 0 ? 0 : (outer_is_null ? 3 : row
% 3 + 1);
+ for (size_t item = 0; item < item_count; ++item) {
+ const bool item_is_null = (row + item) % 5 == 0;
+ nested_item_null_map.push_back(item_is_null);
+ const size_t inner_count = item_is_null ? 2 : (row + item) % 3;
+ for (size_t inner = 0; inner < inner_count; ++inner) {
+ int_values.push_back(static_cast<int32_t>(row * 10 + item * 3
+ inner));
+ int_null_map.push_back((row + item + inner) % 4 == 0);
+ }
+ nested_offsets.push_back(int_values.size());
+ }
+ outer_offsets[row + 1] = nested_item_null_map.size();
+ }
+
+ std::vector<uint64_t> item_data {static_cast<uint64_t>(int_values.size()),
+
reinterpret_cast<uint64_t>(nested_offsets.data()),
+
reinterpret_cast<uint64_t>(int_values.data()),
+
reinterpret_cast<uint64_t>(int_null_map.data())};
+ const std::string file_name =
+ COLUMN_READER_FILE_TEST_DIR + "/array_read_by_rowids_nested_array";
+ write_array_column(file_name, &meta, array_column, num_rows, item_data,
outer_offsets,
+ nested_item_null_map, &outer_null_map);
+
+ io::FileReaderSPtr file_reader;
+ auto st = io::global_local_filesystem()->open_file(file_name,
&file_reader);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ ColumnReaderOptions reader_options;
+ std::shared_ptr<ColumnReader> reader;
+ st = ColumnReader::create(reader_options, meta, num_rows, file_reader,
&reader);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+
+ DataTypePtr int_type =
std::make_shared<DataTypeNullable>(std::make_shared<DataTypeInt32>());
+ DataTypePtr nested_array_type = std::make_shared<DataTypeArray>(int_type);
+ DataTypePtr outer_item_type =
std::make_shared<DataTypeNullable>(nested_array_type);
+ DataTypePtr outer_array_type =
std::make_shared<DataTypeArray>(outer_item_type);
+ DataTypePtr column_type =
std::make_shared<DataTypeNullable>(outer_array_type);
+ MutableColumnPtr baseline;
+ read_array_baseline(reader, array_column, column_type, file_reader.get(),
&baseline);
+ const std::vector<rowid_t> rowids {0, 1, 2, 7, 11, 4095, 4096, 8191,
11998, 11999};
+ check_array_read_by_rowids_matches_baseline(reader, array_column,
column_type, rowids,
+ *baseline, file_reader.get());
+}
+
+TEST_F(ColumnReaderTest, ArrayReadByRowidsNestedStructMatchesSequentialRead) {
+ constexpr size_t num_rows = 12000;
+ ColumnMetaPB meta;
+ TabletColumn
array_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_ARRAY);
+ TabletColumn
struct_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_STRUCT);
+ TabletColumn
int_a_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_INT, true);
+ TabletColumn
int_b_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_INT, true);
+ struct_column.add_sub_column(int_a_column);
+ struct_column.add_sub_column(int_b_column);
+ array_column.add_sub_column(struct_column);
+ array_column.set_is_nullable(true);
+ struct_column.set_is_nullable(true);
+ array_column.set_name("a");
+ struct_column.set_name("item");
+ int_a_column.set_name("a");
+ int_b_column.set_name("b");
+
+ meta.set_column_id(0);
+ meta.set_unique_id(0);
+ meta.set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_ARRAY));
+ meta.set_length(0);
+ meta.set_encoding(DEFAULT_ENCODING);
+ meta.set_compression(CompressionTypePB::LZ4F);
+ meta.set_is_nullable(true);
+ auto* struct_meta = meta.add_children_columns();
+ struct_meta->set_column_id(1);
+ struct_meta->set_unique_id(1);
+
struct_meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_STRUCT));
+ struct_meta->set_length(0);
+ struct_meta->set_encoding(DEFAULT_ENCODING);
+ struct_meta->set_compression(CompressionTypePB::LZ4F);
+ struct_meta->set_is_nullable(true);
+ auto* a_meta = struct_meta->add_children_columns();
+ a_meta->set_column_id(2);
+ a_meta->set_unique_id(2);
+ a_meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_INT));
+ a_meta->set_length(0);
+ a_meta->set_encoding(BIT_SHUFFLE);
+ a_meta->set_compression(CompressionTypePB::LZ4F);
+ a_meta->set_is_nullable(true);
+ auto* b_meta = struct_meta->add_children_columns();
+ b_meta->set_column_id(3);
+ b_meta->set_unique_id(3);
+ b_meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_INT));
+ b_meta->set_length(0);
+ b_meta->set_encoding(BIT_SHUFFLE);
+ b_meta->set_compression(CompressionTypePB::LZ4F);
+ b_meta->set_is_nullable(true);
+
+ std::vector<int32_t> a_values;
+ std::vector<int32_t> b_values;
+ std::vector<uint8_t> a_null_map;
+ std::vector<uint8_t> b_null_map;
+ std::vector<uint8_t> struct_item_null_map;
+ std::vector<uint64_t> outer_offsets(num_rows + 1, 0);
+ std::vector<uint8_t> outer_null_map(num_rows, 0);
+ for (size_t row = 0; row < num_rows; ++row) {
+ const bool outer_is_null = row % 7 == 1;
+ outer_null_map[row] = outer_is_null;
+ const size_t item_count = row % 11 == 0 ? 0 : (outer_is_null ? 3 : row
% 3 + 1);
+ for (size_t item = 0; item < item_count; ++item) {
+ struct_item_null_map.push_back((row + item) % 5 == 0);
+ a_values.push_back(static_cast<int32_t>(row * 10 + item));
+ b_values.push_back(static_cast<int32_t>(row * 100 + item));
+ a_null_map.push_back((row + item) % 4 == 0);
+ b_null_map.push_back((row + item * 2) % 5 == 0);
+ }
+ outer_offsets[row + 1] = struct_item_null_map.size();
+ }
+
+ std::vector<uint64_t> item_data
{reinterpret_cast<uint64_t>(a_values.data()),
+
reinterpret_cast<uint64_t>(b_values.data()),
+
reinterpret_cast<uint64_t>(a_null_map.data()),
+
reinterpret_cast<uint64_t>(b_null_map.data())};
+ const std::string file_name =
+ COLUMN_READER_FILE_TEST_DIR +
"/array_read_by_rowids_nested_struct";
+ write_array_column(file_name, &meta, array_column, num_rows, item_data,
outer_offsets,
+ struct_item_null_map, &outer_null_map);
+
+ io::FileReaderSPtr file_reader;
+ auto st = io::global_local_filesystem()->open_file(file_name,
&file_reader);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ ColumnReaderOptions reader_options;
+ std::shared_ptr<ColumnReader> reader;
+ st = ColumnReader::create(reader_options, meta, num_rows, file_reader,
&reader);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+
+ DataTypePtr a_type =
std::make_shared<DataTypeNullable>(std::make_shared<DataTypeInt32>());
+ DataTypePtr b_type =
std::make_shared<DataTypeNullable>(std::make_shared<DataTypeInt32>());
+ DataTypePtr struct_type = std::make_shared<DataTypeStruct>(
+ std::vector<DataTypePtr> {a_type, b_type},
std::vector<std::string> {"a", "b"});
+ DataTypePtr outer_item_type =
std::make_shared<DataTypeNullable>(struct_type);
+ DataTypePtr outer_array_type =
std::make_shared<DataTypeArray>(outer_item_type);
+ DataTypePtr column_type =
std::make_shared<DataTypeNullable>(outer_array_type);
+ MutableColumnPtr baseline;
+ read_array_baseline(reader, array_column, column_type, file_reader.get(),
&baseline);
+ const std::vector<rowid_t> rowids {0, 1, 2, 7, 11, 4095, 4096, 8191,
11998, 11999};
+ check_array_read_by_rowids_matches_baseline(reader, array_column,
column_type, rowids,
+ *baseline, file_reader.get());
+}
+
+TEST_F(ColumnReaderTest, ArrayReadByRowidsNestedMapMatchesSequentialRead) {
+ constexpr size_t num_rows = 12000;
+ ColumnMetaPB meta;
+ TabletColumn
array_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_ARRAY);
+ TabletColumn
map_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_MAP);
+ TabletColumn
key_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_INT);
+ TabletColumn
value_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_INT, true);
+ map_column.add_sub_column(key_column);
+ map_column.add_sub_column(value_column);
+ array_column.add_sub_column(map_column);
+ array_column.set_is_nullable(true);
+ map_column.set_is_nullable(true);
+ array_column.set_name("a");
+ map_column.set_name("item");
+ key_column.set_name("key");
+ value_column.set_name("value");
+
+ meta.set_column_id(0);
+ meta.set_unique_id(0);
+ meta.set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_ARRAY));
+ meta.set_length(0);
+ meta.set_encoding(DEFAULT_ENCODING);
+ meta.set_compression(CompressionTypePB::LZ4F);
+ meta.set_is_nullable(true);
+ auto* map_meta = meta.add_children_columns();
+ map_meta->set_column_id(1);
+ map_meta->set_unique_id(1);
+ map_meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_MAP));
+ map_meta->set_length(0);
+ map_meta->set_encoding(DEFAULT_ENCODING);
+ map_meta->set_compression(CompressionTypePB::LZ4F);
+ map_meta->set_is_nullable(true);
+ auto* key_meta = map_meta->add_children_columns();
+ key_meta->set_column_id(2);
+ key_meta->set_unique_id(2);
+ key_meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_INT));
+ key_meta->set_length(0);
+ key_meta->set_encoding(BIT_SHUFFLE);
+ key_meta->set_compression(CompressionTypePB::LZ4F);
+ key_meta->set_is_nullable(false);
+ auto* value_meta = map_meta->add_children_columns();
+ value_meta->set_column_id(3);
+ value_meta->set_unique_id(3);
+ value_meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_INT));
+ value_meta->set_length(0);
+ value_meta->set_encoding(BIT_SHUFFLE);
+ value_meta->set_compression(CompressionTypePB::LZ4F);
+ value_meta->set_is_nullable(true);
+
+ std::vector<int32_t> key_values;
+ std::vector<int32_t> value_values;
+ std::vector<uint8_t> key_null_map;
+ std::vector<uint8_t> value_null_map;
+ std::vector<uint64_t> map_offsets {0};
+ std::vector<uint8_t> map_item_null_map;
+ std::vector<uint64_t> outer_offsets(num_rows + 1, 0);
+ std::vector<uint8_t> outer_null_map(num_rows, 0);
+ for (size_t row = 0; row < num_rows; ++row) {
+ const bool outer_is_null = row % 7 == 1;
+ outer_null_map[row] = outer_is_null;
+ const size_t item_count = row % 11 == 0 ? 0 : (outer_is_null ? 3 : row
% 3 + 1);
+ for (size_t item = 0; item < item_count; ++item) {
+ map_item_null_map.push_back((row + item) % 5 == 0);
+ const size_t kv_count = (row + item) % 3;
+ for (size_t kv = 0; kv < kv_count; ++kv) {
+ key_values.push_back(static_cast<int32_t>(row * 10 + item * 3
+ kv));
+ value_values.push_back(static_cast<int32_t>(row * 100 + item *
7 + kv));
+ key_null_map.push_back(0);
+ value_null_map.push_back((row + item + kv) % 4 == 0);
+ }
+ map_offsets.push_back(key_values.size());
+ }
+ outer_offsets[row + 1] = map_item_null_map.size();
+ }
+
+ std::vector<uint64_t> item_data {static_cast<uint64_t>(key_values.size()),
+
reinterpret_cast<uint64_t>(map_offsets.data()),
+
reinterpret_cast<uint64_t>(key_values.data()),
+
reinterpret_cast<uint64_t>(value_values.data()),
+
reinterpret_cast<uint64_t>(key_null_map.data()),
+
reinterpret_cast<uint64_t>(value_null_map.data())};
+ const std::string file_name = COLUMN_READER_FILE_TEST_DIR +
"/array_read_by_rowids_nested_map";
+ write_array_column(file_name, &meta, array_column, num_rows, item_data,
outer_offsets,
+ map_item_null_map, &outer_null_map);
+
+ io::FileReaderSPtr file_reader;
+ auto st = io::global_local_filesystem()->open_file(file_name,
&file_reader);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ ColumnReaderOptions reader_options;
+ std::shared_ptr<ColumnReader> reader;
+ st = ColumnReader::create(reader_options, meta, num_rows, file_reader,
&reader);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+
+ DataTypePtr key_type = std::make_shared<DataTypeInt32>();
+ DataTypePtr value_type =
std::make_shared<DataTypeNullable>(std::make_shared<DataTypeInt32>());
+ DataTypePtr map_type = std::make_shared<DataTypeMap>(key_type, value_type);
+ DataTypePtr outer_item_type = std::make_shared<DataTypeNullable>(map_type);
+ DataTypePtr outer_array_type =
std::make_shared<DataTypeArray>(outer_item_type);
+ DataTypePtr column_type =
std::make_shared<DataTypeNullable>(outer_array_type);
+ MutableColumnPtr baseline;
+ read_array_baseline(reader, array_column, column_type, file_reader.get(),
&baseline);
+ const std::vector<rowid_t> rowids {0, 1, 2, 7, 11, 4095, 4096, 8191,
11998, 11999};
+ check_array_read_by_rowids_matches_baseline(reader, array_column,
column_type, rowids,
+ *baseline, file_reader.get());
+}
+
+TEST_F(ColumnReaderTest,
ArrayReadByRowidsNestedArraySchemaEvolutionFromNonNullSource) {
+ constexpr size_t num_rows = 8;
+ const std::array<size_t, num_rows> item_counts {2, 0, 3, 1, 4, 0, 2, 3};
+ ColumnMetaPB meta;
+ TabletColumn
array_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_ARRAY);
+ TabletColumn
nested_array_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_ARRAY);
+ TabletColumn
int_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_INT, true);
+ nested_array_column.add_sub_column(int_column);
+ array_column.add_sub_column(nested_array_column);
+ array_column.set_name("a");
+ nested_array_column.set_name("item");
+ int_column.set_name("item");
+
+ meta.set_column_id(0);
+ meta.set_unique_id(0);
+ meta.set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_ARRAY));
+ meta.set_length(0);
+ meta.set_encoding(DEFAULT_ENCODING);
+ meta.set_compression(CompressionTypePB::LZ4F);
+ meta.set_is_nullable(false);
+ auto* nested_meta = meta.add_children_columns();
+ nested_meta->set_column_id(1);
+ nested_meta->set_unique_id(1);
+
nested_meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_ARRAY));
+ nested_meta->set_length(0);
+ nested_meta->set_encoding(DEFAULT_ENCODING);
+ nested_meta->set_compression(CompressionTypePB::LZ4F);
+ nested_meta->set_is_nullable(true);
+ auto* int_meta = nested_meta->add_children_columns();
+ int_meta->set_column_id(2);
+ int_meta->set_unique_id(2);
+ int_meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_INT));
+ int_meta->set_length(0);
+ int_meta->set_encoding(BIT_SHUFFLE);
+ int_meta->set_compression(CompressionTypePB::LZ4F);
+ int_meta->set_is_nullable(true);
+
+ std::vector<int32_t> int_values;
+ std::vector<uint8_t> int_null_map;
+ std::vector<uint64_t> nested_offsets {0};
+ std::vector<uint8_t> nested_item_null_map;
+ std::vector<uint64_t> outer_offsets(num_rows + 1, 0);
+ for (size_t row = 0; row < num_rows; ++row) {
+ for (size_t item = 0; item < item_counts[row]; ++item) {
+ nested_item_null_map.push_back((row + item) % 5 == 0);
+ const size_t inner_count = (row + item) % 3;
+ for (size_t inner = 0; inner < inner_count; ++inner) {
+ int_values.push_back(static_cast<int32_t>(row * 10 + item * 3
+ inner));
+ int_null_map.push_back((row + item + inner) % 4 == 0);
+ }
+ nested_offsets.push_back(int_values.size());
+ }
+ outer_offsets[row + 1] = nested_item_null_map.size();
+ }
+
+ std::vector<uint64_t> item_data {static_cast<uint64_t>(int_values.size()),
+
reinterpret_cast<uint64_t>(nested_offsets.data()),
+
reinterpret_cast<uint64_t>(int_values.data()),
+
reinterpret_cast<uint64_t>(int_null_map.data())};
+ const std::string file_name =
+ COLUMN_READER_FILE_TEST_DIR +
"/array_read_by_rowids_nested_array_schema_evolution";
+ write_array_column(file_name, &meta, array_column, num_rows, item_data,
outer_offsets,
+ nested_item_null_map, nullptr);
+
+ io::FileReaderSPtr file_reader;
+ auto st = io::global_local_filesystem()->open_file(file_name,
&file_reader);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ ColumnReaderOptions reader_options;
+ std::shared_ptr<ColumnReader> reader;
+ st = ColumnReader::create(reader_options, meta, num_rows, file_reader,
&reader);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+
+ DataTypePtr int_type =
std::make_shared<DataTypeNullable>(std::make_shared<DataTypeInt32>());
+ DataTypePtr nested_array_type = std::make_shared<DataTypeArray>(int_type);
+ DataTypePtr outer_item_type =
std::make_shared<DataTypeNullable>(nested_array_type);
+ DataTypePtr outer_array_type =
std::make_shared<DataTypeArray>(outer_item_type);
+ DataTypePtr column_type =
std::make_shared<DataTypeNullable>(outer_array_type);
+ MutableColumnPtr baseline;
+ read_array_baseline(reader, array_column, column_type, file_reader.get(),
&baseline);
+ const std::vector<rowid_t> rowids {0, 2, 3, 6, 7};
+ check_array_read_by_rowids_matches_baseline(reader, array_column,
column_type, rowids,
+ *baseline, file_reader.get());
+}
+
+TEST_F(ColumnReaderTest, ArrayReadByRowidsSchemaEvolutionFromNonNullSource) {
+ constexpr size_t num_rows = 8;
+ const std::array<size_t, num_rows> item_counts {2, 0, 3, 1, 4, 0, 2, 3};
+ ColumnMetaPB meta;
+ TabletColumn
array_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_ARRAY);
+ TabletColumn
item_column(FieldAggregationMethod::OLAP_FIELD_AGGREGATION_NONE,
+ FieldType::OLAP_FIELD_TYPE_INT, false);
+ array_column.add_sub_column(item_column);
+ array_column.set_name("a");
+
+ std::vector<int32_t> item_values;
+ std::vector<uint64_t> array_offsets(num_rows + 1, 0);
+ for (size_t row = 0; row < num_rows; ++row) {
+ for (size_t item = 0; item < item_counts[row]; ++item) {
+ item_values.push_back(static_cast<int32_t>(row * 100 + item));
+ }
+ array_offsets[row + 1] = item_values.size();
+ }
+
+ const std::string file_name =
+ COLUMN_READER_FILE_TEST_DIR +
"/array_read_by_rowids_schema_evolution";
+ auto fs = io::global_local_filesystem();
+ {
+ io::FileWriterPtr file_writer;
+ auto st = fs->create_file(file_name, &file_writer);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+
+ ColumnWriterOptions writer_options;
+ writer_options.meta = &meta;
+ writer_options.meta->set_column_id(0);
+ writer_options.meta->set_unique_id(0);
+
writer_options.meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_ARRAY));
+ writer_options.meta->set_length(0);
+ writer_options.meta->set_encoding(DEFAULT_ENCODING);
+ writer_options.meta->set_compression(CompressionTypePB::LZ4F);
+ writer_options.meta->set_is_nullable(false);
+
+ auto* child_meta = meta.add_children_columns();
+ child_meta->set_column_id(1);
+ child_meta->set_unique_id(1);
+
child_meta->set_type(static_cast<int32_t>(FieldType::OLAP_FIELD_TYPE_INT));
+ child_meta->set_length(0);
+ child_meta->set_encoding(BIT_SHUFFLE);
+ child_meta->set_compression(CompressionTypePB::LZ4F);
+ child_meta->set_is_nullable(false);
+
+ std::unique_ptr<ColumnWriter> writer;
+ st = ColumnWriter::create(writer_options, &array_column,
file_writer.get(), &writer);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ ASSERT_TRUE(writer->init().ok());
+ const std::array<uint64_t, 4> array_data
{static_cast<uint64_t>(item_values.size()),
+
reinterpret_cast<uint64_t>(array_offsets.data()),
+
reinterpret_cast<uint64_t>(item_values.data()),
+ 0};
+ ASSERT_TRUE(writer->append(nullptr, array_data.data(), num_rows).ok());
+ ASSERT_TRUE(writer->finish().ok());
+ ASSERT_TRUE(writer->write_data().ok());
+ ASSERT_TRUE(writer->write_ordinal_index().ok());
+ ASSERT_TRUE(file_writer->close().ok());
+ }
+
+ io::FileReaderSPtr file_reader;
+ auto st = fs->open_file(file_name, &file_reader);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+ ColumnReaderOptions reader_options;
+ std::shared_ptr<ColumnReader> reader;
+ st = ColumnReader::create(reader_options, meta, num_rows, file_reader,
&reader);
+ ASSERT_TRUE(st.ok()) << st.to_string();
+
+ DataTypePtr array_type =
std::make_shared<DataTypeArray>(std::make_shared<DataTypeInt32>());
+ DataTypePtr column_type = std::make_shared<DataTypeNullable>(array_type);
+ auto create_iterator = [&](OlapReaderStatistics* stats,
+ ColumnIteratorUPtr* iterator) -> Status {
+ RETURN_IF_ERROR(reader->new_iterator(iterator, &array_column));
+ ColumnIteratorOptions options;
+ options.stats = stats;
+ options.file_reader = file_reader.get();
+ return (*iterator)->init(options);
+ };
+
+ MutableColumnPtr baseline = column_type->create_column();
+ {
+ ColumnIteratorUPtr iterator;
+ OlapReaderStatistics stats;
+ ASSERT_TRUE(create_iterator(&stats, &iterator).ok());
+ ASSERT_TRUE(iterator->seek_to_ordinal(0).ok());
+ size_t rows_to_read = num_rows;
+ bool has_null = false;
+ ASSERT_TRUE(iterator->next_batch(&rows_to_read, baseline,
&has_null).ok());
+ ASSERT_EQ(num_rows, rows_to_read);
+ EXPECT_FALSE(has_null);
+ }
+
+ const std::vector<rowid_t> rowids {0, 2, 3, 6, 7};
+ MutableColumnPtr actual = column_type->create_column();
+ {
+ ColumnIteratorUPtr iterator;
+ OlapReaderStatistics stats;
+ ASSERT_TRUE(create_iterator(&stats, &iterator).ok());
+ ASSERT_TRUE(iterator->read_by_rowids(rowids.data(), rowids.size(),
actual).ok());
+ }
+
+ const auto& baseline_nullable = assert_cast<const
ColumnNullable&>(*baseline);
+ const auto& baseline_array = assert_cast<const ColumnArray&,
TypeCheckOnRelease::DISABLE>(
+ baseline_nullable.get_nested_column());
+ const auto& actual_nullable = assert_cast<const ColumnNullable&>(*actual);
+ const auto& actual_array = assert_cast<const ColumnArray&,
TypeCheckOnRelease::DISABLE>(
+ actual_nullable.get_nested_column());
+ ASSERT_EQ(rowids.size(), actual->size());
+ EXPECT_THAT(actual_nullable.get_null_map_data(), ::testing::Each(0));
+
+ size_t expected_item_count = 0;
+ size_t item_cursor = 0;
+ for (size_t i = 0; i < rowids.size(); ++i) {
+ const auto rowid = rowids[i];
+ EXPECT_EQ(baseline_nullable.get_null_map_data()[rowid],
+ actual_nullable.get_null_map_data()[i]);
+ const size_t item_count = baseline_array.size_at(rowid);
+ expected_item_count += item_count;
+ EXPECT_EQ(expected_item_count, actual_array.get_offsets()[i]);
+ for (size_t item = 0; item < item_count; ++item) {
+ EXPECT_EQ(0, actual_array.get_data().compare_at(item_cursor,
+
baseline_array.offset_at(rowid) + item,
+
baseline_array.get_data(), 1));
+ ++item_cursor;
+ }
+ }
+ EXPECT_EQ(item_cursor, actual_array.get_data().size());
+
+ {
+ ColumnIteratorUPtr iterator;
+ OlapReaderStatistics stats;
+ ASSERT_TRUE(create_iterator(&stats, &iterator).ok());
+ iterator->set_column_name("a");
+ TColumnAccessPaths null_path {create_meta_access_path({"a",
ColumnIterator::ACCESS_NULL})};
+ ASSERT_TRUE(iterator->set_access_paths(null_path, null_path).ok());
+ EXPECT_TRUE(iterator->read_null_map_only());
+
+ MutableColumnPtr null_only = column_type->create_column();
+ ASSERT_TRUE(iterator->read_by_rowids(rowids.data(), rowids.size(),
null_only).ok());
+ const auto& null_only_nullable = assert_cast<const
ColumnNullable&>(*null_only);
+ const auto& null_only_array = assert_cast<const ColumnArray&,
TypeCheckOnRelease::DISABLE>(
+ null_only_nullable.get_nested_column());
+ ASSERT_EQ(rowids.size(), null_only->size());
+ EXPECT_THAT(null_only_nullable.get_null_map_data(),
::testing::Each(0));
+ EXPECT_THAT(null_only_array.get_offsets(), ::testing::Each(0));
+ EXPECT_TRUE(null_only_array.get_data().empty());
+ }
+}
+
TEST_F(ColumnReaderTest, StructAccessPaths) {
auto create_struct_iterator = []() {
auto null_reader = std::make_shared<ColumnReader>();
@@ -2982,6 +3876,162 @@ TEST_F(ColumnReaderTest,
StructNullMapOnlyNextBatchSkipsSubColumns) {
EXPECT_EQ(2, nested_struct.get_column(0).size());
}
+TEST_F(ColumnReaderTest, ArrayReadByRowidsBatchesOffsetsAndMergesItemRanges) {
+ auto offset_file_iterator =
std::make_unique<RowidOffsetFileColumnIterator>(
+ std::vector<ordinal_t> {0, 2, 5, 5, 7, 10, 12, 12, 15, 18, 20});
+ auto* offset_tracker = offset_file_iterator.get();
+ auto offset_iterator =
+
std::make_unique<OffsetFileColumnIterator>(std::move(offset_file_iterator));
+ auto item_iterator = std::make_unique<TrackingColumnIterator>();
+ auto* item_tracker = item_iterator.get();
+
+ ArrayFileColumnIterator array_iterator(
+ create_test_reader(false, 10, FieldType::OLAP_FIELD_TYPE_ARRAY),
+ std::move(offset_iterator), std::move(item_iterator), nullptr);
+ MutableColumnPtr dst =
+ ColumnArray::create(ColumnInt32::create(),
ColumnArray::ColumnOffsets::create());
+
+ // Row 2 is a selected empty array. Unselected row 6 is also empty, so
rows 5 and 7
+ // must share an item read even though their parent rowids are not
adjacent. Unselected
+ // rows 3 and 8 contain items and must still separate the three physical
item ranges.
+ const rowid_t rowids[] = {0, 1, 2, 4, 5, 7, 9};
+ auto st = array_iterator.read_by_rowids(rowids, std::size(rowids), dst);
+ ASSERT_TRUE(st.ok()) << "array read_by_rowids failed: " << st.to_string();
+
+ ASSERT_EQ(1, offset_tracker->read_by_rowids_batches.size());
+ EXPECT_THAT(offset_tracker->read_by_rowids_batches[0],
+ ::testing::ElementsAre(0, 1, 2, 3, 4, 5, 6, 7, 8, 9));
+ EXPECT_THAT(offset_tracker->seek_ordinals, ::testing::ElementsAre(9));
+ EXPECT_THAT(offset_tracker->next_batch_sizes, ::testing::ElementsAre(1));
+ EXPECT_THAT(item_tracker->seek_ordinals, ::testing::ElementsAre(0, 7, 18));
+ EXPECT_THAT(item_tracker->next_batch_sizes, ::testing::ElementsAre(5, 8,
2));
+
+ const auto& array = assert_cast<const ColumnArray&,
TypeCheckOnRelease::DISABLE>(*dst);
+ EXPECT_EQ(7, array.size());
+ EXPECT_EQ(15, array.get_data().size());
+ EXPECT_EQ(2, array.get_offsets()[0]);
+ EXPECT_EQ(5, array.get_offsets()[1]);
+ EXPECT_EQ(5, array.get_offsets()[2]);
+ EXPECT_EQ(8, array.get_offsets()[3]);
+ EXPECT_EQ(10, array.get_offsets()[4]);
+ EXPECT_EQ(13, array.get_offsets()[5]);
+ EXPECT_EQ(15, array.get_offsets()[6]);
+}
+
+TEST_F(ColumnReaderTest, ArrayReadByRowidsSingleLastRowAppendsToExistingData) {
+ auto offset_file_iterator =
std::make_unique<RowidOffsetFileColumnIterator>(
+ std::vector<ordinal_t> {0, 2, 5, 5, 7, 10, 12, 12, 15, 18, 20});
+ auto* offset_tracker = offset_file_iterator.get();
+ auto offset_iterator =
+
std::make_unique<OffsetFileColumnIterator>(std::move(offset_file_iterator));
+ auto item_iterator = std::make_unique<TrackingColumnIterator>();
+ auto* item_tracker = item_iterator.get();
+ ArrayFileColumnIterator array_iterator(
+ create_test_reader(false, 10, FieldType::OLAP_FIELD_TYPE_ARRAY),
+ std::move(offset_iterator), std::move(item_iterator), nullptr);
+
+ auto values = ColumnInt32::create();
+ values->insert_value(42);
+ auto offsets = ColumnArray::ColumnOffsets::create();
+ offsets->insert_value(1);
+ MutableColumnPtr dst = ColumnArray::create(std::move(values),
std::move(offsets));
+ const rowid_t rowid = 9;
+ ASSERT_TRUE(array_iterator.read_by_rowids(&rowid, 1, dst).ok());
+
+ ASSERT_EQ(1, offset_tracker->read_by_rowids_batches.size());
+ EXPECT_THAT(offset_tracker->read_by_rowids_batches[0],
::testing::ElementsAre(9));
+ EXPECT_THAT(offset_tracker->seek_ordinals, ::testing::ElementsAre(9));
+ EXPECT_THAT(offset_tracker->next_batch_sizes, ::testing::ElementsAre(1));
+ EXPECT_THAT(item_tracker->seek_ordinals, ::testing::ElementsAre(18));
+ EXPECT_THAT(item_tracker->next_batch_sizes, ::testing::ElementsAre(2));
+ const auto& array = assert_cast<const ColumnArray&>(*dst);
+ EXPECT_EQ(2, array.size());
+ EXPECT_THAT(array.get_offsets(), ::testing::ElementsAre(1, 3));
+ EXPECT_THAT(assert_cast<const ColumnInt32&>(array.get_data()).get_data(),
+ ::testing::ElementsAre(42, 0, 0));
+
+ ASSERT_TRUE(array_iterator.read_by_rowids(nullptr, 0, dst).ok());
+ EXPECT_EQ(2, dst->size());
+ EXPECT_EQ(1, offset_tracker->read_by_rowids_batches.size());
+ EXPECT_EQ(1, item_tracker->seek_ordinals.size());
+}
+
+TEST_F(ColumnReaderTest,
ArrayOffsetOnlyReadByRowidsFillsItemsWithoutReadingThem) {
+ auto offset_file_iterator =
std::make_unique<RowidOffsetFileColumnIterator>(
+ std::vector<ordinal_t> {0, 2, 5, 5, 7, 10, 12, 12, 15, 18, 20});
+ auto* offset_tracker = offset_file_iterator.get();
+ auto offset_iterator =
+
std::make_unique<OffsetFileColumnIterator>(std::move(offset_file_iterator));
+ auto item_iterator = std::make_unique<TrackingColumnIterator>();
+ auto* item_tracker = item_iterator.get();
+
+ ArrayFileColumnIterator array_iterator(
+ create_test_reader(false, 10, FieldType::OLAP_FIELD_TYPE_ARRAY),
+ std::move(offset_iterator), std::move(item_iterator), nullptr);
+ array_iterator.set_column_name("a");
+ TColumnAccessPaths offset_path {create_meta_access_path({"a",
ColumnIterator::ACCESS_OFFSET})};
+ auto st = array_iterator.set_access_paths(offset_path, {});
+ ASSERT_TRUE(st.ok()) << "set_access_paths failed: " << st.to_string();
+ ASSERT_TRUE(array_iterator.read_offset_only());
+
+ MutableColumnPtr dst =
+ ColumnArray::create(ColumnInt32::create(),
ColumnArray::ColumnOffsets::create());
+ const rowid_t rowids[] = {0, 1, 4};
+ st = array_iterator.read_by_rowids(rowids, std::size(rowids), dst);
+ ASSERT_TRUE(st.ok()) << "array offset-only read_by_rowids failed: " <<
st.to_string();
+
+ const auto& array = assert_cast<const ColumnArray&,
TypeCheckOnRelease::DISABLE>(*dst);
+ EXPECT_EQ(3, array.size());
+ EXPECT_EQ(8, array.get_data().size());
+ EXPECT_EQ(2, array.get_offsets()[0]);
+ EXPECT_EQ(5, array.get_offsets()[1]);
+ EXPECT_EQ(8, array.get_offsets()[2]);
+ EXPECT_TRUE(item_tracker->seek_ordinals.empty());
+ EXPECT_TRUE(item_tracker->next_batch_sizes.empty());
+ ASSERT_EQ(1, offset_tracker->read_by_rowids_batches.size());
+ EXPECT_THAT(offset_tracker->read_by_rowids_batches[0],
::testing::ElementsAre(0, 1, 2, 4, 5));
+}
+
+TEST_F(ColumnReaderTest, ArrayLazyReadByRowidsPreservesMaterializedOffsets) {
+ auto offset_file_iterator =
std::make_unique<RowidOffsetFileColumnIterator>(
+ std::vector<ordinal_t> {0, 2, 5, 5, 7, 10, 12, 12, 15, 18, 20});
+ auto offset_iterator =
+
std::make_unique<OffsetFileColumnIterator>(std::move(offset_file_iterator));
+ auto item_iterator = std::make_unique<TrackingColumnIterator>();
+ auto* item_tracker = item_iterator.get();
+
+ ArrayFileColumnIterator array_iterator(
+ create_test_reader(false, 10, FieldType::OLAP_FIELD_TYPE_ARRAY),
+ std::move(offset_iterator), std::move(item_iterator), nullptr);
+
array_iterator.set_read_requirement_self(ColumnIterator::ReadRequirement::PREDICATE);
+
item_tracker->set_read_requirement(ColumnIterator::ReadRequirement::LAZY_OUTPUT);
+ array_iterator.set_read_phase(ColumnIterator::ReadPhase::PREDICATE);
+
+ MutableColumnPtr dst =
+ ColumnArray::create(ColumnInt32::create(),
ColumnArray::ColumnOffsets::create());
+ auto& array = assert_cast<ColumnArray&, TypeCheckOnRelease::DISABLE>(*dst);
+ array.get_offsets().push_back(2);
+ array.get_offsets().push_back(5);
+ auto items = IColumn::mutate(std::move(array.get_data_ptr()));
+ item_tracker->convert_to_place_holder_column(items, 5);
+ array.get_data_ptr() = std::move(items);
+
+ array_iterator.set_read_phase(ColumnIterator::ReadPhase::LAZY);
+ ASSERT_TRUE(array_iterator.need_to_read());
+ ASSERT_FALSE(array_iterator.need_to_read_meta_columns());
+
+ const rowid_t rowids[] = {0, 4};
+ auto st = array_iterator.read_by_rowids(rowids, std::size(rowids), dst);
+ ASSERT_TRUE(st.ok()) << "lazy array read_by_rowids failed: " <<
st.to_string();
+
+ EXPECT_EQ(2, array.size());
+ EXPECT_EQ(5, array.get_data().size());
+ EXPECT_EQ(2, array.get_offsets()[0]);
+ EXPECT_EQ(5, array.get_offsets()[1]);
+ EXPECT_THAT(item_tracker->seek_ordinals, ::testing::ElementsAre(0, 7));
+ EXPECT_THAT(item_tracker->next_batch_sizes, ::testing::ElementsAre(2, 3));
+}
+
TEST_F(ColumnReaderTest, ArrayNullMapOnlyNextBatchAndReadByRowidsSkipItems) {
auto null_iterator = std::make_unique<TrackingColumnIterator>();
auto* null_iterator_ptr = null_iterator.get();
@@ -3018,8 +4068,10 @@ TEST_F(ColumnReaderTest,
ArrayNullMapOnlyNextBatchAndReadByRowidsSkipItems) {
st = array_iterator.read_by_rowids(rowids, std::size(rowids), dst);
ASSERT_TRUE(st.ok()) << "array null-map-only read_by_rowids failed: " <<
st.to_string();
EXPECT_EQ(5, dst->size());
- EXPECT_THAT(null_iterator_ptr->seek_ordinals, ::testing::ElementsAre(1,
3));
- EXPECT_THAT(null_iterator_ptr->next_batch_sizes, ::testing::ElementsAre(1,
1));
+ EXPECT_TRUE(null_iterator_ptr->seek_ordinals.empty());
+ EXPECT_TRUE(null_iterator_ptr->next_batch_sizes.empty());
+ ASSERT_EQ(1, null_iterator_ptr->read_by_rowids_batches.size());
+ EXPECT_THAT(null_iterator_ptr->read_by_rowids_batches[0],
::testing::ElementsAre(1, 3));
EXPECT_TRUE(item_iterator_ptr->next_batch_sizes.empty());
EXPECT_TRUE(offset_iterator.tracker->next_batch_sizes.empty());
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]