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]