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 0f6b3e91482 [refactor](storage) introduce the block transform chain 
(#66528)
0f6b3e91482 is described below

commit 0f6b3e91482e0b6e5073c58456ed7e4fcf5eccef
Author: Chenyang Sun <[email protected]>
AuthorDate: Mon Aug 10 10:14:37 2026 +0800

    [refactor](storage) introduce the block transform chain (#66528)
    
    Third PR of the #64674 split stack (#65492, #66151). Adds
    storage/transform/block_transform.{h,cpp}: an immutable
    BlockTransformChain applied to every flushed block at the two
    segment_creator seams before the block reaches a segment writer.
    
    ### What problem does this PR solve?
    
    Issue Number: close #xxx
    
    Related PR: #xxx
    
    Problem Summary:
    
    ### Release note
    
    None
    
    ### Check List (For Author)
    
    - Test <!-- At least one of them must be included. -->
        - [ ] Regression test
        - [x] Unit Test
        - [ ] Manual test (add detailed scripts or steps below)
        - [ ] No need to test or manual test. Explain why:
    - [ ] This is a refactor/code format and no logic has been changed.
            - [ ] Previous test can cover this change.
            - [ ] No code files have been changed.
            - [ ] Other reason <!-- Add your reason?  -->
    
    - Behavior changed:
        - [ ] No.
        - [ ] Yes. <!-- Explain the behavior change -->
    
    - Does this need documentation?
        - [ ] No.
    - [ ] Yes. <!-- Add document PR link here. eg:
    https://github.com/apache/doris-website/pull/1214 -->
    
    ### Check List (For Reviewer who merge this PR)
    
    - [ ] Confirm the release note
    - [ ] Confirm test cases
    - [ ] Confirm document
    - [ ] Add branch pick label <!-- Add branch pick label that this PR
    should merge into -->
    
    Co-authored-by: Claude Fable 5 <[email protected]>
---
 be/src/storage/rowset/segment_creator.cpp          |  44 ++
 be/src/storage/rowset/segment_creator.h            |  11 +
 be/src/storage/segment/segment_writer.cpp          |  15 +-
 be/src/storage/segment/vertical_segment_writer.cpp |  75 +--
 be/src/storage/segment/vertical_segment_writer.h   |  14 +
 be/src/storage/transform/block_transform.cpp       | 233 +++++++++
 be/src/storage/transform/block_transform.h         | 129 +++++
 be/test/storage/mow/mow_transform_test_base.h      | 123 ++++-
 be/test/storage/transform/validate_stage_test.cpp  | 405 +++++++++++++++
 .../storage/transform/variant_rowstore_test.cpp    | 547 +++++++++++++++++++++
 10 files changed, 1548 insertions(+), 48 deletions(-)

diff --git a/be/src/storage/rowset/segment_creator.cpp 
b/be/src/storage/rowset/segment_creator.cpp
index d320d31256b..90078eca355 100644
--- a/be/src/storage/rowset/segment_creator.cpp
+++ b/be/src/storage/rowset/segment_creator.cpp
@@ -49,6 +49,7 @@
 #include "storage/segment/segment_writer.h"
 #include "storage/segment/vertical_segment_writer.h"
 #include "storage/tablet/tablet_schema.h"
+#include "storage/transform/block_transform.h"
 #include "storage/utils.h"
 #include "util/debug_points.h"
 #include "util/json/json_parser.h"
@@ -58,6 +59,23 @@
 namespace doris {
 using namespace ErrorCode;
 
+namespace {
+
+segment_v2::TransformExecContext 
make_transform_exec_context(RowsetWriterContext& context,
+                                                             int32_t 
segment_id) {
+    return {.tablet_schema = context.tablet_schema,
+            .write_type = context.write_type,
+            .tablet = context.tablet,
+            .mow_context = context.mow_context,
+            .partial_update_info = context.partial_update_info,
+            .rowset_ctx = &context,
+            .rowset_id = context.rowset_id,
+            .segment_id = segment_id,
+            .derived_column = {}};
+}
+
+} // namespace
+
 SegmentFlusher::SegmentFlusher(RowsetWriterContext& context, 
SegmentFileCollection& seg_files,
                                InvertedIndexFileCollection& idx_files)
         : _context(context), _seg_files(seg_files), _idx_files(idx_files) {}
@@ -72,14 +90,21 @@ Status SegmentFlusher::flush_single_block(const Block* 
block, int32_t segment_id
     }
     Block flush_block(*block);
     bool no_compression = flush_block.bytes() <= 
config::segment_compression_threshold_kb * 1024;
+    segment_v2::DerivedColumn derived_column;
+    RETURN_IF_ERROR(transform_block(&flush_block, segment_id, 
&derived_column));
     bool use_vertical_segment_writer =
             config::enable_vertical_segment_writer && 
!_context.write_binlog_opt().enable;
     if (use_vertical_segment_writer) {
         std::unique_ptr<segment_v2::VerticalSegmentWriter> writer;
         RETURN_IF_ERROR(_create_segment_writer(writer, segment_id, 
no_compression));
+        // the vertical writer feeds the derived column in small fixed-size 
batches
+        writer->set_derived_column(std::move(derived_column));
         RETURN_IF_ERROR_OR_CATCH_EXCEPTION(_add_rows(writer, &flush_block, 0, 
flush_block.rows()));
         RETURN_IF_ERROR(_flush_segment_writer(writer, flush_size));
     } else {
+        // the horizontal writer has no streaming feed, build it all up front
+        RETURN_IF_ERROR_OR_CATCH_EXCEPTION(
+                segment_v2::materialize_derived_columns(derived_column, 
&flush_block));
         std::unique_ptr<segment_v2::SegmentWriter> writer;
         RETURN_IF_ERROR(_create_segment_writer(writer, segment_id, 
no_compression));
         RETURN_IF_ERROR_OR_CATCH_EXCEPTION(_add_rows(writer, &flush_block, 0, 
flush_block.rows()));
@@ -88,6 +113,15 @@ Status SegmentFlusher::flush_single_block(const Block* 
block, int32_t segment_id
     return Status::OK();
 }
 
+Status SegmentFlusher::transform_block(Block* block, int32_t segment_id,
+                                       segment_v2::DerivedColumn* 
derived_column) {
+    auto transform_ctx = make_transform_exec_context(_context, segment_id);
+    RETURN_IF_ERROR_OR_CATCH_EXCEPTION(
+            segment_v2::build_transform_chain(_context).apply(transform_ctx, 
block));
+    *derived_column = std::move(transform_ctx.derived_column);
+    return Status::OK();
+}
+
 Status SegmentFlusher::close() {
     RETURN_IF_ERROR(_seg_files.close());
     RETURN_IF_ERROR(_preload_segment_indexes_to_file_cache());
@@ -406,6 +440,15 @@ Status SegmentCreator::add_block(const Block* block) {
     size_t block_row_num = block->rows();
     size_t row_avg_size_in_bytes = std::max((size_t)1, block_size_in_bytes / 
block_row_num);
     size_t row_offset = 0;
+    // This seam always feeds the horizontal writer, so the derived column is
+    // materialized up front, like flush_single_block's horizontal branch.
+    Block* shared_block = const_cast<Block*>(block);
+    auto transform_block = [&]() -> Status {
+        segment_v2::DerivedColumn derived_column;
+        RETURN_IF_ERROR(
+                _segment_flusher.transform_block(shared_block, 
/*segment_id=*/-1, &derived_column));
+        return segment_v2::materialize_derived_columns(derived_column, 
shared_block);
+    };
 
     if (_flush_writer == nullptr) {
         RETURN_IF_ERROR(_segment_flusher.create_writer(_flush_writer, 
allocate_segment_id()));
@@ -421,6 +464,7 @@ Status SegmentCreator::add_block(const Block* block) {
             DCHECK(max_row_add > 0);
         }
         size_t input_row_num = std::min(block_row_num - row_offset, 
size_t(max_row_add));
+        RETURN_IF_ERROR(transform_block());
         RETURN_IF_ERROR(_flush_writer->add_rows(block, row_offset, 
input_row_num));
         row_offset += input_row_num;
     } while (row_offset < block_row_num);
diff --git a/be/src/storage/rowset/segment_creator.h 
b/be/src/storage/rowset/segment_creator.h
index 2f161ccf276..63312dd47f4 100644
--- a/be/src/storage/rowset/segment_creator.h
+++ b/be/src/storage/rowset/segment_creator.h
@@ -20,7 +20,9 @@
 #include <gen_cpp/internal_service.pb.h>
 #include <gen_cpp/olap_file.pb.h>
 
+#include <memory>
 #include <mutex>
+#include <utility>
 #include <vector>
 
 #include "common/status.h"
@@ -37,6 +39,10 @@ class Block;
 namespace segment_v2 {
 class SegmentWriter;
 class VerticalSegmentWriter;
+class DerivedColumnGenerator;
+// Matches block_transform.h: at most one derived column (the row-store column)
+// for each flush, held as a {cid, generator} pair; null generator means none.
+using DerivedColumn = std::pair<uint32_t, std::shared_ptr<const 
DerivedColumnGenerator>>;
 } // namespace segment_v2
 
 struct SegmentStatistics;
@@ -101,6 +107,11 @@ public:
 
     ~SegmentFlusher();
 
+    // Runs the block transform chain on `block` and hands back the derived 
(row-store)
+    // column for the caller to feed into its writer.
+    Status transform_block(Block* block, int32_t segment_id,
+                           segment_v2::DerivedColumn* derived_column);
+
     // Return the file size flushed to disk in "flush_size"
     // This method is thread-safe.
     Status flush_single_block(const Block* block, int32_t segment_id,
diff --git a/be/src/storage/segment/segment_writer.cpp 
b/be/src/storage/segment/segment_writer.cpp
index ecb45c7ed57..e7b61896f16 100644
--- a/be/src/storage/segment/segment_writer.cpp
+++ b/be/src/storage/segment/segment_writer.cpp
@@ -634,19 +634,8 @@ Status SegmentWriter::append_block(const Block* block, 
size_t row_pos, size_t nu
             << ", block->columns()=" << block->columns()
             << ", _column_writers.size()=" << _column_writers.size()
             << ", _tablet_schema->dump_structure()=" << 
_tablet_schema->dump_structure();
-    // Row column should be filled here when it's a directly write from 
memtable
-    // or it's schema change write(since column data type maybe changed, so we 
should reubild)
-    if (_opts.write_type == DataWriteType::TYPE_DIRECT ||
-        _opts.write_type == DataWriteType::TYPE_SCHEMA_CHANGE) {
-        _serialize_block_to_row_column(*const_cast<Block*>(block));
-    }
-
-    if (_opts.rowset_ctx->write_type != DataWriteType::TYPE_COMPACTION &&
-        _tablet_schema->num_variant_columns() > 0) {
-        RETURN_IF_ERROR(variant_util::parse_and_materialize_variant_columns(
-                const_cast<Block&>(*block), *_tablet_schema, _column_ids));
-    }
-
+    // Blocks from the seams arrive already transformed (variants parsed, 
row-store
+    // column materialized); compaction-family callers bring rows that are 
already final.
     _olap_data_convertor->set_source_content(block, row_pos, num_rows);
 
     // convert column data from engine format to storage layer format
diff --git a/be/src/storage/segment/vertical_segment_writer.cpp 
b/be/src/storage/segment/vertical_segment_writer.cpp
index 3bc95833de9..633069bd9f0 100644
--- a/be/src/storage/segment/vertical_segment_writer.cpp
+++ b/be/src/storage/segment/vertical_segment_writer.cpp
@@ -78,6 +78,7 @@
 #include "storage/segment/variant/variant_ext_meta_writer.h"
 #include "storage/tablet/base_tablet.h"
 #include "storage/tablet/tablet_schema.h"
+#include "storage/transform/block_transform.h"
 #include "storage/utils.h"
 #include "util/coding.h"
 #include "util/debug_points.h"
@@ -360,6 +361,39 @@ Status 
VerticalSegmentWriter::_append_row_store_column(const Block& block, size_
     return Status::OK();
 }
 
+Status VerticalSegmentWriter::_append_generated_column(const 
DerivedColumnGenerator& generator,
+                                                       const Block& block, 
size_t row_pos,
+                                                       size_t num_rows, 
uint32_t cid) {
+    if (num_rows == 0) {
+        return Status::OK();
+    }
+    DCHECK_LE(row_pos + num_rows, block.rows());
+
+    size_t end_pos = row_pos + num_rows;
+    size_t batch_rows = _opts.num_rows_per_block;
+    static constexpr size_t kDerivedColumnBatchBytes = 4 * 1024 * 1024;
+    DCHECK_GT(batch_rows, 0);
+    for (size_t pos = row_pos; pos < end_pos;) {
+        size_t max_rows = std::min(batch_rows, end_pos - pos);
+        auto generated_column = 
block.get_by_position(cid).column->clone_empty();
+        size_t rows = generator.generate(block, pos, max_rows, 
kDerivedColumnBatchBytes,
+                                         generated_column.get());
+        DCHECK_GT(rows, 0);
+
+        auto typed_column = block.get_by_position(cid);
+        typed_column.column = std::move(generated_column);
+        
RETURN_IF_ERROR(_olap_data_convertor->set_source_content_with_specifid_column(
+                typed_column, 0, rows, cid));
+        auto [status, column] = _olap_data_convertor->convert_column_data(cid);
+        RETURN_IF_ERROR(status);
+        RETURN_IF_ERROR(
+                _column_writers[cid]->append(column->get_nullmap(), 
column->get_data(), rows));
+        _olap_data_convertor->clear_source_content(cid);
+        pos += rows;
+    }
+    return Status::OK();
+}
+
 Status VerticalSegmentWriter::_probe_key_for_mow(
         const MowKeyProbe& probe, std::string key, std::size_t segment_pos,
         bool have_input_seq_column, bool have_delete_sign,
@@ -966,36 +1000,17 @@ Status VerticalSegmentWriter::write_batch() {
         }
         return Status::OK();
     }
-    // Row column should be filled here when it's a directly write from 
memtable
-    // or it's schema change write(since column data type maybe changed, so we 
should reubild)
-    bool should_write_row_store_column = _opts.write_type == 
DataWriteType::TYPE_DIRECT ||
-                                         _opts.write_type == 
DataWriteType::TYPE_SCHEMA_CHANGE;
-    if (should_write_row_store_column) {
-        for (uint32_t cid = 0; cid < _tablet_schema->num_columns(); ++cid) {
-            if (!_tablet_schema->column(cid).is_row_store_column()) {
-                continue;
-            }
-            RETURN_IF_ERROR(
-                    _create_column_writer(cid, _tablet_schema->column(cid), 
_tablet_schema));
-            for (auto& data : _batched_blocks) {
-                RETURN_IF_ERROR(
-                        _append_row_store_column(*data.block, data.row_pos, 
data.num_rows, cid));
-            }
-            RETURN_IF_ERROR(_check_column_writer_disk_capacity(cid));
-            RETURN_IF_ERROR(_finalize_column_writer_and_update_meta(cid));
-        }
-    }
-
-    std::vector<uint32_t> column_ids;
-    for (uint32_t i = 0; i < _tablet_schema->num_columns(); ++i) {
-        column_ids.emplace_back(i);
-    }
-    if (_opts.rowset_ctx->write_type != DataWriteType::TYPE_COMPACTION &&
-        _tablet_schema->num_variant_columns() > 0) {
+    // The transform chain already validated, parsed variants and decided the 
derived
+    // (row-store) column; this writer only pumps the generator in bounded 
batches.
+    if (_derived_column.second) {
+        const auto& [cid, generator] = _derived_column;
+        RETURN_IF_ERROR(_create_column_writer(cid, 
_tablet_schema->column(cid), _tablet_schema));
         for (auto& data : _batched_blocks) {
-            
RETURN_IF_ERROR(variant_util::parse_and_materialize_variant_columns(
-                    const_cast<Block&>(*data.block), *_tablet_schema, 
column_ids));
+            RETURN_IF_ERROR(_append_generated_column(*generator, *data.block, 
data.row_pos,
+                                                     data.num_rows, cid));
         }
+        RETURN_IF_ERROR(_check_column_writer_disk_capacity(cid));
+        RETURN_IF_ERROR(_finalize_column_writer_and_update_meta(cid));
     }
 
     std::vector<IOlapColumnDataAccessor*> key_columns;
@@ -1003,7 +1018,7 @@ Status VerticalSegmentWriter::write_batch() {
     // the key is cluster key column unique id
     std::map<uint32_t, IOlapColumnDataAccessor*> cid_to_column;
     for (uint32_t cid = 0; cid < _tablet_schema->num_columns(); ++cid) {
-        if (should_write_row_store_column && 
_tablet_schema->column(cid).is_row_store_column()) {
+        if (_derived_column.second && _derived_column.first == cid) {
             continue;
         }
         RETURN_IF_ERROR(_create_column_writer(cid, 
_tablet_schema->column(cid), _tablet_schema));
@@ -1045,6 +1060,8 @@ Status VerticalSegmentWriter::write_batch() {
     }
 
     _batched_blocks.clear();
+    // The generator snapshots the batched blocks' rows; it must not survive 
them.
+    _derived_column = {};
     return Status::OK();
 }
 
diff --git a/be/src/storage/segment/vertical_segment_writer.h 
b/be/src/storage/segment/vertical_segment_writer.h
index 25fd4ef0c40..a68c3ef0569 100644
--- a/be/src/storage/segment/vertical_segment_writer.h
+++ b/be/src/storage/segment/vertical_segment_writer.h
@@ -71,6 +71,11 @@ struct VerticalSegmentWriterOptions {
     std::shared_ptr<MowContext> mow_ctx;
 };
 
+class DerivedColumnGenerator;
+// Matches block_transform.h: at most one derived column (the row-store column)
+// for each flush, held as a {cid, generator} pair; null generator means none.
+using DerivedColumn = std::pair<uint32_t, std::shared_ptr<const 
DerivedColumnGenerator>>;
+
 struct RowsInBlock {
     const Block* block;
     size_t row_pos;
@@ -96,6 +101,10 @@ public:
     Status batch_block(const Block* block, size_t row_pos, size_t num_rows);
     Status write_batch();
 
+    void set_derived_column(DerivedColumn derived_column) {
+        _derived_column = std::move(derived_column);
+    }
+
     [[nodiscard]] std::string data_dir_path() const {
         return _data_dir == nullptr ? "" : _data_dir->path();
     }
@@ -153,6 +162,8 @@ private:
     void _set_max_key(const Slice& key);
     Status _append_row_store_column(const Block& block, size_t row_pos, size_t 
num_rows,
                                     uint32_t cid);
+    Status _append_generated_column(const DerivedColumnGenerator& generator, 
const Block& block,
+                                    size_t row_pos, size_t num_rows, uint32_t 
cid);
     // Thin wrapper over MowKeyProbe that translates a ProbeOutcome back into 
the out-parameters the
     // partial update fill loops use. `found_cb` receives the rowset that 
holds `loc`: the fixed
     // path pins it in its HistoricalRowFetcher, the flexible path in 
`_rsid_to_rowset`, which its
@@ -255,6 +266,9 @@ private:
 
     std::vector<RowsInBlock> _batched_blocks;
 
+    // the derived column the transform chain hands off to this writer's 
bounded pump
+    DerivedColumn _derived_column;
+
     BlockAggregator _block_aggregator;
 };
 
diff --git a/be/src/storage/transform/block_transform.cpp 
b/be/src/storage/transform/block_transform.cpp
new file mode 100644
index 00000000000..93fc2be4515
--- /dev/null
+++ b/be/src/storage/transform/block_transform.cpp
@@ -0,0 +1,233 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#include "storage/transform/block_transform.h"
+
+#include <limits>
+#include <numeric>
+#include <unordered_set>
+
+#include "common/cast_set.h"
+#include "common/logging.h"
+#include "core/block/block.h"
+#include "core/column/column_string.h"
+#include "exec/common/variant_util.h"
+#include "storage/partial_update_info.h"
+#include "storage/rowset/rowset_writer_context.h"
+#include "storage/tablet/tablet_schema.h"
+#include "util/jsonb/serialize.h"
+
+namespace doris::segment_v2 {
+
+namespace {
+
+// Parses raw variant columns into subcolumn form, in place. The heavy variant
+// work stays in the variant ColumnWriter; this stage only reshapes the block
+// before conversion.
+class VariantParseStage : public BlockTransform {
+public:
+    Status apply(TransformExecContext& ctx, Block* block) const override {
+        const auto& schema = *ctx.tablet_schema;
+        if (schema.num_variant_columns() == 0) {
+            return Status::OK();
+        }
+        std::vector<uint32_t> column_ids(block->columns());
+        std::iota(column_ids.begin(), column_ids.end(), 0);
+        return variant_util::parse_and_materialize_variant_columns(*block, 
schema, column_ids);
+    }
+    std::string_view name() const override { return "VariantParse"; }
+};
+
+// Checks schema rules and block width for every block entering a seam. The
+// writers keep transitional duplicates of these checks until later changes
+// remove them: non-seam callers (compaction, index change) still rely on them.
+class ValidateStage : public BlockTransform {
+public:
+    Status apply(TransformExecContext& ctx, Block* block) const override {
+        const TabletSchema& schema = *ctx.tablet_schema;
+        if (schema.cluster_key_uids().empty()) {
+            DCHECK(schema.num_key_columns() >= schema.num_short_key_columns())
+                    << ", table_id=" << schema.table_id()
+                    << ", num_key_columns=" << schema.num_key_columns()
+                    << ", num_short_key_columns=" << 
schema.num_short_key_columns();
+        }
+        const auto* info = ctx.partial_update_info.get();
+        const bool is_partial_update_load = info != nullptr && 
info->is_partial_update() &&
+                                            ctx.write_type == 
DataWriteType::TYPE_DIRECT &&
+                                            
!ctx.rowset_ctx->is_transient_rowset_writer;
+        if (!is_partial_update_load) {
+            if (block->columns() != schema.num_columns()) {
+                return Status::InvalidArgument(
+                        "illegal block columns, block columns = {}, 
tablet_schema columns = {}",
+                        block->dump_structure(), schema.dump_structure());
+            }
+            return Status::OK();
+        }
+
+        // No tablet context (e.g. the streaming BetaRowsetWriterV2) means this
+        // path can't do partial update: return a clear error instead of
+        // crashing in the probe.
+        if (ctx.tablet == nullptr || ctx.mow_context == nullptr) {
+            return Status::NotSupported(
+                    "partial update is not supported on this write path (no 
tablet context)");
+        }
+        if (!(schema.keys_type() == UNIQUE_KEYS &&
+              ctx.rowset_ctx->enable_unique_key_merge_on_write)) {
+            auto msg = fmt::format(
+                    "Can only do partial update on merge-on-write unique 
table, but found: "
+                    "keys_type={}, enable_unique_key_merge_on_write={}, 
tablet_id={}",
+                    schema.keys_type(), 
ctx.rowset_ctx->enable_unique_key_merge_on_write,
+                    ctx.tablet->tablet_id());
+            DCHECK(false) << msg;
+            return Status::InternalError<false>(msg);
+        }
+        // partial update needs the segment id, which only flush_single_block 
sets
+        if (ctx.segment_id < 0) {
+            return Status::InternalError<false>(
+                    "partial update blocks must be flushed through 
flush_single_block, "
+                    "tablet_id={}",
+                    ctx.tablet->tablet_id());
+        }
+        if (info->is_flexible_partial_update()) {
+            if (block->columns() != schema.num_columns()) {
+                return Status::InvalidArgument(
+                        "illegal flexible partial update block columns, block 
columns = {}, "
+                        "tablet_schema columns = {}",
+                        block->dump_structure(), schema.dump_structure());
+            }
+        } else {
+            DCHECK(info->is_fixed_partial_update());
+            if (block->columns() < schema.num_key_columns() ||
+                block->columns() >= schema.num_columns()) {
+                return Status::InvalidArgument(fmt::format(
+                        "illegal partial update block columns: {}, num key 
columns: {}, total "
+                        "schema columns: {}",
+                        block->columns(), schema.num_key_columns(), 
schema.num_columns()));
+            }
+        }
+        return Status::OK();
+    }
+    std::string_view name() const override { return "Validate"; }
+};
+
+// Generates the hidden row-store column (each row as JSONB). A
+// DerivedColumnGenerator so the vertical writer can stream it in batches.
+class RowStoreColumnGenerator : public DerivedColumnGenerator {
+public:
+    RowStoreColumnGenerator(TabletSchemaSPtr schema, Block source_block)
+            : _schema(std::move(schema)),
+              _source_block(std::move(source_block)),
+              _serdes(create_data_type_serdes(_source_block.get_data_types())),
+              _row_store_cids(_schema->row_columns_uids().begin(),
+                              _schema->row_columns_uids().end()) {}
+
+    size_t generate(const Block& block, size_t row_pos, size_t max_rows, 
size_t max_bytes,
+                    IColumn* dst) const override {
+        // Rows are read from the COW snapshot, indexed by the caller's 
positions in `block`;
+        // the two must describe the same rows or the serialized cells 
silently mismatch.
+        DCHECK_EQ(_source_block.rows(), block.rows());
+        auto* dst_str = static_cast<ColumnString*>(dst);
+        return JsonbSerializeUtil::block_to_jsonb(*_schema, _source_block, 
*dst_str,
+                                                  
cast_set<int>(_schema->num_columns()), _serdes,
+                                                  _row_store_cids, row_pos, 
max_rows, max_bytes);
+    }
+
+private:
+    TabletSchemaSPtr _schema;
+    Block _source_block;
+    DataTypeSerDeSPtrs _serdes;
+    std::unordered_set<int32_t> _row_store_cids;
+};
+
+// Registers a row-store generator over a COW snapshot of the block at this
+// exact stage. Variant parsing can change its JSONB representation, so the
+// snapshot preserves the legacy writer's RowStore/Variant ordering while the
+// vertical writer still materializes the column in bounded batches.
+class RowStoreFillStage : public BlockTransform {
+public:
+    Status apply(TransformExecContext& ctx, Block* block) const override {
+        if (block->rows() == 0) {
+            return Status::OK();
+        }
+        const auto& schema = *ctx.tablet_schema;
+        for (size_t i = 0; i < schema.num_columns(); ++i) {
+            if (!schema.column(i).is_row_store_column()) {
+                continue;
+            }
+            std::shared_ptr<const DerivedColumnGenerator> generator =
+                    
std::make_shared<RowStoreColumnGenerator>(ctx.tablet_schema, *block);
+            ctx.derived_column = std::make_pair(cast_set<uint32_t>(i), 
std::move(generator));
+            break;
+        }
+        return Status::OK();
+    }
+    std::string_view name() const override { return "RowStoreFill"; }
+};
+
+} // namespace
+
+BlockTransformChain build_transform_chain(const RowsetWriterContext& context) {
+    if (context.write_type == DataWriteType::TYPE_COMPACTION) {
+        return BlockTransformChain {};
+    }
+    if (context.write_binlog_opt().enable) {
+        // RowBinlogSegmentWriter still derives the binlog rows itself, so its
+        // chain stays empty until that derivation moves in here.
+        return BlockTransformChain {};
+    }
+    std::vector<std::shared_ptr<const BlockTransform>> stages;
+    stages.push_back(std::make_shared<ValidateStage>());
+    const bool is_partial_update_load = context.partial_update_info != nullptr 
&&
+                                        
context.partial_update_info->is_partial_update() &&
+                                        context.write_type == 
DataWriteType::TYPE_DIRECT &&
+                                        !context.is_transient_rowset_writer;
+    if (is_partial_update_load) {
+        // Partial update loads only get validated here for now: the segment
+        // writers still do their own fill, parse and row-store work until the
+        // fill stages move into the chain.
+        return BlockTransformChain {std::move(stages)};
+    }
+    const bool rebuild_row_store = context.write_type == 
DataWriteType::TYPE_DIRECT ||
+                                   context.write_type == 
DataWriteType::TYPE_SCHEMA_CHANGE;
+    // Direct and schema-change writers rebuilt RowStore from the raw Variant
+    // representation, then parsed Variant for its column writer.
+    if (rebuild_row_store) {
+        stages.push_back(std::make_shared<RowStoreFillStage>());
+    }
+    stages.push_back(std::make_shared<VariantParseStage>());
+    return BlockTransformChain {std::move(stages)};
+}
+
+Status materialize_derived_columns(const DerivedColumn& derived_column, Block* 
block) {
+    if (!derived_column.second) {
+        return Status::OK();
+    }
+    const auto& [cid, generator] = derived_column;
+    auto column_ptr = block->get_by_position(cid).column->clone_empty();
+    size_t num_rows = block->rows();
+    size_t pos = 0;
+    while (pos < num_rows) {
+        size_t rows = generator->generate(*block, pos, num_rows - pos,
+                                          std::numeric_limits<size_t>::max(), 
column_ptr.get());
+        DCHECK_GT(rows, 0);
+        pos += rows;
+    }
+    block->replace_by_position(cid, std::move(column_ptr));
+    return Status::OK();
+}
+
+} // namespace doris::segment_v2
diff --git a/be/src/storage/transform/block_transform.h 
b/be/src/storage/transform/block_transform.h
new file mode 100644
index 00000000000..755459277cb
--- /dev/null
+++ b/be/src/storage/transform/block_transform.h
@@ -0,0 +1,129 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+#pragma once
+
+#include <map>
+#include <memory>
+#include <string_view>
+#include <utility>
+#include <vector>
+
+#include "common/status.h"
+#include "storage/olap_common.h"
+#include "storage/olap_define.h"
+#include "storage/partial_update_info.h"
+#include "storage/tablet/tablet_fwd.h"
+
+namespace doris {
+class Block;
+class IColumn;
+struct RowsetWriterContext;
+
+namespace segment_v2 {
+
+// Computes a derived column (e.g. the hidden row-store column) in bounded 
batches.
+class DerivedColumnGenerator {
+public:
+    virtual ~DerivedColumnGenerator() = default;
+    // Appends a bounded batch (<= max_rows, stops near max_bytes, but always 
>= 1
+    // row); returns the rows produced. A generator may snapshot its source 
rows at
+    // registration time, so it must only be driven with the block it was 
built over:
+    // `block` and `row_pos` index the same rows the snapshot holds.
+    virtual size_t generate(const Block& block, size_t row_pos, size_t 
max_rows, size_t max_bytes,
+                            IColumn* dst) const = 0;
+};
+
+using DerivedColumn = std::pair<uint32_t, std::shared_ptr<const 
DerivedColumnGenerator>>;
+
+// For the horizontal writer, which can't feed derived columns in bounded 
batches.
+Status materialize_derived_columns(const DerivedColumn& derived_column, Block* 
block);
+
+// State for each flush, kept out of the chain so the chain stays immutable and
+// shareable across concurrent flushes. Filled in for each flush; the tablet
+// schema in particular differs between flushes for variant tables.
+struct TransformExecContext {
+    TabletSchemaSPtr tablet_schema;
+    DataWriteType write_type = DataWriteType::TYPE_DEFAULT;
+
+    // --- partial update inputs (filled by the call sites from 
RowsetWriterContext) ---
+    BaseTabletSPtr tablet;
+    std::shared_ptr<MowContext> mow_context;
+    std::shared_ptr<PartialUpdateInfo> partial_update_info;
+    RowsetWriterContext* rowset_ctx = nullptr;
+    // identifies the segment this block lands in (for the self-marks)
+    RowsetId rowset_id;
+    // -1 until flush_single_block sets it; partial update needs it set.
+    int32_t segment_id = -1;
+
+    // --- outputs ---
+    // the derived column for the writer to generate in batches; null 
generator = none
+    DerivedColumn derived_column;
+};
+
+// One Block -> Block step run before the segment writers. apply() may mutate
+// the block in place or swap its columns, but must not keep references to it.
+class BlockTransform {
+public:
+    virtual ~BlockTransform() = default;
+    virtual Status apply(TransformExecContext& ctx, Block* block) const = 0;
+    // Stable name used for debugging and to check how the chain is built.
+    virtual std::string_view name() const = 0;
+};
+
+// An ordered list of transforms. Building one is cheap, so the seams build it
+// per flush; it is immutable and shareable across concurrent flushes.
+class BlockTransformChain {
+public:
+    BlockTransformChain() = default;
+    explicit BlockTransformChain(std::vector<std::shared_ptr<const 
BlockTransform>> stages)
+            : _stages(std::move(stages)) {}
+
+    Status apply(TransformExecContext& ctx, Block* block) const {
+        for (const auto& stage : _stages) {
+            RETURN_IF_ERROR(stage->apply(ctx, block));
+        }
+        return Status::OK();
+    }
+
+    bool empty() const { return _stages.empty(); }
+
+    std::vector<std::string_view> stage_names() const {
+        std::vector<std::string_view> names;
+        names.reserve(_stages.size());
+        for (const auto& stage : _stages) {
+            names.push_back(stage->name());
+        }
+        return names;
+    }
+
+private:
+    std::vector<std::shared_ptr<const BlockTransform>> _stages;
+};
+
+// The single place that decides which transforms a write path gets:
+//   - compaction: empty (rows are already final)
+//   - binlog sub-writer: empty for now (RowBinlogSegmentWriter still derives
+//     the binlog rows itself; a later change moves that in here)
+//   - partial update: [Validate] for now (the segment writers still do their
+//     own fill, parse and row-store work; later changes move the fill in here)
+//   - direct / schema change / transient flush: [Validate, RowStoreFill, 
VariantParse]
+// RowStoreFill is omitted when the write type does not rebuild the row-store 
column.
+BlockTransformChain build_transform_chain(const RowsetWriterContext& context);
+
+} // namespace segment_v2
+} // namespace doris
diff --git a/be/test/storage/mow/mow_transform_test_base.h 
b/be/test/storage/mow/mow_transform_test_base.h
index a7e3fff6630..4b49552d9e4 100644
--- a/be/test/storage/mow/mow_transform_test_base.h
+++ b/be/test/storage/mow/mow_transform_test_base.h
@@ -139,6 +139,103 @@ protected:
         return schema;
     }
 
+    // (k INT key, v VARIANT, delete-sign) UNIQUE_KEYS MoW schema, for the 
variant parse stage.
+    TabletSchemaSPtr create_variant_schema() {
+        TabletSchemaPB pb;
+        pb.set_keys_type(UNIQUE_KEYS);
+        pb.set_num_short_key_columns(1);
+        pb.set_num_rows_per_row_block(1024);
+        pb.set_compress_kind(COMPRESS_LZ4);
+        pb.set_next_column_unique_id(10);
+        {
+            ColumnPB* c = pb.add_column();
+            c->set_unique_id(0);
+            c->set_name("k");
+            c->set_type("INT");
+            c->set_is_key(true);
+            c->set_length(4);
+            c->set_index_length(4);
+            c->set_is_nullable(false);
+            c->set_aggregation("NONE");
+        }
+        {
+            ColumnPB* c = pb.add_column();
+            c->set_unique_id(1);
+            c->set_name("v");
+            c->set_type("VARIANT");
+            c->set_is_key(false);
+            c->set_length(2147483643);
+            c->set_index_length(4);
+            c->set_is_nullable(false);
+            c->set_aggregation("NONE");
+            c->set_variant_max_subcolumns_count(3);
+        }
+        {
+            ColumnPB* c = pb.add_column();
+            c->set_unique_id(2);
+            c->set_name(DELETE_SIGN);
+            c->set_type("TINYINT");
+            c->set_is_key(false);
+            c->set_length(1);
+            c->set_index_length(1);
+            c->set_is_nullable(false);
+            c->set_aggregation("NONE");
+            c->set_default_value(std::to_string(0));
+        }
+        pb.set_delete_sign_idx(2);
+        auto schema = std::make_shared<TabletSchema>();
+        schema->init_from_pb(pb);
+        return schema;
+    }
+
+    // (k INT key, v INT, delete-sign, __DORIS_SKIP_BITMAP_COL__) flexible 
partial update MoW
+    // schema: flexible loads carry a full-width block plus the per-row skip 
bitmap.
+    TabletSchemaSPtr create_flexible_mow_schema() {
+        TabletSchemaPB pb;
+        pb.set_keys_type(UNIQUE_KEYS);
+        pb.set_num_short_key_columns(1);
+        pb.set_num_rows_per_row_block(1024);
+        pb.set_compress_kind(COMPRESS_LZ4);
+        pb.set_next_column_unique_id(10);
+
+        auto type_length = [](const std::string& type) -> int32_t {
+            if (type == "TINYINT") {
+                return 1;
+            }
+            if (type == "BITMAP") {
+                return 16;
+            }
+            return 4;
+        };
+        auto add_col = [&](int uid, const std::string& name, const 
std::string& type, bool is_key,
+                           bool nullable, const std::string& def = "") {
+            ColumnPB* c = pb.add_column();
+            c->set_unique_id(uid);
+            c->set_name(name);
+            c->set_type(type);
+            c->set_is_key(is_key);
+            c->set_length(type_length(type));
+            c->set_index_length(type_length(type));
+            c->set_is_nullable(nullable);
+            c->set_aggregation("NONE");
+            if (!def.empty()) {
+                c->set_default_value(def);
+            }
+        };
+        add_col(0, "k", "INT", true, false);
+        add_col(1, "v", "INT", false, true, std::to_string(0));
+        add_col(2, DELETE_SIGN, "TINYINT", false, false, std::to_string(0));
+        add_col(3, SKIP_BITMAP_COL, "BITMAP", false, false);
+        // init_from_pb reads these hidden-column indices straight from the PB
+        // fields (it does not scan by name), so they must be set explicitly.
+        pb.set_delete_sign_idx(2);
+        pb.set_skip_bitmap_col_idx(3);
+
+        auto schema = std::make_shared<TabletSchema>();
+        schema->init_from_pb(pb);
+        return schema;
+    }
+
     // (k INT key, v INT, [seq INT], delete-sign, __DORIS_ROW_STORE_COL__ 
STRING) with the hidden
     // full row-store column enabled -- the only shape for which the write 
path touches the row
     // cache.
@@ -187,6 +284,16 @@ protected:
         return schema;
     }
 
+    // A MoW tablet over the engine's data dir; no rowsets are registered with 
it.
+    TabletSharedPtr make_tablet(const TabletSchemaSPtr& schema, int64_t 
tablet_id) {
+        TabletMetaSharedPtr tablet_meta = std::make_shared<TabletMeta>();
+        tablet_meta->_tablet_id = tablet_id;
+        static_cast<void>(tablet_meta->set_partition_id(10));
+        tablet_meta->_schema = schema;
+        tablet_meta->_enable_unique_key_merge_on_write = true;
+        return std::make_shared<Tablet>(*_engine, tablet_meta, 
_data_dir.get(), "mow_ut");
+    }
+
     void make_rowset_ctx(const TabletSchemaSPtr& schema, int64_t 
rowset_numeric_id, int64_t version,
                          RowsetWriterContext* ctx, TabletSharedPtr* 
out_tablet) {
         RowsetId rid;
@@ -203,12 +310,7 @@ protected:
         ctx->enable_unique_key_merge_on_write = true;
         ctx->write_type = DataWriteType::TYPE_DIRECT;
 
-        TabletMetaSharedPtr tablet_meta = std::make_shared<TabletMeta>();
-        tablet_meta->_tablet_id = kTabletId;
-        static_cast<void>(tablet_meta->set_partition_id(10));
-        tablet_meta->_schema = schema;
-        tablet_meta->_enable_unique_key_merge_on_write = true;
-        auto tablet = std::make_shared<Tablet>(*_engine, tablet_meta, 
_data_dir.get(), "mow_ut");
+        auto tablet = make_tablet(schema, kTabletId);
         ctx->tablet = tablet;
         *out_tablet = tablet;
     }
@@ -323,6 +425,15 @@ protected:
         return assert_cast<const ColumnInt32&>(*col).get_data()[row];
     }
 
+    // Reads an int8 cell from a (possibly nullable) TINYINT column of `block`.
+    static int8_t read_tinyint(const Block& block, size_t col_pos, size_t row) 
{
+        const IColumn* col = block.get_by_position(col_pos).column.get();
+        if (col->is_nullable()) {
+            col = &assert_cast<const 
ColumnNullable&>(*col).get_nested_column();
+        }
+        return assert_cast<const ColumnInt8&>(*col).get_data()[row];
+    }
+
     // True iff the cell is SQL NULL. A non-nullable column is never null. Use 
alongside read_int to
     // distinguish "value X" from "null with X left in the nested column".
     static bool read_is_null(const Block& block, size_t col_pos, size_t row) {
diff --git a/be/test/storage/transform/validate_stage_test.cpp 
b/be/test/storage/transform/validate_stage_test.cpp
new file mode 100644
index 00000000000..793614eedec
--- /dev/null
+++ b/be/test/storage/transform/validate_stage_test.cpp
@@ -0,0 +1,405 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+// build_transform_chain composition per write path, and every ValidateStage
+// branch that is reachable through the built chain.
+
+#include <gtest/gtest.h>
+
+#include <memory>
+#include <string_view>
+#include <vector>
+
+#include "common/status.h"
+#include "storage/mow/mow_transform_test_base.h"
+#include "storage/partial_update_info.h"
+#include "storage/rowset/beta_rowset_writer.h"
+#include "storage/rowset/rowset_writer_context.h"
+#include "storage/rowset/segment_creator.h"
+#include "storage/transform/block_transform.h"
+
+namespace doris {
+
+using segment_v2::build_transform_chain;
+using segment_v2::TransformExecContext;
+
+// Own fixture subclass so the TEST_F names here never clash with the other
+// files sharing MowTransformTestBase.
+class ValidateStageTest : public MowTransformTestBase {
+protected:
+    // A minimal direct-write rowset context: enough for build_transform_chain 
to
+    // pick the non-binlog, non-compaction stage set.
+    RowsetWriterContext direct_rwc(const TabletSchemaSPtr& schema) {
+        RowsetWriterContext c;
+        c.tablet_schema = schema;
+        c.write_type = DataWriteType::TYPE_DIRECT;
+        c.enable_unique_key_merge_on_write = true;
+        return c;
+    }
+    // The per-flush exec context the chain runs against. ValidateStage reads
+    // write_type / segment_id straight from here and 
is_transient_rowset_writer
+    // through rowset_ctx, so the two must agree with the rwc used to build.
+    TransformExecContext exec_ctx(const TabletSchemaSPtr& schema, 
RowsetWriterContext* rwc,
+                                  int32_t segment_id = 0) {
+        TransformExecContext ctx;
+        ctx.tablet_schema = schema;
+        ctx.write_type = rwc->write_type;
+        ctx.rowset_ctx = rwc;
+        ctx.segment_id = segment_id;
+        return ctx;
+    }
+};
+
+// 
=============================================================================
+// chain composition per write path -- assert the exact ordered stage list.
+// 
=============================================================================
+
+// TYPE_COMPACTION -> empty chain (rows are already final).
+TEST_F(ValidateStageTest, CompositionCompactionEmpty) {
+    auto schema = create_mow_schema(/*has_seq=*/false);
+    RowsetWriterContext c = direct_rwc(schema);
+    c.write_type = DataWriteType::TYPE_COMPACTION;
+    EXPECT_TRUE(build_transform_chain(c).empty());
+    EXPECT_TRUE(build_transform_chain(c).stage_names().empty());
+}
+
+// TYPE_DIRECT, no PU -> Validate, RowStoreFill, VariantParse.
+TEST_F(ValidateStageTest, CompositionDirectNonPartialUpdate) {
+    using V = std::vector<std::string_view>;
+    auto schema = create_mow_schema(/*has_seq=*/false);
+    EXPECT_EQ(build_transform_chain(direct_rwc(schema)).stage_names(),
+              (V {"Validate", "RowStoreFill", "VariantParse"}));
+}
+
+// TYPE_SCHEMA_CHANGE, no PU -> same shape as a direct write.
+TEST_F(ValidateStageTest, CompositionSchemaChange) {
+    using V = std::vector<std::string_view>;
+    auto schema = create_mow_schema(/*has_seq=*/false);
+    RowsetWriterContext c = direct_rwc(schema);
+    c.write_type = DataWriteType::TYPE_SCHEMA_CHANGE;
+    EXPECT_EQ(build_transform_chain(c).stage_names(),
+              (V {"Validate", "RowStoreFill", "VariantParse"}));
+}
+
+// Non-binlog TYPE_DEFAULT does not rebuild the row store: Validate and parse 
only.
+TEST_F(ValidateStageTest, CompositionDefaultOmitsRowStoreFill) {
+    using V = std::vector<std::string_view>;
+    auto schema = create_mow_schema(/*has_seq=*/false);
+    RowsetWriterContext c = direct_rwc(schema);
+    c.write_type = DataWriteType::TYPE_DEFAULT;
+    EXPECT_EQ(build_transform_chain(c).stage_names(), (V {"Validate", 
"VariantParse"}));
+}
+
+// A transient-rowset-writer PU degrades to the plain direct chain: the PU
+// predicate is false because is_transient_rowset_writer is set, so the block
+// is treated as a full-width direct write.
+TEST_F(ValidateStageTest, CompositionTransientPartialUpdateDegradesToNoFill) {
+    using V = std::vector<std::string_view>;
+    auto schema = create_mow_schema(/*has_seq=*/false);
+    auto pui = std::make_shared<PartialUpdateInfo>();
+    ASSERT_TRUE(pui->init(kTabletId, 1, *schema, 
UniqueKeyUpdateModePB::UPDATE_FIXED_COLUMNS,
+                          PartialUpdateNewRowPolicyPB::APPEND, {"k"}, false, 
0, 0, "UTC", "")
+                        .ok());
+    RowsetWriterContext c = direct_rwc(schema);
+    c.partial_update_info = pui;
+    c.is_transient_rowset_writer = true; // degrade: PU predicate becomes false
+    EXPECT_EQ(build_transform_chain(c).stage_names(),
+              (V {"Validate", "RowStoreFill", "VariantParse"}));
+}
+
+// Partial update loads only get validated by the chain for now: the segment
+// writers still own the fill, parse and row-store work. The fill stages take
+// this slot when they move into the chain.
+TEST_F(ValidateStageTest, CompositionPartialUpdateValidateOnly) {
+    using V = std::vector<std::string_view>;
+    auto schema = create_mow_schema(/*has_seq=*/false);
+    auto fixed = std::make_shared<PartialUpdateInfo>();
+    ASSERT_TRUE(fixed->init(kTabletId, 1, *schema, 
UniqueKeyUpdateModePB::UPDATE_FIXED_COLUMNS,
+                            PartialUpdateNewRowPolicyPB::APPEND, {"k"}, false, 
0, 0, "UTC", "")
+                        .ok());
+    RowsetWriterContext c = direct_rwc(schema);
+    c.partial_update_info = fixed;
+    EXPECT_EQ(build_transform_chain(c).stage_names(), (V {"Validate"}));
+
+    auto fschema = create_flexible_mow_schema();
+    auto flexible = std::make_shared<PartialUpdateInfo>();
+    ASSERT_TRUE(flexible->init(kTabletId, 1, *fschema,
+                               UniqueKeyUpdateModePB::UPDATE_FLEXIBLE_COLUMNS,
+                               PartialUpdateNewRowPolicyPB::APPEND, {}, false, 
0, 0, "UTC", "")
+                        .ok());
+    RowsetWriterContext fc = direct_rwc(fschema);
+    fc.partial_update_info = flexible;
+    EXPECT_EQ(build_transform_chain(fc).stage_names(), (V {"Validate"}));
+}
+
+// Binlog sub-writers keep deriving inside RowBinlogSegmentWriter for now, so
+// their chain stays empty for every write type. The derive stage takes this
+// slot when it moves into the chain.
+TEST_F(ValidateStageTest, CompositionBinlogEmpty) {
+    auto schema = create_mow_schema(/*has_seq=*/false);
+    RowsetWriterContext rwc = direct_rwc(schema);
+    rwc.write_binlog_opt().enable = true;
+
+    for (auto write_type : {DataWriteType::TYPE_DIRECT, 
DataWriteType::TYPE_DEFAULT,
+                            DataWriteType::TYPE_SCHEMA_CHANGE}) {
+        rwc.write_type = write_type;
+        EXPECT_TRUE(build_transform_chain(rwc).empty());
+    }
+}
+
+// 
=============================================================================
+// ValidateStage branches (V1-V9, without V5's fill which is not in the chain
+// yet). ValidateStage is the chain's first stage on every non-compaction,
+// non-binlog path, so we build the real chain and drive it with a block that
+// already fails / passes validate.
+// 
=============================================================================
+
+// V1: non-PU direct, full width (columns == num_columns) -> accepted.
+TEST_F(ValidateStageTest, V1_DirectAcceptsGoodWidth) {
+    auto schema = create_mow_schema(/*has_seq=*/false);
+    RowsetWriterContext c = direct_rwc(schema);
+    auto chain = build_transform_chain(c);
+    TransformExecContext ctx = exec_ctx(schema, &c);
+
+    Block block = schema->create_block(); // full width, 0 rows
+    EXPECT_TRUE(chain.apply(ctx, &block).ok());
+}
+
+// V2: non-PU direct, wrong width (columns != num_columns) -> InvalidArgument.
+TEST_F(ValidateStageTest, V2_DirectRejectsBadWidth) {
+    auto schema = create_mow_schema(/*has_seq=*/false); // 3 columns
+    RowsetWriterContext c = direct_rwc(schema);
+    auto chain = build_transform_chain(c);
+    TransformExecContext ctx = exec_ctx(schema, &c);
+
+    Block block = schema->create_block_by_cids({0}); // 1 column != 
num_columns(3)
+    auto st = chain.apply(ctx, &block);
+    EXPECT_FALSE(st.ok());
+    EXPECT_EQ(st.code(), ErrorCode::INVALID_ARGUMENT) << st;
+    EXPECT_NE(st.to_string().find("illegal block columns"), std::string::npos) 
<< st;
+}
+
+// V3: PU on a write path with no tablet context -> NotSupported instead of a
+// crash inside the probe (e.g. the streaming BetaRowsetWriterV2).
+TEST_F(ValidateStageTest, V3_PartialUpdateRejectsNoTabletContext) {
+    auto schema = create_mow_schema(/*has_seq=*/false);
+    auto pui = std::make_shared<PartialUpdateInfo>();
+    ASSERT_TRUE(pui->init(kTabletId, 1, *schema, 
UniqueKeyUpdateModePB::UPDATE_FIXED_COLUMNS,
+                          PartialUpdateNewRowPolicyPB::APPEND, {"k"}, false, 
0, 0, "UTC", "")
+                        .ok());
+    RowsetWriterContext rwc = direct_rwc(schema);
+    rwc.partial_update_info = pui;
+    auto chain = build_transform_chain(rwc);
+
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+    ctx.partial_update_info = pui;
+    ctx.tablet = nullptr; // no tablet context
+    ctx.mow_context = nullptr;
+
+    Block block = schema->create_block_by_cids({0});
+    block.get_by_position(0).column->assert_mutable()->insert_default();
+    auto st = chain.apply(ctx, &block);
+    EXPECT_FALSE(st.ok());
+    EXPECT_EQ(st.code(), ErrorCode::NOT_IMPLEMENTED_ERROR) << st;
+    EXPECT_NE(st.to_string().find("no tablet context"), std::string::npos) << 
st;
+}
+
+// V4: PU flushed without a segment id (segment_id == -1, the add_block seam) 
->
+// InternalError naming flush_single_block.
+TEST_F(ValidateStageTest, V4_PartialUpdateRejectsWithoutSegmentId) {
+    auto schema = create_mow_schema(/*has_seq=*/false);
+    TabletSharedPtr tablet;
+    auto rowset = write_rowset(schema, 8201, 2, {{1, 11}}, &tablet);
+    auto mow = make_mow_context(100, {rowset});
+
+    auto pui = std::make_shared<PartialUpdateInfo>();
+    ASSERT_TRUE(pui->init(kTabletId, 1, *schema, 
UniqueKeyUpdateModePB::UPDATE_FIXED_COLUMNS,
+                          PartialUpdateNewRowPolicyPB::APPEND, {"k"}, false, 
0, 0, "UTC", "")
+                        .ok());
+
+    RowsetWriterContext rwc = direct_rwc(schema);
+    rwc.tablet_id = kTabletId;
+    rwc.tablet = tablet;
+    rwc.partial_update_info = pui;
+
+    auto chain = build_transform_chain(rwc);
+    TransformExecContext ctx;
+    ctx.tablet_schema = schema;
+    ctx.write_type = DataWriteType::TYPE_DIRECT;
+    ctx.tablet = tablet;
+    ctx.mow_context = mow;
+    ctx.partial_update_info = pui;
+    ctx.rowset_ctx = &rwc;
+    ctx.segment_id = -1; // add_block seam: no segment id
+
+    Block block = schema->create_block_by_cids({0});
+    IColumn* kc = block.get_by_position(0).column->assert_mutable().get();
+    int32_t k = 1;
+    kc->insert_data(reinterpret_cast<const char*>(&k), sizeof(int32_t));
+
+    auto st = chain.apply(ctx, &block);
+    EXPECT_FALSE(st.ok());
+    EXPECT_EQ(st.code(), ErrorCode::INTERNAL_ERROR) << st;
+    EXPECT_NE(st.to_string().find("flush_single_block"), std::string::npos) << 
st;
+}
+
+// V6 + V7: fixed PU rejects both a too-wide block (columns >= num_columns) 
and a
+// too-narrow one (columns < num_key_columns) with the same InvalidArgument.
+TEST_F(ValidateStageTest, V6V7_FixedPartialUpdateRejectsBadWidth) {
+    auto schema = create_mow_schema(/*has_seq=*/false); // 1 key, 3 cols
+    TabletSharedPtr tablet;
+    auto rowset = write_rowset(schema, 8401, 2, {{1, 11}}, &tablet);
+    auto mow = make_mow_context(100, {rowset});
+    auto pui = std::make_shared<PartialUpdateInfo>();
+    ASSERT_TRUE(pui->init(kTabletId, 1, *schema, 
UniqueKeyUpdateModePB::UPDATE_FIXED_COLUMNS,
+                          PartialUpdateNewRowPolicyPB::APPEND, {"k"}, false, 
0, 0, "UTC", "")
+                        .ok());
+    RowsetId new_rsid;
+    new_rsid.init(8402);
+    RowsetWriterContext rwc = direct_rwc(schema);
+    rwc.tablet_id = kTabletId;
+    rwc.tablet = tablet;
+    rwc.partial_update_info = pui;
+    rwc.rowset_id = new_rsid;
+    auto chain = build_transform_chain(rwc);
+
+    auto make_ctx = [&] {
+        TransformExecContext ctx;
+        ctx.tablet_schema = schema;
+        ctx.write_type = DataWriteType::TYPE_DIRECT;
+        ctx.tablet = tablet;
+        ctx.mow_context = mow;
+        ctx.partial_update_info = pui;
+        ctx.rowset_ctx = &rwc;
+        ctx.rowset_id = new_rsid;
+        ctx.segment_id = 0;
+        return ctx;
+    };
+
+    // V6 too wide: full width (3 == num_columns) is not a partial update 
block.
+    {
+        TransformExecContext ctx = make_ctx();
+        Block block = schema->create_block();
+        auto st = chain.apply(ctx, &block);
+        EXPECT_FALSE(st.ok());
+        EXPECT_EQ(st.code(), ErrorCode::INVALID_ARGUMENT) << st;
+        EXPECT_NE(st.to_string().find("illegal partial update block columns"), 
std::string::npos)
+                << st;
+    }
+    // V7 too narrow: fewer columns than the key (0 < 1 key column).
+    {
+        TransformExecContext ctx = make_ctx();
+        Block block = schema->create_block_by_cids({});
+        auto st = chain.apply(ctx, &block);
+        EXPECT_FALSE(st.ok());
+        EXPECT_EQ(st.code(), ErrorCode::INVALID_ARGUMENT) << st;
+        EXPECT_NE(st.to_string().find("illegal partial update block columns"), 
std::string::npos)
+                << st;
+    }
+}
+
+// V8: flexible PU requires a full-width block; any width mismatch is rejected.
+TEST_F(ValidateStageTest, V8_FlexiblePartialUpdateRejectsBadWidth) {
+    auto schema = create_flexible_mow_schema(); // k v delete_sign 
skip_bitmap: 4 cols
+    auto tablet = make_tablet(schema, 8501);
+    auto mow = make_mow_context(100, {});
+    auto pui = std::make_shared<PartialUpdateInfo>();
+    ASSERT_TRUE(pui->init(kTabletId, 1, *schema, 
UniqueKeyUpdateModePB::UPDATE_FLEXIBLE_COLUMNS,
+                          PartialUpdateNewRowPolicyPB::APPEND, {}, false, 0, 
0, "UTC", "")
+                        .ok());
+    RowsetWriterContext rwc = direct_rwc(schema);
+    rwc.tablet = tablet;
+    rwc.partial_update_info = pui;
+    auto chain = build_transform_chain(rwc);
+
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+    ctx.tablet = tablet;
+    ctx.mow_context = mow;
+    ctx.partial_update_info = pui;
+
+    Block block = schema->create_block_by_cids({0}); // 1 col != num_columns(4)
+    auto st = chain.apply(ctx, &block);
+    EXPECT_FALSE(st.ok());
+    EXPECT_EQ(st.code(), ErrorCode::INVALID_ARGUMENT) << st;
+    EXPECT_NE(st.to_string().find("illegal flexible partial update block 
columns"),
+              std::string::npos)
+            << st;
+}
+
+// V9: a transient PU is validated as a plain direct write (the PU predicate is
+// false because is_transient_rowset_writer is set). The full-width block is
+// accepted; a narrow one is rejected with the non-PU "illegal block columns"
+// error, proving the transient path does NOT use the partial-update 
validation.
+TEST_F(ValidateStageTest, V9_TransientPartialUpdateValidatedAsDirect) {
+    auto schema = create_mow_schema(/*has_seq=*/false); // 3 columns
+    auto pui = std::make_shared<PartialUpdateInfo>();
+    ASSERT_TRUE(pui->init(kTabletId, 1, *schema, 
UniqueKeyUpdateModePB::UPDATE_FIXED_COLUMNS,
+                          PartialUpdateNewRowPolicyPB::APPEND, {"k"}, false, 
0, 0, "UTC", "")
+                        .ok());
+    RowsetWriterContext rwc = direct_rwc(schema);
+    rwc.partial_update_info = pui;
+    rwc.is_transient_rowset_writer = true; // degrade to direct
+    auto chain = build_transform_chain(rwc);
+
+    // full-width block is accepted as a direct write
+    {
+        TransformExecContext ctx = exec_ctx(schema, &rwc);
+        ctx.partial_update_info = pui;
+        Block block = schema->create_block(); // 3 cols == num_columns
+        EXPECT_TRUE(chain.apply(ctx, &block).ok());
+    }
+    // a narrow block is rejected with the non-PU width error -- not the PU one
+    {
+        TransformExecContext ctx = exec_ctx(schema, &rwc);
+        ctx.partial_update_info = pui;
+        Block block = schema->create_block_by_cids({0}); // 1 col != 
num_columns(3)
+        auto st = chain.apply(ctx, &block);
+        EXPECT_FALSE(st.ok());
+        EXPECT_EQ(st.code(), ErrorCode::INVALID_ARGUMENT) << st;
+        EXPECT_NE(st.to_string().find("illegal block columns"), 
std::string::npos) << st;
+    }
+}
+
+// The flush seam runs the chain before it creates any segment writer: a block
+// ValidateStage rejects must surface the chain's error with no segment file
+// touched. The context deliberately has no file_writer_creator, so reaching
+// writer creation at all would fail very differently from InvalidArgument.
+// Today the writers still duplicate this width check; once later changes 
remove
+// those duplicates, this seam is the only guard, so pin it now.
+TEST_F(ValidateStageTest, FlushSeamRejectsBeforeCreatingAWriter) {
+    auto schema = create_mow_schema(/*has_seq=*/false); // 3 columns
+    RowsetWriterContext rwc = direct_rwc(schema);
+
+    SegmentFileCollection segment_files;
+    InvertedIndexFileCollection index_files;
+    SegmentFlusher flusher(rwc, segment_files, index_files);
+
+    Block block = schema->create_block_by_cids({0}); // 1 column != 
num_columns(3)
+    block.get_by_position(0).column->assert_mutable()->insert_default();
+    auto st = flusher.flush_single_block(&block, /*segment_id=*/0);
+    EXPECT_FALSE(st.ok());
+    EXPECT_EQ(st.code(), ErrorCode::INVALID_ARGUMENT) << st;
+    EXPECT_NE(st.to_string().find("illegal block columns"), std::string::npos) 
<< st;
+}
+
+// The sort-key invariant (num_key_columns >= num_short_key_columns) is the 
only
+// branch a plain unit test cannot exercise: pure DCHECK, aborts a debug build
+// and compiles out in release. The "Can only do partial update on 
merge-on-write
+// unique table" branch does return a deterministic InternalError in release,
+// but its DCHECK still aborts debug test runs.
+
+} // namespace doris
diff --git a/be/test/storage/transform/variant_rowstore_test.cpp 
b/be/test/storage/transform/variant_rowstore_test.cpp
new file mode 100644
index 00000000000..4e4c5bd7b70
--- /dev/null
+++ b/be/test/storage/transform/variant_rowstore_test.cpp
@@ -0,0 +1,547 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+// VariantParseStage, RowStoreFillStage and the RowStoreColumnGenerator pump
+// contract (bounded batches by rows and by bytes, always >= 1 row).
+
+#include <gtest/gtest.h>
+
+#include <limits>
+#include <string>
+#include <string_view>
+#include <unordered_map>
+#include <unordered_set>
+#include <vector>
+
+#include "core/block/block.h"
+#include "core/column/column_string.h"
+#include "core/column/column_variant.h"
+#include "core/data_type_serde/data_type_serde.h"
+#include "core/field.h"
+#include "storage/mow/mow_transform_test_base.h"
+#include "storage/rowset/rowset_writer_context.h"
+#include "storage/transform/block_transform.h"
+#include "testutil/variant_util.h"
+#include "util/jsonb/serialize.h"
+#include "util/jsonb_document.h"
+
+namespace doris {
+
+using segment_v2::build_transform_chain;
+using segment_v2::materialize_derived_columns;
+using segment_v2::TransformExecContext;
+
+class VariantRowStoreTest : public MowTransformTestBase {
+protected:
+    RowsetWriterContext direct_rwc(const TabletSchemaSPtr& schema) {
+        RowsetWriterContext c;
+        c.tablet_schema = schema;
+        c.write_type = DataWriteType::TYPE_DIRECT;
+        c.enable_unique_key_merge_on_write = true;
+        return c;
+    }
+    TransformExecContext exec_ctx(const TabletSchemaSPtr& schema, 
RowsetWriterContext* rwc,
+                                  int32_t segment_id = 0) {
+        TransformExecContext ctx;
+        ctx.tablet_schema = schema;
+        ctx.write_type = rwc->write_type;
+        ctx.rowset_ctx = rwc;
+        ctx.segment_id = segment_id;
+        return ctx;
+    }
+
+    // create_variant_schema() plus the hidden row-store column, so one table
+    // carries both a variant and the whole-row store.
+    TabletSchemaSPtr create_variant_row_store_schema() {
+        TabletSchemaPB pb;
+        create_variant_schema()->to_schema_pb(&pb);
+        pb.set_store_row_column(true);
+        pb.set_next_column_unique_id(11);
+        ColumnPB* row_store = pb.add_column();
+        row_store->set_unique_id(10);
+        row_store->set_name(BeConsts::ROW_STORE_COL);
+        row_store->set_type("STRING");
+        row_store->set_is_key(false);
+        row_store->set_length(2147483643);
+        row_store->set_index_length(4);
+        row_store->set_is_nullable(false);
+        row_store->set_aggregation("NONE");
+
+        auto schema = std::make_shared<TabletSchema>();
+        schema->init_from_pb(pb);
+        return schema;
+    }
+
+    // Inserts one root-scalar JSON object string into a block's variant 
column.
+    static void insert_variant_json(Block& block, size_t variant_pos, 
std::string_view json) {
+        auto* variant = assert_cast<ColumnVariant*>(
+                
block.get_by_position(variant_pos).column->assert_mutable().get());
+        VariantUtil::insert_root_scalar_field(
+                *variant, 
Field::create_field<TYPE_STRING>(String(std::string(json))));
+    }
+
+    // Round-trips one finalized variant row back to canonical JSON (spaces 
stripped).
+    static std::string variant_row_json(const Block& block, size_t 
variant_pos, size_t row) {
+        const auto* parsed =
+                assert_cast<const 
ColumnVariant*>(block.get_by_position(variant_pos).column.get());
+        DataTypeSerDe::FormatOptions options;
+        std::string json;
+        parsed->serialize_one_row_to_string(static_cast<int64_t>(row), &json, 
options);
+        std::erase(json, ' ');
+        return json;
+    }
+
+    // Decodes one row-store JSONB cell back into a 1-row block of the 
non-row-store
+    // columns, the way BaseTablet::fetch_value_through_row_column does on 
read.
+    // Returns the decoded block; `block` is keyed by the schema's logical 
position.
+    Block decode_row_store_cell(const TabletSchemaSPtr& schema, StringRef 
cell) {
+        // Build a block + serdes for every non-row-store column, keyed by 
unique_id.
+        std::vector<uint32_t> cids;
+        for (size_t i = 0; i < schema->num_columns(); ++i) {
+            if (!schema->column(i).is_row_store_column()) {
+                cids.push_back(static_cast<uint32_t>(i));
+            }
+        }
+        Block dst = schema->create_block_by_cids(cids);
+        DataTypeSerDeSPtrs serdes = 
create_data_type_serdes(dst.get_data_types());
+        std::unordered_map<uint32_t, uint32_t> col_uid_to_idx;
+        std::vector<std::string> default_values(cids.size());
+        for (size_t i = 0; i < cids.size(); ++i) {
+            const TabletColumn& col = schema->column(cids[i]);
+            col_uid_to_idx[static_cast<uint32_t>(col.unique_id())] = 
static_cast<uint32_t>(i);
+            default_values[i] = col.default_value();
+        }
+        EXPECT_TRUE(JsonbSerializeUtil::jsonb_to_block(serdes, cell.data, 
cell.size, col_uid_to_idx,
+                                                       dst, default_values, {})
+                            .ok());
+        return dst;
+    }
+};
+
+// ===========================================================================
+// VariantParseStage
+// ===========================================================================
+
+// A schema with no variant column -> VariantParseStage is a pass-through;
+// column count, row count and every cell value are left untouched.
+TEST_F(VariantRowStoreTest, VariantParseNoVariantPassThrough) {
+    auto schema = create_mow_schema(/*has_seq=*/false); // k v delete_sign: no 
variant
+    ASSERT_EQ(schema->num_variant_columns(), 0U);
+    RowsetWriterContext rwc = direct_rwc(schema);
+    auto chain = build_transform_chain(rwc);
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+
+    Block block = schema->create_block(); // full width, 2 rows
+    IColumn* k = block.get_by_position(0).column->assert_mutable().get();
+    IColumn* v = block.get_by_position(1).column->assert_mutable().get();
+    IColumn* ds = block.get_by_position(2).column->assert_mutable().get();
+    int32_t ks[] = {7, 9};
+    int32_t vs[] = {70, 90};
+    int8_t zero8 = 0;
+    for (int i = 0; i < 2; ++i) {
+        k->insert_data(reinterpret_cast<const char*>(&ks[i]), sizeof(int32_t));
+        v->insert_data(reinterpret_cast<const char*>(&vs[i]), sizeof(int32_t));
+        ds->insert_data(reinterpret_cast<const char*>(&zero8), sizeof(int8_t));
+    }
+
+    ASSERT_TRUE(chain.apply(ctx, &block).ok());
+    // unchanged: same width, same height, same values
+    ASSERT_EQ(block.columns(), schema->num_columns());
+    ASSERT_EQ(block.rows(), 2);
+    EXPECT_EQ(read_int(block, 0, 0), 7);
+    EXPECT_EQ(read_int(block, 0, 1), 9);
+    EXPECT_EQ(read_int(block, 1, 0), 70);
+    EXPECT_EQ(read_int(block, 1, 1), 90);
+    EXPECT_EQ(read_tinyint(block, 2, 0), 0);
+    EXPECT_EQ(read_tinyint(block, 2, 1), 0);
+    // a non-variant table registers no derived column either
+    EXPECT_EQ(ctx.derived_column.second, nullptr);
+}
+
+// A direct write parses the root-only variant in place -- the row finalizes
+// and the original {"a":1,"b":"x"} survives parse + finalize as the same JSON.
+TEST_F(VariantRowStoreTest, VariantParseDirectSingleRow) {
+    auto schema = create_variant_schema(); // k(0) v VARIANT(1) delete_sign(2)
+    ASSERT_EQ(schema->num_variant_columns(), 1U);
+    RowsetWriterContext rwc = direct_rwc(schema);
+    auto chain = build_transform_chain(rwc);
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+
+    Block block = schema->create_block();
+    int32_t k = 1;
+    int8_t z = 0;
+    block.get_by_position(0).column->assert_mutable()->insert_data(
+            reinterpret_cast<const char*>(&k), sizeof(int32_t));
+    insert_variant_json(block, 1, R"({"a":1,"b":"x"})");
+    block.get_by_position(2).column->assert_mutable()->insert_data(
+            reinterpret_cast<const char*>(&z), sizeof(int8_t));
+
+    ASSERT_TRUE(chain.apply(ctx, &block).ok());
+    ASSERT_EQ(block.columns(), schema->num_columns());
+    ASSERT_EQ(block.rows(), 1);
+    const auto* parsed = assert_cast<const 
ColumnVariant*>(block.get_by_position(1).column.get());
+    EXPECT_TRUE(parsed->is_finalized());
+    const std::string json = variant_row_json(block, 1, 0);
+    EXPECT_NE(json.find(R"("a":1)"), std::string::npos) << json;
+    EXPECT_NE(json.find(R"("b":"x")"), std::string::npos) << json;
+}
+
+// Two distinct objects both finalize and each round-trips to its own inserted
+// keys; the column width is unchanged.
+TEST_F(VariantRowStoreTest, VariantParseDirectMultiRow) {
+    auto schema = create_variant_schema();
+    RowsetWriterContext rwc = direct_rwc(schema);
+    auto chain = build_transform_chain(rwc);
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+
+    Block block = schema->create_block();
+    IColumn* k = block.get_by_position(0).column->assert_mutable().get();
+    IColumn* ds = block.get_by_position(2).column->assert_mutable().get();
+    int32_t ks[] = {1, 2};
+    int8_t z = 0;
+    k->insert_data(reinterpret_cast<const char*>(&ks[0]), sizeof(int32_t));
+    insert_variant_json(block, 1, R"({"a":1})");
+    ds->insert_data(reinterpret_cast<const char*>(&z), sizeof(int8_t));
+    k->insert_data(reinterpret_cast<const char*>(&ks[1]), sizeof(int32_t));
+    insert_variant_json(block, 1, R"({"a":2,"c":true})");
+    ds->insert_data(reinterpret_cast<const char*>(&z), sizeof(int8_t));
+
+    ASSERT_TRUE(chain.apply(ctx, &block).ok());
+    ASSERT_EQ(block.columns(), schema->num_columns());
+    ASSERT_EQ(block.rows(), 2);
+    const auto* parsed = assert_cast<const 
ColumnVariant*>(block.get_by_position(1).column.get());
+    EXPECT_TRUE(parsed->is_finalized());
+    const std::string json0 = variant_row_json(block, 1, 0);
+    const std::string json1 = variant_row_json(block, 1, 1);
+    EXPECT_NE(json0.find(R"("a":1)"), std::string::npos) << json0;
+    EXPECT_NE(json1.find(R"("a":2)"), std::string::npos) << json1;
+    // the variant serializes a JSON bool as an integer (true -> 1)
+    EXPECT_NE(json1.find(R"("c":1)"), std::string::npos) << json1;
+    // row 0 did not gain row 1's key
+    EXPECT_EQ(json0.find(R"("c":)"), std::string::npos) << json0;
+}
+
+// An empty variant block parses without crashing and keeps its full width with
+// zero rows.
+TEST_F(VariantRowStoreTest, VariantParseEmptyBlock) {
+    auto schema = create_variant_schema();
+    RowsetWriterContext rwc = direct_rwc(schema);
+    auto chain = build_transform_chain(rwc);
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+
+    Block block = schema->create_block(); // typed columns, 0 rows
+    ASSERT_EQ(block.rows(), 0);
+    ASSERT_TRUE(chain.apply(ctx, &block).ok());
+    EXPECT_EQ(block.columns(), schema->num_columns());
+    EXPECT_EQ(block.rows(), 0);
+}
+
+// ===========================================================================
+// RowStoreFillStage + RowStoreColumnGenerator
+// ===========================================================================
+
+// Builds a full-width row-store block of `rows` rows (k = i+1, v = base+10*i),
+// row-store column left as a placeholder default for the generator to 
overwrite.
+static Block make_row_store_block(const TabletSchemaSPtr& schema, int 
num_rows, int32_t base) {
+    Block block = schema->create_block();
+    IColumn* k = block.get_by_position(0).column->assert_mutable().get();
+    IColumn* v = block.get_by_position(1).column->assert_mutable().get();
+    IColumn* ds = block.get_by_position(2).column->assert_mutable().get();
+    IColumn* rs = block.get_by_position(3).column->assert_mutable().get();
+    int8_t zero8 = 0;
+    for (int i = 0; i < num_rows; ++i) {
+        int32_t kk = i + 1;
+        int32_t vv = base + 10 * i;
+        k->insert_data(reinterpret_cast<const char*>(&kk), sizeof(int32_t));
+        v->insert_data(reinterpret_cast<const char*>(&vv), sizeof(int32_t));
+        ds->insert_data(reinterpret_cast<const char*>(&zero8), sizeof(int8_t));
+        rs->insert_default(); // placeholder, replaced by the generator
+    }
+    return block;
+}
+
+// An empty block returns OK and registers NO generator (early return before
+// the schema scan).
+TEST_F(VariantRowStoreTest, RowStoreFillRowsZeroNoGenerator) {
+    auto schema = create_row_store_schema();
+    RowsetWriterContext rwc = direct_rwc(schema);
+    auto chain = build_transform_chain(rwc);
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+
+    Block block = schema->create_block(); // 0 rows
+    ASSERT_EQ(block.rows(), 0);
+    ASSERT_TRUE(chain.apply(ctx, &block).ok());
+    EXPECT_EQ(ctx.derived_column.second, nullptr);
+}
+
+// A non-empty row-store block registers a generator for the hidden row-store
+// column (cid 3), with a non-null generator.
+TEST_F(VariantRowStoreTest, RowStoreFillRegistersGenerator) {
+    auto schema = create_row_store_schema(); // k(0) v(1) delete_sign(2) 
row_store(3)
+    RowsetWriterContext rwc = direct_rwc(schema);
+    auto chain = build_transform_chain(rwc);
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+
+    Block block = make_row_store_block(schema, 2, /*base=*/10);
+    ASSERT_TRUE(chain.apply(ctx, &block).ok());
+    ASSERT_NE(ctx.derived_column.second, nullptr);
+    EXPECT_EQ(ctx.derived_column.first, 3U);
+}
+
+// Horizontal one-shot materialize fills every row with real, distinct,
+// non-empty JSONB.
+TEST_F(VariantRowStoreTest, RowStoreFillMaterializeHorizontal) {
+    auto schema = create_row_store_schema();
+    RowsetWriterContext rwc = direct_rwc(schema);
+    auto chain = build_transform_chain(rwc);
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+
+    Block block = make_row_store_block(schema, 2, /*base=*/10); // v = 10, 20
+    ASSERT_TRUE(chain.apply(ctx, &block).ok());
+    ASSERT_NE(ctx.derived_column.second, nullptr);
+    ASSERT_EQ(ctx.derived_column.first, 3U);
+
+    ASSERT_TRUE(materialize_derived_columns(ctx.derived_column, &block).ok());
+    ASSERT_EQ(block.rows(), 2);
+    const auto& rs_str = assert_cast<const 
ColumnString&>(*block.get_by_position(3).column);
+    ASSERT_EQ(rs_str.size(), 2U);
+    StringRef row0 = rs_str.get_data_at(0);
+    StringRef row1 = rs_str.get_data_at(1);
+    EXPECT_GT(row0.size, 0U);
+    EXPECT_GT(row1.size, 0U);
+    EXPECT_NE(row0.to_string(), row1.to_string()); // v differs (10 vs 20)
+}
+
+// Decode the materialized JSONB of one row and check the exact uid->value
+// mapping. The whole-row store encodes every non-row-store column keyed by
+// unique_id (uid0=k, uid1=v, uid2=delete_sign) and never the row-store column.
+TEST_F(VariantRowStoreTest, RowStoreFillMaterializeContent) {
+    auto schema = create_row_store_schema();
+    RowsetWriterContext rwc = direct_rwc(schema);
+    auto chain = build_transform_chain(rwc);
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+
+    // one row: k=42, v=-5, delete_sign=0
+    Block block = schema->create_block();
+    int32_t k = 42;
+    int32_t v = -5;
+    int8_t z = 0;
+    block.get_by_position(0).column->assert_mutable()->insert_data(
+            reinterpret_cast<const char*>(&k), sizeof(int32_t));
+    block.get_by_position(1).column->assert_mutable()->insert_data(
+            reinterpret_cast<const char*>(&v), sizeof(int32_t));
+    block.get_by_position(2).column->assert_mutable()->insert_data(
+            reinterpret_cast<const char*>(&z), sizeof(int8_t));
+    block.get_by_position(3).column->assert_mutable()->insert_default();
+
+    ASSERT_TRUE(chain.apply(ctx, &block).ok());
+    ASSERT_NE(ctx.derived_column.second, nullptr);
+    ASSERT_TRUE(materialize_derived_columns(ctx.derived_column, &block).ok());
+
+    const auto& rs_str = assert_cast<const 
ColumnString&>(*block.get_by_position(3).column);
+    ASSERT_EQ(rs_str.size(), 1U);
+    StringRef cell = rs_str.get_data_at(0);
+    ASSERT_GT(cell.size, 0U);
+
+    Block decoded = decode_row_store_cell(schema, cell); // positions: k, v, 
delete_sign
+    ASSERT_EQ(decoded.rows(), 1);
+    EXPECT_EQ(read_int(decoded, 0, 0), 42);    // uid 0 = k
+    EXPECT_EQ(read_int(decoded, 1, 0), -5);    // uid 1 = v
+    EXPECT_EQ(read_tinyint(decoded, 2, 0), 0); // uid 2 = delete_sign
+
+    // the JSONB object holds exactly the three non-row-store uids {0,1,2} and
+    // never the row-store column's own uid (3).
+    const JsonbDocument* doc = nullptr;
+    ASSERT_TRUE(JsonbDocument::checkAndCreateDocument(cell.data, cell.size, 
&doc).ok());
+    std::unordered_set<int> key_ids;
+    // JsonbDocument's object iterator is not a standard range; the explicit
+    // begin/end loop is intentional.
+    // NOLINTNEXTLINE(modernize-loop-convert)
+    for (auto it = (*doc)->begin(); it != (*doc)->end(); ++it) {
+        key_ids.insert(static_cast<int>(it->getKeyId()));
+    }
+    EXPECT_TRUE(key_ids.count(0)) << "missing uid 0 (k)";
+    EXPECT_TRUE(key_ids.count(1)) << "missing uid 1 (v)";
+    EXPECT_TRUE(key_ids.count(2)) << "missing uid 2 (delete_sign)";
+    EXPECT_FALSE(key_ids.count(3)) << "row-store column uid 3 must not be 
encoded";
+}
+
+// RowStore must preserve the raw Variant representation that existed before
+// VariantParse. Parsing normalizes a JSON boolean to an integer in the Variant
+// column, but the row-store JSONB must still contain the original boolean.
+TEST_F(VariantRowStoreTest, RowStoreSnapshotsVariantBeforeParse) {
+    auto schema = create_variant_row_store_schema();
+    RowsetWriterContext rwc = direct_rwc(schema);
+    auto chain = build_transform_chain(rwc);
+    EXPECT_EQ(chain.stage_names(),
+              (std::vector<std::string_view> {"Validate", "RowStoreFill", 
"VariantParse"}));
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+
+    Block block = schema->create_block();
+    int32_t key = 1;
+    int8_t delete_sign = 0;
+    block.get_by_position(0).column->assert_mutable()->insert_data(
+            reinterpret_cast<const char*>(&key), sizeof(key));
+    insert_variant_json(block, 1, R"({"flag":true})");
+    block.get_by_position(2).column->assert_mutable()->insert_data(
+            reinterpret_cast<const char*>(&delete_sign), sizeof(delete_sign));
+    block.get_by_position(3).column->assert_mutable()->insert_default();
+
+    ASSERT_TRUE(chain.apply(ctx, &block).ok());
+    EXPECT_NE(variant_row_json(block, 1, 0).find(R"("flag":1)"), 
std::string::npos);
+    ASSERT_NE(ctx.derived_column.second, nullptr);
+    ASSERT_TRUE(materialize_derived_columns(ctx.derived_column, &block).ok());
+
+    const auto& row_store = assert_cast<const 
ColumnString&>(*block.get_by_position(3).column);
+    ASSERT_EQ(row_store.size(), 1U);
+    Block decoded = decode_row_store_cell(schema, row_store.get_data_at(0));
+    const std::string stored_variant = variant_row_json(decoded, 1, 0);
+    EXPECT_NE(stored_variant.find(R"("flag":true)"), std::string::npos) << 
stored_variant;
+    EXPECT_EQ(stored_variant.find(R"("flag":1)"), std::string::npos) << 
stored_variant;
+}
+
+// Drive the registered generator directly as the vertical writer does -- a
+// fresh clone_empty() dst per batch, max_bytes huge, batch_rows = 2. Over 5
+// rows this yields 2,2,1 and walks pos 0->2->4->5, and the concatenation
+// matches the one-shot materialize.
+TEST_F(VariantRowStoreTest, RowStorePumpByRows) {
+    auto schema = create_row_store_schema();
+    RowsetWriterContext rwc = direct_rwc(schema);
+    auto chain = build_transform_chain(rwc);
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+
+    Block block = make_row_store_block(schema, 5, /*base=*/100);
+    ASSERT_TRUE(chain.apply(ctx, &block).ok());
+    ASSERT_NE(ctx.derived_column.second, nullptr);
+    const auto& gen = *ctx.derived_column.second;
+    const uint32_t cid = ctx.derived_column.first;
+
+    // oracle: one-shot materialize on a copy
+    Block oracle = make_row_store_block(schema, 5, /*base=*/100);
+    ASSERT_TRUE(materialize_derived_columns(ctx.derived_column, &oracle).ok());
+    const auto& oracle_str = assert_cast<const 
ColumnString&>(*oracle.get_by_position(cid).column);
+
+    const size_t num_rows = block.rows();
+    const size_t batch_rows = 2;
+    const size_t big_bytes = std::numeric_limits<size_t>::max();
+    std::vector<size_t> batch_sizes;
+    std::vector<std::string> produced;
+    size_t pos = 0;
+    while (pos < num_rows) {
+        auto dst = block.get_by_position(cid).column->clone_empty();
+        size_t max_rows = std::min(batch_rows, num_rows - pos);
+        size_t rows = gen.generate(block, pos, max_rows, big_bytes, dst.get());
+        ASSERT_GT(rows, 0U);
+        batch_sizes.push_back(rows);
+        const auto& dst_str = assert_cast<const ColumnString&>(*dst);
+        ASSERT_EQ(dst_str.size(), rows);
+        for (size_t r = 0; r < rows; ++r) {
+            produced.push_back(dst_str.get_data_at(r).to_string());
+        }
+        pos += rows;
+    }
+    EXPECT_EQ(pos, num_rows);
+    EXPECT_EQ(batch_sizes, (std::vector<size_t> {2, 2, 1}));
+    ASSERT_EQ(produced.size(), num_rows);
+    for (size_t r = 0; r < num_rows; ++r) {
+        EXPECT_EQ(produced[r], oracle_str.get_data_at(r).to_string()) << "row 
" << r;
+    }
+}
+
+// Batch by bytes: batch_rows unbounded, max_bytes set just above one real 
row's
+// byte size (measured at runtime, never hardcoded). Each batch uses a fresh 
dst,
+// so block_to_jsonb stops once that batch's accumulated byte_size >= 
max_bytes:
+// it writes a row, then breaks -> one row per batch -> 5 batches summing to 5.
+TEST_F(VariantRowStoreTest, RowStorePumpByBytes) {
+    auto schema = create_row_store_schema();
+    RowsetWriterContext rwc = direct_rwc(schema);
+    auto chain = build_transform_chain(rwc);
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+
+    Block block = make_row_store_block(schema, 5, /*base=*/100);
+    ASSERT_TRUE(chain.apply(ctx, &block).ok());
+    ASSERT_NE(ctx.derived_column.second, nullptr);
+    const auto& gen = *ctx.derived_column.second;
+    const uint32_t cid = ctx.derived_column.first;
+
+    // measure a real single-row size first
+    size_t single_row_bytes = 0;
+    {
+        auto probe = block.get_by_position(cid).column->clone_empty();
+        size_t rows = gen.generate(block, 0, 1, 
std::numeric_limits<size_t>::max(), probe.get());
+        ASSERT_EQ(rows, 1U);
+        single_row_bytes = assert_cast<const 
ColumnString&>(*probe).byte_size();
+        ASSERT_GT(single_row_bytes, 0U);
+    }
+    // threshold at exactly one row -> block_to_jsonb writes one row, sees
+    // byte_size() >= max_bytes, and breaks, so each batch yields one row.
+    const size_t max_bytes = single_row_bytes;
+
+    const size_t num_rows = block.rows();
+    std::vector<size_t> batch_sizes;
+    size_t pos = 0;
+    size_t total = 0;
+    while (pos < num_rows) {
+        auto dst = block.get_by_position(cid).column->clone_empty();
+        // max_rows must be finite (<= remaining): the generator passes it
+        // straight to block_to_jsonb as num_rows; the byte cap forces the 
early
+        // break within the batch. (The real vertical writer caps it at
+        // num_rows_per_block.)
+        size_t rows = gen.generate(block, pos, num_rows - pos, max_bytes, 
dst.get());
+        ASSERT_GT(rows, 0U);
+        const auto& dst_str = assert_cast<const ColumnString&>(*dst);
+        ASSERT_EQ(dst_str.size(), rows);
+        batch_sizes.push_back(rows);
+        total += rows;
+        pos += rows;
+    }
+    EXPECT_EQ(total, num_rows);
+    // byte threshold ~ one row -> at least 3 batches (here exactly 5, one per 
row)
+    EXPECT_GE(batch_sizes.size(), 3U);
+    for (size_t s : batch_sizes) {
+        EXPECT_EQ(s, 1U);
+    }
+}
+
+// A single oversize row with max_bytes = 1 still produces exactly 1 row --
+// block_to_jsonb writes the row before testing the byte budget, so it never
+// returns 0.
+TEST_F(VariantRowStoreTest, RowStorePumpSingleOversize) {
+    auto schema = create_row_store_schema();
+    RowsetWriterContext rwc = direct_rwc(schema);
+    auto chain = build_transform_chain(rwc);
+    TransformExecContext ctx = exec_ctx(schema, &rwc);
+
+    Block block = make_row_store_block(schema, 1, /*base=*/100);
+    ASSERT_TRUE(chain.apply(ctx, &block).ok());
+    ASSERT_NE(ctx.derived_column.second, nullptr);
+    const auto& gen = *ctx.derived_column.second;
+    const uint32_t cid = ctx.derived_column.first;
+
+    // max_rows is the batch cap (= remaining rows; the generator forwards it 
to
+    // block_to_jsonb as num_rows, so it must be <= block.rows()). max_bytes=1 
is
+    // smaller than the single row, but block_to_jsonb writes one row before it
+    // checks the byte cap, so it still yields exactly one row (the >=1 
guarantee).
+    auto dst = block.get_by_position(cid).column->clone_empty();
+    size_t rows = gen.generate(block, 0, /*max_rows=*/1, /*max_bytes=*/1, 
dst.get());
+    EXPECT_EQ(rows, 1U);
+    const auto& dst_str = assert_cast<const ColumnString&>(*dst);
+    ASSERT_EQ(dst_str.size(), 1U);
+    EXPECT_GT(dst_str.get_data_at(0).size, 0U);
+}
+
+} // namespace doris


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to