You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在DolphinDB中实现实时传感器数据流的时间邻近匹配?

基于DolphinDB实现多频率传感器数据流的时间邻近匹配

问题描述

工厂内一台电机配备温度(每500ms采集一次)和振动(每200ms采集一次)两个传感器,需实时关联两个数据流计算综合健康评分。使用EquiJoinEngine基于设备ID和时间戳精确匹配时,仅时间戳完全一致的数据能成功匹配,大部分数据因时间戳不对应丢失,导致输出结果极为稀疏。

原测试代码:

baseTime = 2026.03.11T14:00:00.000

vibRawData = table(
    (baseTime + (0..49) * 200) as ts,      
    take(`motor_01, 50) as device_id,
    3.0 + rand(2.0, 50) as vibration
)

tempRawData = table(
    (baseTime + (0..19) * 500) as ts,       
    take(`motor_01, 20) as device_id,
    60.0 + rand(10.0, 20) as temperature
)

share streamTable(1:0,
    `ts`device_id`temperature,
    [TIMESTAMP, SYMBOL, DOUBLE]
) as tempForEQ

share streamTable(1:0,
    `ts`device_id`vibration,
    [TIMESTAMP, SYMBOL, DOUBLE]
) as vibForEQ

share streamTable(1:0,
    `ts`device_id`temperature`vibration`health_score,
    [TIMESTAMP, SYMBOL, DOUBLE, DOUBLE, DOUBLE]
) as eqResult

eqTempVib = createEquiJoinEngine(
    name="eqTempVib",
    leftTable=tempForEQ,
    rightTable=vibForEQ,
    outputTable=eqResult,
    metrics=<[
        tempForEQ.temperature,
        vibForEQ.vibration,
        100.0 - tempForEQ.temperature * 0.5 - vibForEQ.vibration * 10.0
    ]>,
    matchingColumn=`device_id,
    timeColumn=`ts               
)

subscribeTable(tableName="tempForEQ", actionName="eqLeft",
    offset=0, handler=getLeftStream(eqTempVib), msgAsTable=true, hash=0)
subscribeTable(tableName="vibForEQ", actionName="eqRight",
    offset=0, handler=getRightStream(eqTempVib), msgAsTable=true, hash=1)

vibForEQ.append!(vibRawData)
sleep(500)
tempForEQ.append!(tempRawData)

select * from eqResult order by ts

解决方案

EquiJoinEngine仅支持精确时间戳匹配,无法处理多频率数据流的时间邻近关联需求。推荐以下两种适配场景的实现方式:

方案1:使用LookupJoinEngine匹配最近邻数据

LookupJoinEngine支持在指定时间范围内,为左表每条数据匹配右表中时间最接近的记录,完美适配多频率传感器数据的实时关联。

修改后的代码:

baseTime = 2026.03.11T14:00:00.000

vibRawData = table(
    (baseTime + (0..49) * 200) as ts,      
    take(`motor_01, 50) as device_id,
    3.0 + rand(2.0, 50) as vibration
)

tempRawData = table(
    (baseTime + (0..19) * 500) as ts,       
    take(`motor_01, 20) as device_id,
    60.0 + rand(10.0, 20) as temperature
)

share streamTable(1:0,
    `ts`device_id`temperature,
    [TIMESTAMP, SYMBOL, DOUBLE]
) as tempForLookup

share streamTable(1:0,
    `ts`device_id`vibration,
    [TIMESTAMP, SYMBOL, DOUBLE]
) as vibForLookup

share streamTable(1:0,
    `ts`device_id`temperature`vibration`health_score,
    [TIMESTAMP, SYMBOL, DOUBLE, DOUBLE, DOUBLE]
) as lookupResult

// 创建LookupJoinEngine,设置时间范围为温度数据前后300ms,匹配最近的振动数据
lookupTempVib = createLookupJoinEngine(
    name="lookupTempVib",
    leftTable=tempForLookup,
    rightTable=vibForLookup,
    outputTable=lookupResult,
    metrics=<[
        tempForLookup.temperature,
        vibForLookup.vibration,
        100.0 - tempForLookup.temperature * 0.5 - vibForLookup.vibration * 10.0
    ]>,
    matchingColumn=`device_id,
    timeColumn=`ts,
    timeRange=[-300, 300],  // 匹配温度数据前后300ms内的振动数据
    closest=true  // 选择时间最接近的记录
)

subscribeTable(tableName="tempForLookup", actionName="lookupLeft",
    offset=0, handler=getLeftStream(lookupTempVib), msgAsTable=true, hash=0)
subscribeTable(tableName="vibForLookup", actionName="lookupRight",
    offset=0, handler=getRightStream(lookupTempVib), msgAsTable=true, hash=1)

vibForLookup.append!(vibRawData)
sleep(500)
tempForLookup.append!(tempRawData)

select * from lookupResult order by ts

方案2:滑动窗口聚合后关联

若需要对高频振动数据做预处理(如取窗口平均值),可先通过TimeSeriesEngine将振动数据按温度采集频率(500ms)做滑动窗口聚合,再用EquiJoinEngine关联温度数据。

示例代码:

baseTime = 2026.03.11T14:00:00.000

vibRawData = table(
    (baseTime + (0..49) * 200) as ts,      
    take(`motor_01, 50) as device_id,
    3.0 + rand(2.0, 50) as vibration
)

tempRawData = table(
    (baseTime + (0..19) * 500) as ts,       
    take(`motor_01, 20) as device_id,
    60.0 + rand(10.0, 20) as temperature
)

// 振动数据滑动窗口聚合:500ms窗口,与温度采集频率对齐
share streamTable(1:0, `ts`device_id`vibration_avg, [TIMESTAMP, SYMBOL, DOUBLE]) as vibAgg

vibAggEngine = createTimeSeriesEngine(
    name="vibAggEngine",
    windowSize=500,
    step=500,
    metrics=<[avg(vibration)]>,
    dummyTable=vibRawData,
    outputTable=vibAgg,
    timeColumn=`ts,
    keyColumn=`device_id
)

subscribeTable(tableName="vibRawData", actionName="vibAgg",
    offset=0, handler=getStream(vibAggEngine), msgAsTable=true, hash=0)

// 关联温度数据与聚合后的振动数据
share streamTable(1:0, `ts`device_id`temperature`vibration_avg`health_score, [TIMESTAMP, SYMBOL, DOUBLE, DOUBLE, DOUBLE]) as windowJoinResult

joinEngine = createEquiJoinEngine(
    name="joinEngine",
    leftTable=tempRawData,
    rightTable=vibAgg,
    outputTable=windowJoinResult,
    metrics=<[
        tempRawData.temperature,
        vibAgg.vibration_avg,
        100.0 - tempRawData.temperature * 0.5 - vibAgg.vibration_avg * 10.0
    ]>,
    matchingColumn=`device_id,
    timeColumn=`ts
)

subscribeTable(tableName="tempRawData", actionName="joinLeft",
    offset=0, handler=getLeftStream(joinEngine), msgAsTable=true, hash=0)
subscribeTable(tableName="vibAgg", actionName="joinRight",
    offset=0, handler=getRightStream(joinEngine), msgAsTable=true, hash=1)

vibRawData.append!(vibRawData)
sleep(500)
tempRawData.append!(tempRawData)

select * from windowJoinResult order by ts

内容的提问来源于stack exchange,提问作者Lambert

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.11 10:37:06