diff --git a/be/src/format/parquet/vparquet_column_reader.cpp b/be/src/format/parquet/vparquet_column_reader.cpp index 5fd9adcff06918..8797b54c0a6471 100644 --- a/be/src/format/parquet/vparquet_column_reader.cpp +++ b/be/src/format/parquet/vparquet_column_reader.cpp @@ -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, @@ -539,11 +540,26 @@ Status ScalarColumnReader::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(); @@ -552,8 +568,7 @@ Status ScalarColumnReader::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; @@ -631,8 +646,7 @@ Status ScalarColumnReader::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(); } diff --git a/be/test/format/parquet/parquet_column_chunk_reader_test.cpp b/be/test/format/parquet/parquet_column_chunk_reader_test.cpp index be9616c523f638..d8df938095446d 100644 --- a/be/test/format/parquet/parquet_column_chunk_reader_test.cpp +++ b/be/test/format/parquet/parquet_column_chunk_reader_test.cpp @@ -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" @@ -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(std::move(fixture.data)); + auto string_type = std::make_shared(); + 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 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()) << 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);