Flux中使用变量在filter函数实现动态过滤时的报错问题
解决Flux中基于分位数过滤记录的错误问题
问题场景
通过quantile函数计算出"PurchaseOrder"类型文档处理时间的95%分位数(结果为999),但尝试过滤出处理时间低于该分位数的记录时触发错误:
error @116:41-116:51: expected [{A with pt: B}] (array) but found stream[{A with pt: B}]
分位数计算代码
percentile = totalTimeByDoc |> filter(fn: (r) => r["documentType"] == "PurchaseOrder") |> group(columns:["documentType"]) |> quantile(column: "processTime", q: 0.95, method: "estimate_tdigest", compression: 9999.0) |> limit(n: 1) |> rename(columns: {processTime: "pt"})
报错的过滤代码
percentile_filered = totalTimeByDoc |> filter(fn: (r) => r["documentType"] == "PurchaseOrder") |> filter(fn: (r) => r.processTime < percentile[0]["pt"]) |> yield()
示例数据(totalTimeByDoc)
| 序号 | documentType | 文档ID | processTime |
|---|---|---|---|
| 0 | PurchaseOrder | testpass22PID230207222747-1 | 1200 |
| 1 | PurchaseOrder | testpass22PID230207222747-2 | 807 |
| 2 | PurchaseOrder | testpass22PID230207222934-1 | 671 |
| 3 | PurchaseOrder | testpass22PID230207222934-2 | 670 |
错误原因
percentile变量存储的是Flux的**流(stream)**对象,而非内存中的数组。Flux中大多数数据处理操作返回的都是流,无法直接通过数组索引[0]["pt"]的方式访问流中的元素。
解决方案
需要将分位数流转换为可直接引用的数值,推荐两种实现方式:
方法1:用findRecord提取标量值
通过findRecord从分位数流中取出唯一记录,提取pt字段作为标量值,再用于过滤:
// 提取分位数标量值 percentileValue = totalTimeByDoc |> filter(fn: (r) => r["documentType"] == "PurchaseOrder") |> group(columns:["documentType"]) |> quantile(column: "processTime", q: 0.95, method: "estimate_tdigest", compression: 9999.0) |> limit(n: 1) |> findRecord(fn: (key) => true, idx: 0) |> get(key: "pt") // 基于标量值过滤数据 percentile_filtered = totalTimeByDoc |> filter(fn: (r) => r["documentType"] == "PurchaseOrder") |> filter(fn: (r) => r.processTime < percentileValue) |> yield()
方法2:用join关联流过滤
将原始数据流和分位数流通过documentType关联,再基于关联后的字段进行过滤:
// 保留分位数流 percentileStream = totalTimeByDoc |> filter(fn: (r) => r["documentType"] == "PurchaseOrder") |> group(columns:["documentType"]) |> quantile(column: "processTime", q: 0.95, method: "estimate_tdigest", compression: 9999.0) |> limit(n: 1) |> rename(columns: {processTime: "pt"}) |> group() // 取消分组,方便后续关联 // 关联后过滤 percentile_filtered = totalTimeByDoc |> filter(fn: (r) => r["documentType"] == "PurchaseOrder") |> join( tables: {raw: _, percentile: percentileStream}, on: ["documentType"] ) |> filter(fn: (r) => r.raw_processTime < r.percentile_pt) |> yield()
内容的提问来源于stack exchange,提问作者Rahul Bhardwaj
相关产品推荐
相关产品推荐

