enableTableShareAndCachePurge
语法
enableTableShareAndCachePurge(table, tableName,
[cacheSize],[cachePurgeTimeColumn],[cachePurgeInterval],[cacheRetentionTime])
详情
将非持久化的流数据表共享,并设置定时清理。
通过以下两种方式之一来清理内存中的数据:
-
配置 cacheSize 参数时,如果插入的数据使内存中流数据表的行数达到阈值,系统将清理内存中较旧的已发布记录。阈值确定方法如下:
-
每次 append 的数据都不超过 cacheSize 时,阈值为 cacheSize 的 2.5 倍。
-
当 append 的数据超过 cacheSize 时,阈值为追加行数和 cacheSize 之和的 1.2 倍。
-
-
同时配置 cachePurgeTimeColumn, cachePurgeInterval 和 cacheRetentionTime,系统将根据时间列清理数据。每次插入新数据时,系统会计算新数据与内存中第一条数据的时间戳差值,当差值大于等于 cachePurgeInterval 时,系统仅保留时间戳与新数据时间戳差值小于等于 cacheRetentionTime 的数据,清理其它数据。
- 按时间清理时,系统以单次批量写入的数据批次作为清理粒度,批次内的数据不会被拆分清理。因此,当单个批次覆盖的时间范围较长时,实际保留的数据可能超过
cacheRetentionTime 指定的时间范围。如果对内存上限有严格要求,建议使用
cacheSize,或缩小单个批次覆盖的时间范围。例如,设置
cachePurgeInterval=40m、cacheRetentionTime=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 是一个空的流数据表。
tableName 是一个字符串,表示 table 共享后的名称。
cacheSize 是一个正整数,可选参数,表示流数据表在内存中最多保留的记录数。
cachePurgeTimeColumn 字符串标量,可选参数。需要指定为持久化流表中的时间列名称。
cachePurgeInterval DURATION 类型标量,表示触发清理内存中数据的时间间隔。
cacheRetentionTime 可选参数,DURATION 类型标量,表示按时间清理时用于判断是否保留数据的时间范围。系统以单次批量写入的数据批次作为清理粒度,因此实际保留的数据可能超过该参数指定的时间范围。
返回值
无。
例子
例1. 配置 cacheSize 参数,根据内存中的数据量进行清理。
t = streamTable(1000:0, `time`sym`volume, [DATETIME, SYMBOL, INT])
enableTableShareAndCachePurge(table=t, tableName=`st, 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, cachePurgeInterval 和 cacheRetentionTime,根据时间列清理数据。
t = streamTable(1000:0, `time`sym`volume, [DATETIME, SYMBOL, INT])
enableTableShareAndCachePurge(table=t, tableName=`st, 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
