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

liaoxin01 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 3c459fc1fde [fix](filecache) Unify the file cache TTL deadline across 
all paths (#68090)
3c459fc1fde is described below

commit 3c459fc1fde70ef0940fa2e1912a0be10ed016f4
Author: Xin Liao <[email protected]>
AuthorDate: Thu Sep 17 21:19:04 2026 +0800

    [fix](filecache) Unify the file cache TTL deadline across all paths (#68090)
    
    ### What problem does this PR solve?
    
    Related PR: #57922
    
    Problem Summary:
    
    #57922 moved file cache TTL management from block granularity to tablet
    granularity. As part of that change every path stopped computing an
    expiration
    timestamp and started passing the raw ttl_seconds instead:
    
        expiration_time = tablet_meta->ttl_seconds();
    
    The receiving side was never updated. KeyMeta::expiration_time still
    documents
    itself as an absolute time, and that value is what gets persisted into
    the block
    meta store and compared by the cache consistency check. Those fields now
    hold a
    duration such as 3600 rather than a timestamp, so the only information
    left in
    them is "non-zero means this is a TTL block".
    
    Meanwhile the expiration decision itself moved into
    BlockFileCacheTtlMgr, which
    sweeps by tablet_ctime + ttl < now, anchored at the tablet creation
    time.
    
    Two consequences:
    
    1. Blocks are created as TTL blocks regardless of whether the tablet is
    still
    within its TTL window. Once a tablet is past creation_time +
    ttl_seconds, the
    load, compaction and query paths keep putting blocks into the TTL queue
    and
    the background sweep keeps pulling them straight back out. Each
    conversion
    writes the block meta store and takes the cache lock to move the block
       between LRU queues, and the loop never terminates for the rest of the
       tablet's life.
    
    2. The two writer paths disagreed with each other. The packed file path
    in
    RowsetWriterContext still computed an absolute newest_write_timestamp +
    ttl,
    while the regular file writer path passed the relative ttl_seconds, so
    segment files and packed small files of the same rowset were stamped
    with
    different kinds of value. check_file_cache_consistency reports that as
       EXPIRATION_TIME_INCONSISTENT.
    
    This PR keeps the tablet creation time as the anchor and gives that
    deadline a
    single definition, TabletMeta::file_cache_ttl_expiration_time(), which
    returns
    the absolute creation_time + ttl_seconds, or 0 when the tablet has no
    TTL or is
    already past the deadline. The load, compaction, schema change, query
    and warm
    up paths all stamp the blocks they create with that value, so a block's
    recorded
    expiration time now agrees with the sweep that acts on it, and a tablet
    past its
    deadline has its blocks created as NORMAL directly instead of cycling
    through
    the TTL queue.
    
    BlockFileCacheTtlMgr additionally no longer promotes the blocks of an
    already
    expired tablet on the edge where it first sees that tablet, which
    previously
    caused one full promote/demote pass per TTL tablet after every BE
    restart.
    
    The existing regression tests could not catch this: they all create a
    table and
    load into it immediately, so the tablet creation time and the data write
    time
    coincide and the bug is invisible.
    
    ### Release note
    
    Fixed the file cache TTL expiration time being computed inconsistently
    across
    the load, compaction, query and warm up paths. Data belonging to a
    tablet that
    is past its TTL deadline is now written directly into the normal cache
    queue
    instead of repeatedly entering and leaving the TTL queue.
---
 be/src/cloud/cloud_internal_service.cpp            |   2 +-
 be/src/cloud/cloud_rowset_builder.cpp              |   2 +-
 be/src/cloud/cloud_schema_change_job.cpp           |   2 +-
 be/src/cloud/cloud_tablet.cpp                      |   2 +-
 be/src/cloud/cloud_warm_up_manager.cpp             |   2 +-
 be/src/exec/scan/olap_scanner.cpp                  |   2 +-
 be/src/io/fs/file_writer.h                         |   2 +-
 be/src/storage/compaction/compaction.cpp           |   2 +-
 be/src/storage/rowset/beta_rowset_reader.cpp       |   2 +-
 be/src/storage/rowset/rowset_reader_context.h      |   5 +-
 be/src/storage/rowset/rowset_writer_context.h      |  12 +-
 be/src/storage/tablet/base_tablet.h                |   4 +
 be/src/storage/tablet/tablet_meta.cpp              |  19 ++
 be/src/storage/tablet/tablet_meta.h                |  11 ++
 be/src/storage/tablet/tablet_reader.cpp            |   2 +-
 be/test/storage/tablet/tablet_meta_test.cpp        |  39 ++++
 .../cloud_p0/cache/ttl/alter_ttl_seconds.groovy    |   5 +-
 .../cache/ttl/test_ttl_expired_tablet.groovy       | 207 +++++++++++++++++++++
 18 files changed, 305 insertions(+), 17 deletions(-)

diff --git a/be/src/cloud/cloud_internal_service.cpp 
b/be/src/cloud/cloud_internal_service.cpp
index 4c0fd49b8d1..e34fa03134f 100644
--- a/be/src/cloud/cloud_internal_service.cpp
+++ b/be/src/cloud/cloud_internal_service.cpp
@@ -1270,7 +1270,7 @@ void 
CloudInternalServiceImpl::warm_up_rowset(google::protobuf::RpcController* c
                       << " us, tablet_id: " << rs_meta.tablet_id()
                       << ", rowset_id: " << rowset_id.to_string();
         }
-        int64_t expiration_time = tablet_meta->ttl_seconds();
+        int64_t expiration_time = 
tablet_meta->file_cache_ttl_expiration_time();
 
         if (!tablet->add_rowset_warmup_state(rs_meta, 
WarmUpTriggerSource::EVENT_DRIVEN)) {
             LOG(INFO) << "found duplicate warmup task for rowset " << 
rowset_id.to_string()
diff --git a/be/src/cloud/cloud_rowset_builder.cpp 
b/be/src/cloud/cloud_rowset_builder.cpp
index 8730b86a555..d9065c340ec 100644
--- a/be/src/cloud/cloud_rowset_builder.cpp
+++ b/be/src/cloud/cloud_rowset_builder.cpp
@@ -107,7 +107,7 @@ Status CloudRowsetBuilder::init() {
     context.write_file_cache = _req.write_file_cache;
     context.partial_update_info = _partial_update_info;
     context.write_binlog_opt().enable = _req.write_req_type == 
WriteRequestType::ROW_BINLOG;
-    context.file_cache_ttl_sec = _tablet->ttl_seconds();
+    context.file_cache_expiration_time = 
_tablet->file_cache_ttl_expiration_time();
     context.storage_resource = 
_engine.get_storage_resource(_req.storage_vault_id);
     if (!context.storage_resource) {
         return Status::InternalError("vault id not found, maybe not sync, 
vault id {}",
diff --git a/be/src/cloud/cloud_schema_change_job.cpp 
b/be/src/cloud/cloud_schema_change_job.cpp
index 2caeb78bd9f..c269f7dcb1b 100644
--- a/be/src/cloud/cloud_schema_change_job.cpp
+++ b/be/src/cloud/cloud_schema_change_job.cpp
@@ -405,7 +405,7 @@ Status 
CloudSchemaChangeJob::_convert_historical_rowsets(const SchemaChangeParam
         // like the load and compaction output does. Otherwise it is cached in 
the
         // NORMAL/INDEX queues here, while every warm-up path downloads it 
into the TTL
         // queue on the destination cluster.
-        context.file_cache_ttl_sec = _new_tablet->ttl_seconds();
+        context.file_cache_expiration_time = 
_new_tablet->file_cache_ttl_expiration_time();
         context.tablet = _new_tablet;
         if (!context.storage_resource) {
             return Status::InternalError("vault id not found, maybe not sync, 
vault id {}",
diff --git a/be/src/cloud/cloud_tablet.cpp b/be/src/cloud/cloud_tablet.cpp
index de3394688af..aeed9dac426 100644
--- a/be/src/cloud/cloud_tablet.cpp
+++ b/be/src/cloud/cloud_tablet.cpp
@@ -1979,7 +1979,7 @@ void 
CloudTablet::_add_rowsets_directly(std::vector<RowsetSharedPtr>& rowsets,
                     continue;
                 }
 
-                int64_t expiration_time = _tablet_meta->ttl_seconds();
+                int64_t expiration_time = 
_tablet_meta->file_cache_ttl_expiration_time();
                 g_file_cache_cloud_tablet_submitted_segment_num << 1;
                 if (seg.file_size() > 0) {
                     g_file_cache_cloud_tablet_submitted_segment_size << 
seg.file_size();
diff --git a/be/src/cloud/cloud_warm_up_manager.cpp 
b/be/src/cloud/cloud_warm_up_manager.cpp
index 747b406cc7b..bbfc8f5c156 100644
--- a/be/src/cloud/cloud_warm_up_manager.cpp
+++ b/be/src/cloud/cloud_warm_up_manager.cpp
@@ -297,7 +297,7 @@ void CloudWarmUpManager::handle_jobs() {
                     continue;
                 }
 
-                int64_t expiration_time = tablet_meta->ttl_seconds();
+                int64_t expiration_time = 
tablet_meta->file_cache_ttl_expiration_time();
                 if (!tablet->add_rowset_warmup_state(*rs, 
WarmUpTriggerSource::JOB)) {
                     LOG(INFO) << "found duplicate warmup task for rowset " << 
rs->rowset_id()
                               << ", skip it";
diff --git a/be/src/exec/scan/olap_scanner.cpp 
b/be/src/exec/scan/olap_scanner.cpp
index 787a70aaefc..90e9d94a2ae 100644
--- a/be/src/exec/scan/olap_scanner.cpp
+++ b/be/src/exec/scan/olap_scanner.cpp
@@ -281,7 +281,7 @@ Status OlapScanner::_prepare_impl() {
         _tablet_reader_params.collection_statistics = 
std::make_shared<CollectionStatistics>();
 
         auto io_ctx = build_score_runtime_collection_io_context(
-                _state, _tablet_reader_params.reader_type, 
tablet->ttl_seconds(),
+                _state, _tablet_reader_params.reader_type, 
tablet->file_cache_ttl_expiration_time(),
                 &_tablet_reader->mutable_stats()->file_cache_stats);
 
         RETURN_IF_ERROR(_tablet_reader_params.collection_statistics->collect(
diff --git a/be/src/io/fs/file_writer.h b/be/src/io/fs/file_writer.h
index 9402fdef183..92e4475d637 100644
--- a/be/src/io/fs/file_writer.h
+++ b/be/src/io/fs/file_writer.h
@@ -47,7 +47,7 @@ struct FileWriterOptions {
     bool allow_adaptive_file_cache_write = true;
     bool is_cold_data = false;
     bool sync_file_data = true;              // Whether flush data into 
storage system
-    uint64_t file_cache_expiration_time = 0; // Relative time
+    uint64_t file_cache_expiration_time = 0; // Absolute time, 0 means no TTL
     uint64_t approximate_bytes_to_write = 0; // Approximate bytes to write, 
used for file cache
 };
 
diff --git a/be/src/storage/compaction/compaction.cpp 
b/be/src/storage/compaction/compaction.cpp
index 2c91d5ee3a3..0186fa9b0bf 100644
--- a/be/src/storage/compaction/compaction.cpp
+++ b/be/src/storage/compaction/compaction.cpp
@@ -2383,7 +2383,7 @@ Status 
CloudCompactionMixin::construct_output_rowset_writer(RowsetWriterContext&
     // TODO(gavin): Ensure that the retention of hot data is implemented with 
precision.
 
     ctx.write_file_cache = should_cache_compaction_output();
-    ctx.file_cache_ttl_sec = _tablet->ttl_seconds();
+    ctx.file_cache_expiration_time = _tablet->file_cache_ttl_expiration_time();
     ctx.approximate_bytes_to_write = _input_rowsets_total_size;
 
     // Set fine-grained control: only write index files to cache if configured
diff --git a/be/src/storage/rowset/beta_rowset_reader.cpp 
b/be/src/storage/rowset/beta_rowset_reader.cpp
index 209ca64841a..fea4e22d6dc 100644
--- a/be/src/storage/rowset/beta_rowset_reader.cpp
+++ b/be/src/storage/rowset/beta_rowset_reader.cpp
@@ -242,7 +242,7 @@ Status 
BetaRowsetReader::get_segment_iterators(RowsetReaderContext* read_context
         _read_options.condition_cache_digest = 
_read_context->condition_cache_digest;
     }
 
-    _read_options.io_ctx.expiration_time = read_context->ttl_seconds;
+    _read_options.io_ctx.expiration_time = 
read_context->file_cache_expiration_time;
 
     bool enable_segment_cache = true;
     auto* state = read_context->runtime_state;
diff --git a/be/src/storage/rowset/rowset_reader_context.h 
b/be/src/storage/rowset/rowset_reader_context.h
index cc6da6aa2f3..b72c5d58c4c 100644
--- a/be/src/storage/rowset/rowset_reader_context.h
+++ b/be/src/storage/rowset/rowset_reader_context.h
@@ -94,7 +94,10 @@ struct RowsetReaderContext {
     RowsetId rowset_id;
     // slots that cast may be eliminated in storage layer
     std::map<std::string, DataTypePtr> target_cast_type_for_variants;
-    int64_t ttl_seconds = 0;
+    // Absolute timestamp (seconds since epoch) after which cache blocks 
filled by this
+    // read stop being TTL protected; 0 means no TTL.
+    // See TabletMeta::file_cache_ttl_expiration_time().
+    int64_t file_cache_expiration_time = 0;
 
     std::map<ColumnId, VExprContextSPtr> virtual_column_exprs;
 
diff --git a/be/src/storage/rowset/rowset_writer_context.h 
b/be/src/storage/rowset/rowset_writer_context.h
index f3dd1d8bdf3..08e7b31cace 100644
--- a/be/src/storage/rowset/rowset_writer_context.h
+++ b/be/src/storage/rowset/rowset_writer_context.h
@@ -127,7 +127,11 @@ struct RowsetWriterContext {
     /// begin file cache opts
     bool write_file_cache = false;
     bool is_hot_data = false;
-    uint64_t file_cache_ttl_sec = 0;
+    // Absolute timestamp (seconds since epoch) after which the cache blocks 
written by
+    // this rowset stop being TTL protected; 0 means no TTL. Always set it from
+    // BaseTablet::file_cache_ttl_expiration_time() so every writer agrees 
with the
+    // deadline BlockFileCacheTtlMgr sweeps by.
+    uint64_t file_cache_expiration_time = 0;
     uint64_t approximate_bytes_to_write = 0;
     // If true, compaction output only writes index files to file cache, not 
data files
     bool compaction_output_write_index_only = false;
@@ -277,9 +281,7 @@ struct RowsetWriterContext {
             append_info.rowset_id = rowset_id.to_string();
             append_info.first_segment_id = first_segment_id;
             append_info.txn_id = txn_id;
-            append_info.expiration_time = file_cache_ttl_sec > 0 && 
newest_write_timestamp > 0
-                                                  ? newest_write_timestamp + 
file_cache_ttl_sec
-                                                  : 0;
+            append_info.expiration_time = file_cache_expiration_time;
             fs = std::make_shared<io::PackedFileSystem>(fs, append_info);
         }
 
@@ -298,7 +300,7 @@ struct RowsetWriterContext {
     io::FileWriterOptions get_file_writer_options(FileType file_type = 
FileType::SEGMENT_FILE) {
         io::FileWriterOptions opts {.write_file_cache = write_file_cache,
                                     .is_cold_data = is_hot_data,
-                                    .file_cache_expiration_time = 
file_cache_ttl_sec,
+                                    .file_cache_expiration_time = 
file_cache_expiration_time,
                                     .approximate_bytes_to_write = 
approximate_bytes_to_write};
 
         if (config::enable_file_cache_write_index_file_only) {
diff --git a/be/src/storage/tablet/base_tablet.h 
b/be/src/storage/tablet/base_tablet.h
index 769f9bb28e2..5248be6cbac 100644
--- a/be/src/storage/tablet/base_tablet.h
+++ b/be/src/storage/tablet/base_tablet.h
@@ -78,6 +78,10 @@ public:
     KeysType keys_type() const { return 
_tablet_meta->tablet_schema()->keys_type(); }
     size_t num_key_columns() const { return 
_tablet_meta->tablet_schema()->num_key_columns(); }
     int64_t ttl_seconds() const { return _tablet_meta->ttl_seconds(); }
+    // See TabletMeta::file_cache_ttl_expiration_time().
+    int64_t file_cache_ttl_expiration_time() const {
+        return _tablet_meta->file_cache_ttl_expiration_time();
+    }
     // currently used by schema change, inverted index building, and cooldown
     std::timed_mutex& get_schema_change_lock() { return _schema_change_lock; }
     bool enable_unique_key_merge_on_write() const {
diff --git a/be/src/storage/tablet/tablet_meta.cpp 
b/be/src/storage/tablet/tablet_meta.cpp
index db8783c334b..d9bea9673a0 100644
--- a/be/src/storage/tablet/tablet_meta.cpp
+++ b/be/src/storage/tablet/tablet_meta.cpp
@@ -716,6 +716,25 @@ Status TabletMeta::save_meta(DataDir* data_dir) {
     return _save_meta(data_dir);
 }
 
+int64_t TabletMeta::file_cache_ttl_expiration_time() const {
+    int64_t ttl = ttl_seconds();
+    int64_t ctime = creation_time();
+    if (ttl <= 0 || ctime <= 0) {
+        return 0;
+    }
+    // FE caps file_cache_ttl_seconds at Long.MAX_VALUE / 2, so this cannot 
wrap, but a tablet
+    // meta that reached us from anywhere else still must not turn a huge ttl 
into a past
+    // deadline that silently downgrades the tablet to normal cache.
+    if (ctime > std::numeric_limits<int64_t>::max() - ttl) {
+        return std::numeric_limits<int64_t>::max();
+    }
+    int64_t expiration_time = ctime + ttl;
+    // Already past the deadline: report no TTL at all, so callers stamp the 
blocks they
+    // create as NORMAL right away instead of putting them in the TTL queue for
+    // BlockFileCacheTtlMgr to take straight back out again.
+    return expiration_time > UnixSeconds() ? expiration_time : 0;
+}
+
 Status TabletMeta::_save_meta(DataDir* data_dir) {
     // check if tablet uid is valid
     if (_tablet_uid.hi == 0 && _tablet_uid.lo == 0) {
diff --git a/be/src/storage/tablet/tablet_meta.h 
b/be/src/storage/tablet/tablet_meta.h
index 3b01f94df8a..f83b4bb3672 100644
--- a/be/src/storage/tablet/tablet_meta.h
+++ b/be/src/storage/tablet/tablet_meta.h
@@ -335,6 +335,17 @@ public:
         _ttl_seconds = ttl_seconds;
     }
 
+    // Absolute timestamp (seconds since epoch) at which this tablet's data 
stops being kept
+    // in the file cache TTL queue, or 0 when the tablet has no TTL or the 
deadline has
+    // already passed. The deadline is anchored at the tablet creation time, 
so every tablet
+    // of a table shares one deadline regardless of when each rowset was 
written.
+    //
+    // This is the single definition of that deadline. The load, compaction, 
schema change,
+    // query and warm up paths all stamp the cache blocks they create with 
this value, and
+    // BlockFileCacheTtlMgr expires those blocks by the very same value, so a 
block's
+    // recorded expiration time always agrees with the sweep that acts on it.
+    int64_t file_cache_ttl_expiration_time() const;
+
     int64_t avg_rs_meta_serialize_size() const { return 
_avg_rs_meta_serialize_size; }
 
     EncryptionAlgorithmPB encryption_algorithm() const { return 
_encryption_algorithm; }
diff --git a/be/src/storage/tablet/tablet_reader.cpp 
b/be/src/storage/tablet/tablet_reader.cpp
index a942f16d55c..1ac7af12a57 100644
--- a/be/src/storage/tablet/tablet_reader.cpp
+++ b/be/src/storage/tablet/tablet_reader.cpp
@@ -191,7 +191,7 @@ Status TabletReader::_capture_rs_readers(const 
ReaderParams& read_params) {
     _reader_context.output_columns = &read_params.output_columns;
     _reader_context.extra_columns = read_params.extra_columns;
     _reader_context.push_down_agg_type_opt = 
read_params.push_down_agg_type_opt;
-    _reader_context.ttl_seconds = _tablet->ttl_seconds();
+    _reader_context.file_cache_expiration_time = 
_tablet->file_cache_ttl_expiration_time();
     _reader_context.score_runtime = read_params.score_runtime;
     _reader_context.collection_statistics = read_params.collection_statistics;
 
diff --git a/be/test/storage/tablet/tablet_meta_test.cpp 
b/be/test/storage/tablet/tablet_meta_test.cpp
index 035426ac3e3..45c9da37624 100644
--- a/be/test/storage/tablet/tablet_meta_test.cpp
+++ b/be/test/storage/tablet/tablet_meta_test.cpp
@@ -20,6 +20,7 @@
 #include <gtest/gtest-message.h>
 #include <gtest/gtest-test-part.h>
 
+#include <limits>
 #include <memory>
 #include <string>
 #include <utility>
@@ -29,6 +30,7 @@
 #include "storage/rowset/rowset.h"
 #include "storage/tablet/tablet_schema.h"
 #include "testutil/mock_rowset.h"
+#include "util/time.h"
 
 namespace doris {
 
@@ -478,4 +480,41 @@ TEST(TabletMetaTest, 
TestDeleteBitmapSubsetAndAggWithSegmentList) {
     EXPECT_NE(subset_delete_map.get({rowset_id, 0, 9}, 
&segment_delete_bitmap), 0);
 }
 
+TEST(TabletMetaTest, FileCacheTtlExpirationTime) {
+    TabletMeta meta;
+
+    // No TTL configured: nothing to protect.
+    meta.set_creation_time(UnixSeconds() - 10);
+    meta.set_ttl_seconds(0);
+    EXPECT_EQ(0, meta.file_cache_ttl_expiration_time());
+
+    // Unknown creation time cannot anchor a deadline.
+    meta.set_creation_time(0);
+    meta.set_ttl_seconds(3600);
+    EXPECT_EQ(0, meta.file_cache_ttl_expiration_time());
+
+    // Live tablet: the deadline is creation time + ttl, an absolute timestamp.
+    int64_t ctime = UnixSeconds() - 10;
+    meta.set_creation_time(ctime);
+    meta.set_ttl_seconds(3600);
+    EXPECT_EQ(ctime + 3600, meta.file_cache_ttl_expiration_time());
+
+    // Past the deadline: report no TTL, so callers stamp new blocks as NORMAL 
instead of
+    // putting them in the TTL queue for the expiration sweep to take back out.
+    meta.set_creation_time(UnixSeconds() - 3600);
+    meta.set_ttl_seconds(60);
+    EXPECT_EQ(0, meta.file_cache_ttl_expiration_time());
+
+    // Exactly at the deadline counts as expired.
+    int64_t now = UnixSeconds();
+    meta.set_creation_time(now - 60);
+    meta.set_ttl_seconds(60);
+    EXPECT_EQ(0, meta.file_cache_ttl_expiration_time());
+
+    // A ttl large enough to overflow must not wrap into a past deadline.
+    meta.set_creation_time(UnixSeconds());
+    meta.set_ttl_seconds(std::numeric_limits<int64_t>::max());
+    EXPECT_EQ(std::numeric_limits<int64_t>::max(), 
meta.file_cache_ttl_expiration_time());
+}
+
 } // namespace doris
diff --git a/regression-test/suites/cloud_p0/cache/ttl/alter_ttl_seconds.groovy 
b/regression-test/suites/cloud_p0/cache/ttl/alter_ttl_seconds.groovy
index fd5d72fb842..491576e5601 100644
--- a/regression-test/suites/cloud_p0/cache/ttl/alter_ttl_seconds.groovy
+++ b/regression-test/suites/cloud_p0/cache/ttl/alter_ttl_seconds.groovy
@@ -139,7 +139,10 @@ suite("test_ttl_seconds") {
 
     load_customer_once("customer_ttl")
     def tabletIds = getTabletIds.call("customer_ttl")
-    waitForFileCacheType.call(tabletIds, "ttl", 15000L, 500L)
+    // No wait for the "ttl" type here. The TTL deadline is the tablet 
creation time plus
+    // file_cache_ttl_seconds, and with a 5s ttl the load itself outlives it, 
so most of the
+    // data is written straight into the normal queue and the table never has 
all of its
+    // blocks in the TTL queue at once.
     sleep(30000) // 30s
     getMetricsMethod.call() {
         respCode, body ->
diff --git 
a/regression-test/suites/cloud_p0/cache/ttl/test_ttl_expired_tablet.groovy 
b/regression-test/suites/cloud_p0/cache/ttl/test_ttl_expired_tablet.groovy
new file mode 100644
index 00000000000..80b86d63f83
--- /dev/null
+++ b/regression-test/suites/cloud_p0/cache/ttl/test_ttl_expired_tablet.groovy
@@ -0,0 +1,207 @@
+// 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.
+
+import org.codehaus.groovy.runtime.IOGroovyMethods
+
+// The file cache TTL deadline is the tablet creation time plus 
file_cache_ttl_seconds. Once a
+// tablet is past that deadline, the load and query paths must stamp the 
blocks they create as
+// non-TTL right away, instead of creating TTL blocks for the background sweep 
to demote again.
+//
+// This test deliberately runs with the TTL background threads turned down to 
a 10 minute
+// interval. That is what gives it teeth: the sweep would otherwise repair the 
cache type within
+// a second or two and the assertions would hold no matter what the write and 
read paths did.
+// With the sweep out of the way, the cache type observed here is purely the 
one chosen at
+// admission, so the expired-tablet cases below fail if the deadline is not 
applied at the source.
+suite("test_ttl_expired_tablet") {
+    def ttlSeconds = 30
+    def sweepOffMs = 600000
+
+    def custoBeConfig = [
+        enable_evict_file_cache_in_advance : false,
+        file_cache_enter_disk_resource_limit_mode_percent : 99,
+        // Long enough that neither TTL background thread reconciles anything 
while the test runs.
+        file_cache_background_ttl_gc_interval_ms : sweepOffMs,
+        file_cache_background_ttl_info_update_interval_ms : sweepOffMs,
+        file_cache_background_tablet_id_flush_interval_ms : 1000
+    ]
+
+    setBeConfigTemporary(custoBeConfig) {
+    sql "set global enable_auto_analyze = false"
+    sql "set global enable_audit_plugin = false"
+    def clusters = sql " SHOW CLUSTERS; "
+    assertTrue(!clusters.isEmpty())
+    def validCluster = clusters[0][0]
+    sql """use @${validCluster};""";
+
+    String[][] backends = sql """ show backends """
+    String backendId;
+    def backendIdToBackendIP = [:]
+    def backendIdToBackendHttpPort = [:]
+    for (String[] backend in backends) {
+        if (backend[9].equals("true") && 
backend[19].contains("${validCluster}")) {
+            backendIdToBackendIP.put(backend[0], backend[1])
+            backendIdToBackendHttpPort.put(backend[0], backend[4])
+        }
+    }
+    assertEquals(backendIdToBackendIP.size(), 1)
+
+    backendId = backendIdToBackendIP.keySet()[0]
+    def url = backendIdToBackendIP.get(backendId) + ":" + 
backendIdToBackendHttpPort.get(backendId) + 
"""/api/file_cache?op=clear&sync=true"""
+    def clearFileCache = { check_func ->
+        httpTest {
+            endpoint ""
+            uri url
+            op "get"
+            body ""
+            check check_func
+        }
+    }
+
+    def getTabletIds = { String tableName ->
+        def tablets = sql "show tablets from ${tableName}"
+        assertTrue(tablets.size() > 0, "No tablets found for table 
${tableName}")
+        tablets.collect { it[0] as Long }
+    }
+
+    // Counts of each cache type across the given tablets, e.g. [normal: 40, 
index: 8].
+    def cacheTypeCounts = { List<Long> tabletIds ->
+        def counts = [:].withDefault { 0 }
+        for (Long tabletId in tabletIds) {
+            def rows = sql "select type from 
information_schema.file_cache_info where tablet_id = ${tabletId}"
+            for (row in rows) {
+                counts[row[0].toString().toLowerCase()] += 1
+            }
+        }
+        return counts
+    }
+
+    def totalBlocks = { Map counts -> counts.values().sum() ?: 0 }
+
+    // Cache writes are asynchronous, so wait for the blocks to show up before 
judging their
+    // type. Only their existence is waited on; the type assertions stay 
strict.
+    def waitForAnyBlock = { List<Long> tabletIds, long timeoutMs ->
+        long start = System.currentTimeMillis()
+        while (System.currentTimeMillis() - start < timeoutMs) {
+            if (totalBlocks(cacheTypeCounts.call(tabletIds)) > 0) {
+                return
+            }
+            sleep(1000)
+        }
+        assertTrue(false, "timed out waiting for cached blocks of tablets 
${tabletIds}")
+    }
+
+    def waitForEmptyCache = { List<Long> tabletIds, long timeoutMs ->
+        long start = System.currentTimeMillis()
+        while (System.currentTimeMillis() - start < timeoutMs) {
+            if (totalBlocks(cacheTypeCounts.call(tabletIds)) == 0) {
+                return
+            }
+            sleep(1000)
+        }
+        assertTrue(false, "timed out waiting for an empty cache for tablets 
${tabletIds}")
+    }
+
+    def loadCustomerRows = { String table ->
+        def totalRows = 200
+        def batchSize = 100
+        def commentSuffix = ' ' + ('X' * 50)
+        for (int offset = 0; offset < totalRows; offset += batchSize) {
+            def sb = new StringBuilder()
+            int batchEnd = Math.min(totalRows, offset + batchSize)
+            for (int idx = offset; idx < batchEnd; idx++) {
+                def customerId = 10001 + idx
+                def customerName = String.format('Customer#%09d', customerId)
+                sb.append("""INSERT INTO ${table} VALUES (
+                    ${customerId},
+                    '${customerName}',
+                    'Address Line 1',
+                    15,
+                    '123-456-7890',
+                    12345.67,
+                    'AUTOMOBILE',
+                    'This is a test comment for the customer.${commentSuffix}'
+                    );
+                    """)
+            }
+            sql sb.toString()
+        }
+    }
+
+    def createCustomerTable = { String table, long ttl ->
+        def ddl = new 
File("""${context.file.parent}/../ddl/customer_ttl.sql""").text
+        sql (ddl.replace("customer_ttl", table)
+             + """ PROPERTIES("file_cache_ttl_seconds"="${ttl}") """)
+        sql """ alter table ${table} set ("disable_auto_compaction" = "true") 
"""
+    }
+
+    sql """ DROP TABLE IF EXISTS customer_ttl_expired """
+    sql """ DROP TABLE IF EXISTS customer_ttl_live """
+    clearFileCache.call() { respCode, body -> {} }
+    sleep(5000)
+
+    // 
---------------------------------------------------------------------------------------
+    // Case 1: load into a tablet that is already past its deadline.
+    // 
---------------------------------------------------------------------------------------
+    createCustomerTable.call("customer_ttl_expired", ttlSeconds)
+    // Let the deadline pass before a single row is written.
+    sleep((ttlSeconds + 15) * 1000L)
+
+    loadCustomerRows("customer_ttl_expired")
+    def expiredTablets = getTabletIds.call("customer_ttl_expired")
+
+    // Guard against a vacuous pass: the load must actually have cached 
something.
+    waitForAnyBlock.call(expiredTablets, 60000L)
+    def afterLoad = cacheTypeCounts.call(expiredTablets)
+    logger.info("cache types after loading an expired tablet: ${afterLoad}")
+    assertEquals(0, afterLoad['ttl'],
+                 "load path put blocks of an expired tablet into the TTL 
queue: ${afterLoad}")
+
+    // 
---------------------------------------------------------------------------------------
+    // Case 2: read from that tablet on a cold cache, so the query path admits 
the blocks.
+    // 
---------------------------------------------------------------------------------------
+    clearFileCache.call() { respCode, body -> {} }
+    // Without an empty cache the query below would be a hit and would admit 
nothing, leaving
+    // the read path untested.
+    waitForEmptyCache.call(expiredTablets, 60000L)
+
+    // sum() has to read the column data, unlike count(*) which can be 
answered from metadata.
+    sql """ select sum(C_ACCTBAL), count(C_COMMENT) from customer_ttl_expired 
"""
+
+    waitForAnyBlock.call(expiredTablets, 60000L)
+    def afterRead = cacheTypeCounts.call(expiredTablets)
+    logger.info("cache types after reading an expired tablet: ${afterRead}")
+    assertEquals(0, afterRead['ttl'],
+                 "read path put blocks of an expired tablet into the TTL 
queue: ${afterRead}")
+
+    // 
---------------------------------------------------------------------------------------
+    // Case 3: a tablet still inside its window must keep using the TTL queue. 
Without this the
+    // two cases above would also pass if the deadline were simply always 
reported as expired.
+    // 
---------------------------------------------------------------------------------------
+    createCustomerTable.call("customer_ttl_live", 3600)
+    loadCustomerRows("customer_ttl_live")
+    def liveTablets = getTabletIds.call("customer_ttl_live")
+
+    waitForAnyBlock.call(liveTablets, 60000L)
+    def liveCounts = cacheTypeCounts.call(liveTablets)
+    logger.info("cache types for a tablet inside its TTL window: 
${liveCounts}")
+    assertTrue(liveCounts['ttl'] > 0,
+               "tablet inside its TTL window got no TTL blocks: ${liveCounts}")
+
+    sql """ DROP TABLE IF EXISTS customer_ttl_expired """
+    sql """ DROP TABLE IF EXISTS customer_ttl_live """
+    }
+}


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

Reply via email to