enableTableCachePurge

语法

enableTableCachePurge(table, [cacheSize],[cachePurgeTimeColumn],[cachePurgeInterval],[cacheRetentionTime])

详情

为非持久化流表开启自动清理缓存。

通过以下两种方式之一来清理内存中的数据:

  • 配置 cacheSize 参数时,如果插入的数据使内存中流数据表的行数达到阈值,系统将清理内存中较旧的已发布记录。阈值确定方法如下:

    • 每次 append 的数据都不超过 cacheSize 时,阈值为 cacheSize 的 2.5 倍。

    • 当 append 的数据超过 cacheSize 时,阈值为追加行数和 cacheSize 之和的 1.2 倍。

  • 同时配置 cachePurgeTimeColumn, cachePurgeIntervalcacheRetentionTime,系统将根据时间列清理数据。每次插入新数据时,系统会计算新数据与内存中第一条数据的时间戳差值,当差值大于等于 cachePurgeInterval 时,系统仅保留时间戳与新数据时间戳差值小于等于 cacheRetentionTime 的数据,清理其它数据。

  • 按时间清理时,系统以单次批量写入的数据批次作为清理粒度,批次内的数据不会被拆分清理。因此,当单个批次覆盖的时间范围较长时,实际保留的数据可能超过 cacheRetentionTime 指定的时间范围。如果对内存上限有严格要求,建议使用 cacheSize,或缩小单个批次覆盖的时间范围。例如,设置 cachePurgeInterval=40mcacheRetentionTime=20m,假设 cachePurgeTimeColumn 中的数据按时间升序排列,并先后写入以下两个批次的数据:
    批次 时间范围 时间跨度
    批次 1 09:00–09:30 30 分钟
    批次 2 10:00–10:30 30 分钟

    批次 1 写入后

    • 内存中第一条数据的时间戳为 09:00;
    • 本次写入的新数据的最新时间戳为 09:30。

    系统计算:新数据的最新时间戳(09:30) - 内存中第一条数据的时间戳(09:00)= 30m;由于 30m < cachePurgeInterval(40m),本次写入没有达到清理触发条件,系统不执行清理。

    批次 2 写入后

    • 内存中第一条数据仍然是批次 1 的 09:00;
    • 本次写入的新数据的最新时间戳为 10:30。

    系统再次计算:新数据的最新时间戳(10:30) - 内存中第一条数据的时间戳(09:00);由于 90m >= cachePurgeInterval(40m),系统以本次新数据的最新时间戳 10:30 为基准,根据 cacheRetentionTime=20m 计算理论保留边界:10:30 - 20m = 10:10。按照时间条件,时间戳早于 10:10 的数据已经超出保留范围。因此:

    • 批次 1 的时间范围为 09:00–09:30,全部早于 10:10,被清理;
    • 批次 2 的时间范围为 10:00–10:30,其中 10:00–10:09 的数据也早于理论保留边界 10:10;但是,系统以单次批量写入的数据批次作为清理粒度,不会将本次写入的批次 2 拆分成“10:00–10:09”和“10:10–10:30”两部分。因此,批次 2 整体保留。

    清理结果如下:

    批次 与理论保留边界的关系 清理结果
    批次 1:09:00–09:30 整批位于保留范围之外 整批清理
    批次 2:10:00–10:30 批次跨越理论保留边界 整批保留

参数

table 是一个空的流数据表。

cacheSize 是一个正整数,可选参数,表示流数据表在内存中最多保留的记录数。

cachePurgeTimeColumn 字符串标量,需要指定为非持久化流表中的时间列名称。

cachePurgeInterval DURATION 类型标量,表示触发清理内存中数据的时间间隔。

cacheRetentionTime 可选参数,DURATION 类型标量,表示按时间清理时用于判断是否保留数据的时间范围。系统以单次批量写入的数据批次作为清理粒度,因此实际保留的数据可能超过该参数指定的时间范围。

返回值

无。

例子

例1. 配置 cacheSize 参数,根据内存中的数据量进行清理。

t = streamTable(1000:0, `time`sym`volume, [DATETIME, SYMBOL, INT])
enableTableCachePurge(table=t, cacheSize=1000)
time = datetime(2024.01.01T09:00:00) +1..1000*2
sym=take(`a`b`c, 1000)
volume = rand(10,1000)

insert into t values([time, sym, volume])
getStreamTableCacheOffset(t)
//0

time = datetime(2024.01.01T09:35:00) +1..1000*2
sym=take(`a`b`c, 1000)
volume = rand(10,1000)
insert into t values([time, sym, volume])
getStreamTableCacheOffset(t)
//500

例2. 配置 cachePurgeTimeColumn, cachePurgeIntervalcacheRetentionTime,根据时间列清理数据。

t = streamTable(1000:0, `time`sym`volume, [DATETIME, SYMBOL, INT])

enableTableCachePurge(table=t, cachePurgeTimeColumn=`time,
 cachePurgeInterval=30m, cacheRetentionTime=20m)

time = datetime(2024.01.01T09:00:00) +1..1000*2
sym=take(`a`b`c, 1000)
volume = rand(10,1000)

insert into t values([time, sym, volume])
getStreamTableCacheOffset(t)
//0

time = datetime(2024.01.01T09:35:00) +1..1000*2
sym=take(`a`b`c, 1000)
volume = rand(10,1000)
insert into t values([time, sym, volume])
getStreamTableCacheOffset(t)
//999