Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 18 additions & 4 deletions be/src/format/parquet/vparquet_column_reader.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@
#include "format/table/iceberg_default_value.h"
#include "io/fs/tracing_file_reader.h"
#include "runtime/runtime_profile.h"
#include "util/defer_op.h"

namespace doris {
static void fill_struct_null_map(FieldSchema* field, NullMap& null_map,
Expand Down Expand Up @@ -539,11 +540,26 @@ Status ScalarColumnReader<IN_COLLECTION, OFFSET_INDEX>::read_column_data(
ColumnPtr resolved_column =
_converter->get_physical_column(_field_schema->physical_type, _field_schema->data_type,
doris_column, type, is_dict_filter);
// Direct reads transfer the caller's only ColumnPtr so mutate() can avoid cloning it. Restore
// that ownership if any read step returns before convert() moves the column back.
bool restore_doris_column = false;
Defer restore_column([&]() {
if (restore_doris_column) {
doris_column = std::move(resolved_column);
}
});
if (_converter->read_directly_into_dst_logical_column()) {
DCHECK_EQ(resolved_column.get(), doris_column.get());
resolved_column = std::move(doris_column);
restore_doris_column = true;
}
DataTypePtr& resolved_type = _converter->get_physical_type();
auto convert_column = [&]() -> Status {
RETURN_IF_ERROR(_converter->convert(resolved_column, _field_schema->data_type, type,
doris_column, is_dict_filter));
restore_doris_column = false;
return Status::OK();
};

_def_levels.clear();
_rep_levels.clear();
Expand All @@ -552,8 +568,7 @@ Status ScalarColumnReader<IN_COLLECTION, OFFSET_INDEX>::read_column_data(
if (_in_nested) {
RETURN_IF_ERROR(_read_nested_column(resolved_column, resolved_type, filter_map, batch_size,
read_rows, eof, is_dict_filter));
return _converter->convert(resolved_column, _field_schema->data_type, type, doris_column,
is_dict_filter);
return convert_column();
}

int64_t right_row = 0;
Expand Down Expand Up @@ -631,8 +646,7 @@ Status ScalarColumnReader<IN_COLLECTION, OFFSET_INDEX>::read_column_data(

{
SCOPED_RAW_TIMER(&_convert_time);
RETURN_IF_ERROR(_converter->convert(resolved_column, _field_schema->data_type, type,
doris_column, is_dict_filter));
RETURN_IF_ERROR(convert_column());
}
return Status::OK();
}
Expand Down
37 changes: 37 additions & 0 deletions be/test/format/parquet/parquet_column_chunk_reader_test.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -27,11 +27,13 @@

#include "core/assert_cast.h"
#include "core/column/column_string.h"
#include "core/data_type/data_type_string.h"
#include "format/parquet/schema_desc.h"
#include "format/parquet/vparquet_column_chunk_reader.h"
#include "format/parquet/vparquet_column_reader.h"
#include "io/fs/buffered_reader.h"
#include "io/fs/file_reader.h"
#include "io/io_common.h"
#include "runtime/runtime_state.h"
#include "util/coding.h"
#include "util/thrift_util.h"
Expand Down Expand Up @@ -418,6 +420,41 @@ TEST(ParquetColumnChunkReaderTest, ScalarDictionaryReadUsesExplicitProbe) {
EXPECT_EQ(std::string(strings.get_data_at(2)), "carol");
}

TEST(ParquetColumnChunkReaderTest, ScalarNestedReadRestoresColumnWhenStopped) {
ColumnChunkFixture fixture;
ASSERT_TRUE(make_plain_fixture(&fixture).ok());
auto file_reader = std::make_shared<CountingFileReader>(std::move(fixture.data));
auto string_type = std::make_shared<DataTypeString>();
fixture.field_schema.data_type = string_type;
fixture.field_schema.parquet_schema.__set_type(tparquet::Type::BYTE_ARRAY);

RowRanges row_ranges;
row_ranges.add({0, 1});
io::IOContext io_ctx;
ScalarColumnReader<true, false> reader(row_ranges, 1, fixture.chunk, nullptr, nullptr, &io_ctx);
reader.set_column_in_nested();

TQueryOptions query_options;
query_options.__set_enable_parquet_file_page_cache(false);
RuntimeState runtime_state(query_options, TQueryGlobals());
ASSERT_TRUE(reader.init(file_reader, &fixture.field_schema,
/*max_buf_size=*/1024 * 1024, &runtime_state)
.ok());

ColumnPtr column = ColumnString::create();
FilterMap filter_map;
ASSERT_TRUE(filter_map.init(nullptr, 0, false).ok());
size_t read_rows = 0;
bool eof = false;
io_ctx.should_stop = true;

Status status = reader.read_column_data(column, string_type, nullptr, filter_map, 1, &read_rows,
&eof, false);
EXPECT_TRUE(status.is<ErrorCode::END_OF_FILE>()) << status;
ASSERT_NE(column, nullptr);
EXPECT_TRUE(column->empty());
}

void expect_offset_index_skip(ColumnChunkFixture fixture) {
CountingBufferedReader buffered_reader(std::move(fixture.data));
ParquetPageReadContext page_read_ctx(false);
Expand Down
Loading