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

使用DolphinDB MapReduce进行OHLC降采样时周数据重复问题

DolphinDB MapReduce聚合跨月周线产生重复记录的原因及解决方法

问题原因

  • 你的DFS表按月份分区,sqlDS会按分区拆分数据源,每个分区独立执行Map函数calcOHLC1w。
  • 跨月的周(比如2020.01.27的周)会同时存在于两个相邻的月份分区中(1月和2月)。每个分区的Map函数在分组时,会各自基于本分区内的5分钟数据生成该周的一条聚合记录,然后直接append到目标表,最终导致同一周出现两条重复记录。
  • 而标准分布式SQL是全局分组,它会跨分区收集所有属于同一周的数据,统一聚合为一条记录,因此不会产生重复。

解决方法

核心思路:不在Map阶段直接写入目标表,而是将每个Map的聚合结果返回,在Reduce阶段合并相同open_time(周起始时间)+ticker的记录,最后统一写入表。

修正后的代码:

klines = loadTable("dfs://um", "k")
ds = sqlDS(<select open_time, ticker, freq, open, high, low, close, quote_volume, count, taker_buy_quote_volume from klines where freq="5">)

def calcOHLC1w(t){
    // 仅在Map阶段生成当前分区内的周线聚合结果,不写入表
    return select "w" as freq, 
           first(open) as open, 
           max(high) as high, 
           min(low) as low, 
           last(close) as close, 
           sum(quote_volume) as quote_vol, 
           sum(count) as count, 
           sum(taker_buy_quote_volume) as taker_buy_vol 
           from t group by weekBegin(open_time) as open_time, ticker
}

def reduceFunc(tmpTables){
    // Reduce阶段合并所有Map的结果,按周和ticker再次聚合去重
    combined = unionAll(tmpTables)
    finalResult = select first(open) as open, 
                  max(high) as high, 
                  min(low) as low, 
                  last(close) as close, 
                  sum(quote_vol) as quote_vol, 
                  sum(count) as count, 
                  sum(taker_buy_vol) as taker_buy_vol 
                  from combined group by open_time, ticker
    // 统一写入目标表
    loadTable("dfs://um", "k").append!(select "w" as freq, * from finalResult)
    return finalResult.size()
}

def finalFunc(size){
    print(size)
}

mr(ds, calcOHLC1w, reduceFunc, finalFunc, false)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.11 18:42:12