使用mr函数写入不同分区规则的分区表时出现ChunkInTransaction错误
解决DolphinDB分布式计算写入目标表分区冲突问题
错误原因分析
你的代码存在两个核心问题:
- map阶段错误加载全量表:在
myMapFunc中调用loadTable("dfs://k_minute","k_minute")会导致每个map任务都加载整个源表,完全失去分布式计算的意义,且极大浪费资源。 - 多任务并发写入冲突:源表按「日期+股票代码」分区,同一个日期下的多个股票分区会触发多个map任务,这些任务同时向目标表的同一日期分区写入数据,触发分布式锁冲突,导致报错。
修正方案
方案1:map阶段仅计算返回结果,统一合并后写入
修改mapFunc仅处理当前分区数据并返回计算结果,通过mr合并所有分区的计算结果后,一次性写入目标表,避免并发冲突:
def myMapFunc(table){ // 仅对当前map任务分配的分区数据做计算 return select `dayReturnSkew as factorname, skew(ratios(close)) as val from table group by date(tradetime) as tradetime, securityid } // 构建源表的分区数据集 ds = sqlDS(<select * from loadTable("dfs://k_minute","k_minute")>) // 执行分布式计算,合并所有map结果 factorResult = mr(ds, myMapFunc) // 统一写入目标表 loadTable("dfs://K_FACTOR_VERTICAL","factor_k").append!(factorResult)
方案2:利用mr的outputTable参数自动写入
如果数据量极大,合并结果内存无法承载,可以直接指定mr的outputTable参数,让系统自动协调分区写入,避免冲突:
def myMapFunc(table){ return select `dayReturnSkew as factorname, skew(ratios(close)) as val from table group by date(tradetime) as tradetime, securityid } // 加载目标表 targetTable = loadTable("dfs://K_FACTOR_VERTICAL","factor_k") // 构建源表分区数据集 ds = sqlDS(<select * from loadTable("dfs://k_minute","k_minute")>) // 执行分布式计算并直接写入目标表 mr(ds, myMapFunc, outputTable=targetTable)
关键说明
- 必须使用map任务传入的
table参数处理当前分区数据,禁止在map阶段加载全量表,否则分布式计算毫无意义。 - 统一写入或通过
outputTable写入时,系统会以单任务或协调后的多任务方式写入目标分区,不会触发并发锁冲突。
内容的提问来源于stack exchange,提问作者Jane
相关产品推荐
相关产品推荐

