This is an automated email from the ASF dual-hosted git repository.

luwei16 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 efd64dba17d [fix](binlog) Preserve commit TSO across rowset metadata 
rewrites (#68793)
efd64dba17d is described below

commit efd64dba17dc3db307d5a46eb1a24126bb2fc590
Author: Luwei <[email protected]>
AuthorDate: Fri Oct 9 15:19:15 2026 +0800

    [fix](binlog) Preserve commit TSO across rowset metadata rewrites (#68793)
    
    ### What problem does this PR solve?
    
    Issue Number: N/A
    
    Related PR: N/A
    
    Problem Summary: Single-version rowsets store a zero placeholder in
    their physical commit-TSO column and use rowset metadata to provide the
    actual timestamp. Local index build/drop and snapshot rowset-ID
    conversion rebuild metadata without preserving commit_tso. After those
    operations a fresh reader can return zero instead of the committed
    timestamp and incorrectly prune rows. Cloud snapshot metadata conversion
    has the same field omission.
    
    Copy the source commit_tso when present in these three one-to-one
    rewrite paths, reusing the existing manual-build metadata helper for
    index changes. Preserve absent and unassigned fields and complete TSO
    ranges without changing reader compatibility or generic multi-rowset
    writer behavior. Add tests for build and drop, fresh-reader projection
    and pruning, persisted local visible/stale rowsets, and Cloud metadata
    conversion. This prevents future omissions but does not recover
    already-lost timestamps.
    
    ### Release note
    
    Preserve rowset commit timestamps during local index changes and
    snapshot metadata conversion to avoid incorrect timestamp values and
    filtering.
    
    ### Check List (For Author)
    
    - Test: Unit Test
    - ASAN_UT: 8 focused tests reproduced 6 failures before the fix and all
    passed afterward; all 54 related BE tests passed.
    - clang-format v16, check-format, changed-line clang-tidy, build
    hygiene, and git diff --check.
    - No deployment or SQL regression: tests exercise the real rowset
    conversions, header persistence, and fresh segment readers; live Cloud
    RESTORE is not validated.
    - Behavior changed: Yes (preserve source commit TSO when rebuilding
    rowset metadata)
    - Does this need documentation: No
---
 be/src/cloud/cloud_snapshot_mgr.cpp                |   3 +
 be/src/storage/snapshot/snapshot_manager.cpp       |   3 +
 be/src/storage/task/index_builder.cpp              |   3 +
 be/test/cloud/cloud_snapshot_mgr_test.cpp          |  60 ++++++++++
 be/test/storage/index/index_builder_test.cpp       | 121 ++++++++++++++++++++-
 be/test/storage/snapshot/snapshot_manager_test.cpp |  64 +++++++++++
 6 files changed, 253 insertions(+), 1 deletion(-)

diff --git a/be/src/cloud/cloud_snapshot_mgr.cpp 
b/be/src/cloud/cloud_snapshot_mgr.cpp
index 9b121f42da3..375572bbdc7 100644
--- a/be/src/cloud/cloud_snapshot_mgr.cpp
+++ b/be/src/cloud/cloud_snapshot_mgr.cpp
@@ -272,6 +272,9 @@ Status CloudSnapshotMgr::_create_rowset_meta(
     new_rowset_meta_pb->set_num_segments(source_meta_pb.num_segments());
     
new_rowset_meta_pb->mutable_segment_ids()->CopyFrom(source_meta_pb.segment_ids());
     new_rowset_meta_pb->set_rowset_state(source_meta_pb.rowset_state());
+    if (source_meta_pb.has_commit_tso()) {
+        
new_rowset_meta_pb->mutable_commit_tso()->CopyFrom(source_meta_pb.commit_tso());
+    }
     new_rowset_meta_pb->mutable_segment_group_sizes()->CopyFrom(
             source_meta_pb.segment_group_sizes());
 
diff --git a/be/src/storage/snapshot/snapshot_manager.cpp 
b/be/src/storage/snapshot/snapshot_manager.cpp
index dd7ea1cb268..dab2df21e16 100644
--- a/be/src/storage/snapshot/snapshot_manager.cpp
+++ b/be/src/storage/snapshot/snapshot_manager.cpp
@@ -375,6 +375,9 @@ Status SnapshotManager::_rename_rowset_id(const 
RowsetMetaPB& rs_meta_pb,
                                    "failed to build rowset when rename rowset 
id");
     RETURN_IF_ERROR(new_rowset->load(false));
     new_rowset->rowset_meta()->to_rowset_pb(new_rs_meta_pb);
+    if (rs_meta_pb.has_commit_tso()) {
+        
new_rs_meta_pb->mutable_commit_tso()->CopyFrom(rs_meta_pb.commit_tso());
+    }
     RETURN_IF_ERROR(org_rowset->remove());
     return Status::OK();
 }
diff --git a/be/src/storage/task/index_builder.cpp 
b/be/src/storage/task/index_builder.cpp
index 427c30cdcaf..c639e423bd9 100644
--- a/be/src/storage/task/index_builder.cpp
+++ b/be/src/storage/task/index_builder.cpp
@@ -410,6 +410,9 @@ Status IndexBuilder::update_inverted_index_info() {
         rowset_meta->set_num_segments(input_rowset_meta->num_segments());
         
rowset_meta->set_segments_overlap(input_rowset_meta->segments_overlap());
         rowset_meta->set_rowset_state(input_rowset_meta->rowset_state());
+        if (input_rowset_meta->has_commit_tso()) {
+            rowset_meta->set_commit_tso(input_rowset_meta->commit_tso());
+        }
         std::vector<KeyBoundsPB> key_bounds;
         RETURN_IF_ERROR(input_rowset->get_segments_key_bounds(&key_bounds));
         rowset_meta->set_segments_key_bounds_truncated(
diff --git a/be/test/cloud/cloud_snapshot_mgr_test.cpp 
b/be/test/cloud/cloud_snapshot_mgr_test.cpp
index d53157ed168..5f7df2a5981 100644
--- a/be/test/cloud/cloud_snapshot_mgr_test.cpp
+++ b/be/test/cloud/cloud_snapshot_mgr_test.cpp
@@ -19,6 +19,9 @@
 
 #include <gtest/gtest.h>
 
+#include <optional>
+#include <vector>
+
 #include "cloud/cloud_storage_engine.h"
 #include "io/fs/remote_file_system.h"
 #include "io/fs/s3_file_system.h"
@@ -169,6 +172,63 @@ TEST_F(CloudSnapshotMgrTest, TestConvertRowsets) {
     EXPECT_TRUE(status.ok());
 }
 
+TEST_F(CloudSnapshotMgrTest, ConvertRowsetsPreservesCommitTso) {
+    const std::vector<std::optional<TsoRange>> cases = {std::nullopt, 
TsoRange(-1, -1),
+                                                        TsoRange(100, 100), 
TsoRange(100, 200)};
+    for (size_t i = 0; i < cases.size(); ++i) {
+        SCOPED_TRACE(i);
+        TabletMetaPB input;
+        input.set_tablet_id(1000);
+        input.set_schema_hash(123456);
+        *input.mutable_tablet_uid() = TabletUid::gen_uid().to_proto();
+        auto* schema = input.mutable_schema();
+        schema->set_keys_type(DUP_KEYS);
+        schema->set_num_short_key_columns(1);
+        schema->set_num_rows_per_row_block(1024);
+        auto* key = schema->add_column();
+        key->set_unique_id(1);
+        key->set_name("k1");
+        key->set_type("INT");
+        key->set_is_key(true);
+        key->set_aggregation("NONE");
+        auto* source = input.add_rs_metas();
+        source->set_rowset_id(0);
+        source->set_rowset_id_v2(_engine->next_rowset_id().to_string());
+        source->set_tablet_id(1000);
+        source->set_rowset_type(BETA_ROWSET);
+        source->set_rowset_state(VISIBLE);
+        source->set_start_version(7);
+        source->set_end_version(
+                cases[i].has_value() && cases[i]->start_tso() != 
cases[i]->end_tso() ? 8 : 7);
+        source->set_num_segments(1);
+        source->add_segment_ids(7);
+        source->set_num_rows(10);
+        source->set_newest_write_timestamp(1000);
+        source->mutable_tablet_schema()->CopyFrom(*schema);
+        if (cases[i].has_value()) {
+            source->mutable_commit_tso()->set_start_tso(cases[i]->start_tso());
+            source->mutable_commit_tso()->set_end_tso(cases[i]->end_tso());
+        }
+        auto tablet_meta = std::make_shared<TabletMeta>();
+        tablet_meta->init_from_pb(input);
+        auto target = std::make_shared<CloudTablet>(*_engine, tablet_meta);
+        StorageResource resource {_fs};
+        std::unordered_map<std::string, std::string> file_mapping;
+        TabletMetaPB output;
+        auto st = _snapshot_mgr->convert_rowsets(&output, input, 3000, target, 
resource,
+                                                 file_mapping);
+        ASSERT_TRUE(st.ok()) << st;
+        ASSERT_EQ(1, output.rs_metas_size());
+        const auto& converted = output.rs_metas(0);
+        EXPECT_NE(source->rowset_id_v2(), converted.rowset_id_v2());
+        EXPECT_EQ(source->has_commit_tso(), converted.has_commit_tso());
+        EXPECT_EQ(source->commit_tso().SerializeAsString(),
+                  converted.commit_tso().SerializeAsString());
+        EXPECT_EQ(source->start_version(), converted.start_version());
+        EXPECT_EQ(source->end_version(), converted.end_version());
+    }
+}
+
 TEST_F(CloudSnapshotMgrTest, TestRenameIndexIds) {
     TabletSchemaPB source_schema_pb;
     source_schema_pb.set_keys_type(KeysType::DUP_KEYS);
diff --git a/be/test/storage/index/index_builder_test.cpp 
b/be/test/storage/index/index_builder_test.cpp
index 809c1e9fe03..f688d0d1f31 100644
--- a/be/test/storage/index/index_builder_test.cpp
+++ b/be/test/storage/index/index_builder_test.cpp
@@ -21,19 +21,28 @@
 #include <gtest/gtest.h>
 
 #include <filesystem>
+#include <optional>
 #include <set>
 
 #include "common/config.h"
+#include "core/assert_cast.h"
+#include "core/column/column_vector.h"
 #include "storage/index/index_file_reader.h"
 #include "storage/index/index_writer.h"
 #include "storage/index/snii/query/term_query.h"
 #include "storage/olap_common.h"
+#include "storage/predicate/block_column_predicate.h"
+#include "storage/predicate/comparison_predicate.h"
 #include "storage/rowset/beta_rowset.h"
 #include "storage/rowset/rowset_factory.h"
 #include "storage/rowset/rowset_writer_context.h"
+#include "storage/segment/column_reader.h"
+#include "storage/segment/segment.h"
 #include "storage/storage_engine.h"
 #include "storage/tablet/tablet_fwd.h"
 #include "storage/tablet/tablet_schema.h"
+#include "storage/tablet/tablet_schema_helper.h"
+#include "storage/utils.h"
 #include "util/debug_points.h"
 
 namespace doris {
@@ -195,7 +204,10 @@ protected:
         rs_meta->set_tablet_schema(tablet_schema);
     }
 
-    void prepare_single_index_build(int64_t rowset_id) {
+    void prepare_single_index_build(int64_t rowset_id, bool with_commit_tso = 
false) {
+        if (with_commit_tso) {
+            _tablet_schema->append_column(*create_commit_tso_column(3));
+        }
         auto tablet_path = _absolute_dir + "/" + std::to_string(rowset_id);
         _tablet->_tablet_path = tablet_path;
         
ASSERT_TRUE(io::global_local_filesystem()->delete_directory(tablet_path).ok());
@@ -223,6 +235,9 @@ protected:
             int32_t k2 = i;
             columns[0]->insert_data(reinterpret_cast<const char*>(&k1), 
sizeof(k1));
             columns[1]->insert_data(reinterpret_cast<const char*>(&k2), 
sizeof(k2));
+            if (with_commit_tso) {
+                columns[2]->insert_default();
+            }
         }
         block.set_columns(std::move(columns));
         ASSERT_TRUE(rowset_writer->add_block(&block).ok());
@@ -691,6 +706,110 @@ TEST_F(IndexBuilderTest, BasicBuildTest) {
     EXPECT_EQ(builder._alter_index_ids.size(), 1);
 }
 
+class IndexBuilderCommitTsoTest : public IndexBuilderTest,
+                                  public 
testing::WithParamInterface<std::optional<TsoRange>> {
+protected:
+    void check_rowset(const RowsetSharedPtr& rowset) {
+        const auto& tso = GetParam();
+        const int64_t expected_tso = tso.has_value() && tso->end_tso() != -1 ? 
tso->end_tso() : 0;
+        EXPECT_EQ(tso.has_value(), rowset->rowset_meta()->has_commit_tso());
+        if (tso.has_value()) {
+            EXPECT_EQ(*tso, rowset->rowset_meta()->commit_tso());
+        }
+        // Open a fresh segment so an old cached constant reader cannot hide 
lost metadata.
+        auto segment_path = rowset->segment(0).path();
+        ASSERT_TRUE(segment_path.has_value()) << segment_path.error();
+        segment_v2::SegmentSharedPtr segment;
+        auto st = segment_v2::Segment::open(
+                io::global_local_filesystem(), segment_path.value(), 
_tablet->tablet_id(), 0,
+                rowset->rowset_id(), rowset->tablet_schema(), 
io::FileReaderOptions {}, &segment);
+        ASSERT_TRUE(st.ok()) << st;
+        OlapReaderStatistics stats;
+        StorageReadOptions options;
+        options.stats = &stats;
+        options.version = rowset->version();
+        options.commit_tso = rowset->rowset_meta()->commit_tso();
+        options.io_ctx.reader_type = ReaderType::READER_QUERY;
+        segment_v2::ColumnIteratorUPtr iter;
+        st = segment->new_column_iterator(rowset->tablet_schema()->column(2), 
&iter, &options);
+        ASSERT_TRUE(st.ok()) << st;
+        segment_v2::ColumnIteratorOptions iter_options;
+        iter_options.stats = &stats;
+        iter_options.file_reader = segment->file_reader().get();
+        iter_options.io_ctx.reader_type = ReaderType::READER_QUERY;
+        ASSERT_TRUE(iter->init(iter_options).ok());
+        ASSERT_TRUE(iter->seek_to_ordinal(0).ok());
+        MutableColumnPtr values = ColumnInt64::create();
+        size_t rows = 8;
+        bool has_null = true;
+        st = iter->next_batch(&rows, values, &has_null);
+        ASSERT_TRUE(st.ok()) << st;
+        ASSERT_EQ(8, rows);
+        ASSERT_EQ(8, values->size());
+        EXPECT_FALSE(has_null);
+        const auto& column = assert_cast<const ColumnInt64&>(*values);
+        for (size_t i = 0; i < rows; ++i) {
+            EXPECT_EQ(expected_tso, column.get_element(i));
+        }
+        auto read_schema = 
std::make_shared<ReadSchema>(rowset->tablet_schema()->columns());
+        auto predicate = AndBlockColumnPredicate::create_shared();
+        std::shared_ptr<ColumnPredicate> greater_than_zero =
+                std::make_shared<ComparisonPredicateBase<TYPE_BIGINT, 
PredicateType::GT>>(
+                        2, COMMIT_TSO_COL, 
Field::create_field<TYPE_BIGINT>(0));
+        predicate->add_column_predicate(
+                SingleColumnBlockPredicate::create_unique(greater_than_zero));
+        options.col_id_to_predicates.emplace(2, predicate);
+        std::unique_ptr<RowwiseIterator> filtered;
+        st = segment->new_iterator(read_schema, options, &filtered);
+        ASSERT_TRUE(st.ok()) << st;
+        EXPECT_EQ(expected_tso == 0, filtered->empty());
+    }
+};
+
+TEST_P(IndexBuilderCommitTsoTest, BuildPreservesCommitTso) {
+    ASSERT_NO_FATAL_FAILURE(prepare_single_index_build(16605, true));
+    auto source = _tablet->get_rowset_by_version(Version(10, 10));
+    ASSERT_NE(source, nullptr);
+    if (GetParam().has_value()) {
+        source->rowset_meta()->set_commit_tso(*GetParam());
+    }
+
+    ASSERT_NO_FATAL_FAILURE(check_rowset(source));
+    auto st = build_single_index();
+    ASSERT_TRUE(st.ok()) << st;
+    auto built = _tablet->get_rowset_by_version(Version(10, 10));
+    ASSERT_NE(built, nullptr);
+    EXPECT_NE(source->rowset_id(), built->rowset_id());
+    ASSERT_NO_FATAL_FAILURE(check_rowset(built));
+}
+
+TEST_P(IndexBuilderCommitTsoTest, DropPreservesCommitTso) {
+    ASSERT_NO_FATAL_FAILURE(prepare_single_index_build(16606, true));
+    // Seed the TSO after building the index so DROP has a valid input even on 
the old implementation.
+    auto st = build_single_index();
+    ASSERT_TRUE(st.ok()) << st;
+    auto source = _tablet->get_rowset_by_version(Version(10, 10));
+    ASSERT_NE(source, nullptr);
+    if (GetParam().has_value()) {
+        source->rowset_meta()->set_commit_tso(*GetParam());
+    }
+    ASSERT_NO_FATAL_FAILURE(check_rowset(source));
+
+    IndexBuilder drop_builder(*_engine_ref, _tablet, _columns, _alter_indexes, 
true);
+    ASSERT_TRUE(drop_builder.init().ok());
+    st = drop_builder.do_build_inverted_index();
+    ASSERT_TRUE(st.ok()) << st;
+    auto dropped = _tablet->get_rowset_by_version(Version(10, 10));
+    ASSERT_NE(dropped, nullptr);
+    EXPECT_NE(source->rowset_id(), dropped->rowset_id());
+    ASSERT_NO_FATAL_FAILURE(check_rowset(dropped));
+}
+
+INSTANTIATE_TEST_SUITE_P(CommitTso, IndexBuilderCommitTsoTest,
+                         testing::Values(std::optional<TsoRange> {},
+                                         std::optional<TsoRange> {TsoRange(-1, 
-1)},
+                                         std::optional<TsoRange> 
{TsoRange(100, 100)}));
+
 TEST_F(IndexBuilderTest, HandleSingleRowsetPreservesOrdinaryAppendFailure) {
     prepare_single_index_build(16604);
     ScopedIndexBuilderDebugPoints debug_points;
diff --git a/be/test/storage/snapshot/snapshot_manager_test.cpp 
b/be/test/storage/snapshot/snapshot_manager_test.cpp
index 00694eacbc6..b4258860b40 100644
--- a/be/test/storage/snapshot/snapshot_manager_test.cpp
+++ b/be/test/storage/snapshot/snapshot_manager_test.cpp
@@ -22,6 +22,7 @@
 #include <gtest/gtest-test-part.h>
 
 #include <memory>
+#include <optional>
 #include <string>
 #include <vector>
 
@@ -113,6 +114,69 @@ TEST_F(SnapshotManagerTest, 
TestConvertRowsetIdsInvalidDir) {
     EXPECT_EQ(result.error().code(), ErrorCode::DIR_NOT_EXIST);
 }
 
+TEST_F(SnapshotManagerTest, ConvertRowsetIdsPreservesCommitTso) {
+    const std::vector<std::optional<TsoRange>> cases = {std::nullopt, 
TsoRange(-1, -1),
+                                                        TsoRange(100, 100), 
TsoRange(100, 200)};
+    for (size_t i = 0; i < cases.size(); ++i) {
+        SCOPED_TRACE(i);
+        const auto clone_dir = _engine_data_path + "/commit_tso_" + 
std::to_string(i);
+        
ASSERT_TRUE(io::global_local_filesystem()->create_directory(clone_dir).ok());
+        auto tablet_meta = testutil::create_tablet_meta_pb(10006, 12346, 1, 
1000, 100);
+        auto* schema = tablet_meta.mutable_schema();
+        schema->set_keys_type(DUP_KEYS);
+        schema->set_num_short_key_columns(1);
+        schema->set_num_rows_per_row_block(1024);
+        schema->set_compress_kind(COMPRESS_LZ4);
+        testutil::add_column_pb(schema, 0, "k1", "INT", true, false);
+
+        RowsetMetaPB source;
+        source.set_rowset_id(0);
+        source.set_rowset_id_v2(_engine->next_rowset_id().to_string());
+        source.set_tablet_id(10006);
+        source.set_tablet_schema_hash(12346);
+        source.set_rowset_type(BETA_ROWSET);
+        source.set_rowset_state(VISIBLE);
+        source.set_start_version(7);
+        source.set_end_version(
+                cases[i].has_value() && cases[i]->start_tso() != 
cases[i]->end_tso() ? 8 : 7);
+        source.set_num_rows(0);
+        source.set_num_segments(0);
+        source.set_empty(true);
+        source.set_segments_overlap_pb(NONOVERLAPPING);
+        source.mutable_tablet_schema()->CopyFrom(*schema);
+        if (cases[i].has_value()) {
+            source.mutable_commit_tso()->set_start_tso(cases[i]->start_tso());
+            source.mutable_commit_tso()->set_end_tso(cases[i]->end_tso());
+        }
+        tablet_meta.add_rs_metas()->CopyFrom(source);
+        auto* stale = tablet_meta.add_stale_rs_metas();
+        stale->CopyFrom(source);
+        stale->set_rowset_id_v2(_engine->next_rowset_id().to_string());
+        stale->set_start_version(3);
+        stale->set_end_version(source.start_version() == source.end_version() 
? 3 : 4);
+
+        const auto meta_file = clone_dir + "/20006.hdr";
+        ASSERT_TRUE(TabletMeta::save(meta_file, tablet_meta).ok());
+        auto result =
+                _engine->snapshot_mgr()->convert_rowset_ids(clone_dir, 20006, 
2, 2000, 200, 65432);
+        ASSERT_TRUE(result.has_value()) << result.error();
+        TabletMetaPB converted;
+        ASSERT_TRUE(TabletMeta::load_from_file(meta_file, &converted).ok());
+        ASSERT_EQ(1, converted.rs_metas_size());
+        ASSERT_EQ(1, converted.stale_rs_metas_size());
+        auto check = [](const RowsetMetaPB& before, const RowsetMetaPB& after) 
{
+            EXPECT_NE(before.rowset_id_v2(), after.rowset_id_v2());
+            EXPECT_EQ(before.has_commit_tso(), after.has_commit_tso());
+            EXPECT_EQ(before.commit_tso().SerializeAsString(),
+                      after.commit_tso().SerializeAsString());
+            EXPECT_EQ(before.start_version(), after.start_version());
+            EXPECT_EQ(before.end_version(), after.end_version());
+        };
+        check(source, converted.rs_metas(0));
+        check(tablet_meta.stale_rs_metas(0), converted.stale_rs_metas(0));
+    }
+}
+
 TEST_F(SnapshotManagerTest, TestConvertRowsetIdsNormal) {
     std::string clone_dir = _engine_data_path + "/clone_dir_local";
     
EXPECT_TRUE(io::global_local_filesystem()->create_directory(clone_dir).ok());


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

Reply via email to