使用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
相关产品推荐
相关产品推荐

