From 109c5600223e0a1747d0e1511af4acecaf1f7aec Mon Sep 17 00:00:00 2001 From: daidai Date: Fri, 14 Aug 2026 18:23:32 +0800 Subject: [PATCH] [fix](parquet) Restore column ownership after interrupted reads ### What problem does this PR solve? Issue Number: None Related PR: None Problem Summary: The direct Parquet conversion path temporarily moves the caller-owned column into the physical read column. If an interrupted read returns before conversion, the caller is left with a null column pointer. Complex-column cleanup can then dereference that null child and crash the BE. Restore the transferred column on every early return, and disarm the restoration only after conversion moves ownership back successfully. Add a focused test that triggers the IO stop path and verifies the caller retains a valid column. ### Release note Fix an occasional BE crash when a Parquet scan is interrupted. ### Check List (For Author) - Test: Unit Test (added but not run; compilation was stopped at the user request) - Added `ParquetColumnChunkReaderTest.ScalarNestedReadRestoresColumnWhenStopped` - `build-support/check-format.sh` - `git diff --check` - Behavior changed: Yes. Interrupted direct Parquet reads now preserve caller column ownership. - Does this need documentation: No --- .../format/parquet/vparquet_column_reader.cpp | 22 +++++++++-- .../parquet_column_chunk_reader_test.cpp | 37 +++++++++++++++++++ 2 files changed, 55 insertions(+), 4 deletions(-) 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);