This is an automated email from the ASF dual-hosted git repository.
yiguolei pushed a commit to branch branch-4.1
in repository https://gitbox.apache.org/repos/asf/doris.git
The following commit(s) were added to refs/heads/branch-4.1 by this push:
new e6bb8f591c4 branch-4.1: [fix](filecache) Unify the file cache TTL
deadline across all paths (#68313)
e6bb8f591c4 is described below
commit e6bb8f591c44143efd8643cb23b1f6d5fec84b9f
Author: Xin Liao <[email protected]>
AuthorDate: Tue Sep 22 10:20:14 2026 +0800
branch-4.1: [fix](filecache) Unify the file cache TTL deadline across all
paths (#68313)
Pick apache/doris#68090
---
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 | 6 +-
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 | 8 +-
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(+), 27 deletions(-)
diff --git a/be/src/cloud/cloud_internal_service.cpp
b/be/src/cloud/cloud_internal_service.cpp
index 98c6c9790be..a423f8622f9 100644
--- a/be/src/cloud/cloud_internal_service.cpp
+++ b/be/src/cloud/cloud_internal_service.cpp
@@ -1202,7 +1202,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 e03d347d94f..29ece82d320 100644
--- a/be/src/cloud/cloud_rowset_builder.cpp
+++ b/be/src/cloud/cloud_rowset_builder.cpp
@@ -78,7 +78,7 @@ Status CloudRowsetBuilder::init() {
context.mow_context = mow_context;
context.write_file_cache = _req.write_file_cache;
context.partial_update_info = _partial_update_info;
- 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 32b4386e059..8066063dd28 100644
--- a/be/src/cloud/cloud_schema_change_job.cpp
+++ b/be/src/cloud/cloud_schema_change_job.cpp
@@ -391,7 +391,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 f2398128424..b5c5f0425a0 100644
--- a/be/src/cloud/cloud_tablet.cpp
+++ b/be/src/cloud/cloud_tablet.cpp
@@ -1908,11 +1908,7 @@ void
CloudTablet::_add_rowsets_directly(std::vector<RowsetSharedPtr>& rowsets,
continue;
}
- int64_t expiration_time = _tablet_meta->ttl_seconds() == 0 ||
-
rowset_meta->newest_write_timestamp() <= 0
- ? 0
- :
rowset_meta->newest_write_timestamp() +
-
_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 (rs->rowset_meta()->segment_file_size(seg_id) > 0) {
g_file_cache_cloud_tablet_submitted_segment_size
diff --git a/be/src/cloud/cloud_warm_up_manager.cpp
b/be/src/cloud/cloud_warm_up_manager.cpp
index bc9b30b6273..c8b2a8a9eb9 100644
--- a/be/src/cloud/cloud_warm_up_manager.cpp
+++ b/be/src/cloud/cloud_warm_up_manager.cpp
@@ -302,7 +302,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 04ede793ed5..06e6a08c771 100644
--- a/be/src/exec/scan/olap_scanner.cpp
+++ b/be/src/exec/scan/olap_scanner.cpp
@@ -292,7 +292,7 @@ Status OlapScanner::prepare() {
_tablet_reader_params.collection_statistics =
std::make_shared<CollectionStatistics>();
auto io_ctx = build_score_runtime_collection_io_context(
- _state, ReaderType::READER_QUERY, tablet->ttl_seconds(),
+ _state, ReaderType::READER_QUERY,
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 a60ce91da5b..449a49fbf64 100644
--- a/be/src/storage/compaction/compaction.cpp
+++ b/be/src/storage/compaction/compaction.cpp
@@ -1867,7 +1867,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 af3169eaf86..a28572d54dc 100644
--- a/be/src/storage/rowset/beta_rowset_reader.cpp
+++ b/be/src/storage/rowset/beta_rowset_reader.cpp
@@ -241,13 +241,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 > 0 &&
_rowset->rowset_meta()->newest_write_timestamp() > 0
- ? _rowset->rowset_meta()->newest_write_timestamp() +
read_context->ttl_seconds
- : 0;
- if (_read_options.io_ctx.expiration_time <= UnixSeconds()) {
- _read_options.io_ctx.expiration_time = 0;
- }
+ _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 066c5890af0..539a4576d65 100644
--- a/be/src/storage/rowset/rowset_reader_context.h
+++ b/be/src/storage/rowset/rowset_reader_context.h
@@ -88,7 +88,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;
std::map<ColumnId, size_t> vir_cid_to_idx_in_block;
diff --git a/be/src/storage/rowset/rowset_writer_context.h
b/be/src/storage/rowset/rowset_writer_context.h
index b65232f3c1c..11ed54ac433 100644
--- a/be/src/storage/rowset/rowset_writer_context.h
+++ b/be/src/storage/rowset/rowset_writer_context.h
@@ -104,7 +104,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;
@@ -218,9 +222,7 @@ struct RowsetWriterContext {
append_info.tablet_id = tablet_id;
append_info.rowset_id = rowset_id.to_string();
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);
}
@@ -239,7 +241,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 73e797610b9..9abad642144 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 345bf4546c8..15abf78b549 100644
--- a/be/src/storage/tablet/tablet_meta.cpp
+++ b/be/src/storage/tablet/tablet_meta.cpp
@@ -652,6 +652,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 0f6d2bc69c3..a9d06f8827a 100644
--- a/be/src/storage/tablet/tablet_meta.h
+++ b/be/src/storage/tablet/tablet_meta.h
@@ -317,6 +317,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 7e3f3ffd701..7e87419cc4d 100644
--- a/be/src/storage/tablet/tablet_reader.cpp
+++ b/be/src/storage/tablet/tablet_reader.cpp
@@ -184,7 +184,7 @@ Status TabletReader::_capture_rs_readers(const
ReaderParams& read_params) {
_reader_context.common_expr_ctxs_push_down =
read_params.common_expr_ctxs_push_down;
_reader_context.output_columns = &read_params.output_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 33b784f666f..77e2626d7d7 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 {
@@ -403,4 +405,41 @@ TEST(TabletMetaTest, TestDeleteBitmap) {
EXPECT_EQ(d.cardinality(), 500);
}
+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]