【问题标题】:R - doRedis - Overwrite getTask to control the order of execution in parallel foreach loopsR - doRedis - 覆盖 getTask 以控制并行 foreach 循环中的执行顺序
【发布时间】:2017-01-16 09:20:47
【问题描述】:

问题:我需要控制由 foreach 循环并行处理任务的执行顺序。不幸的是,foreach 不支持这一点。

考虑的解决方案:使用 doRedis 使用数据库来保存在 foreach 循环中执行的所有任务。为了控制顺序,我想通过 setGetTask 覆盖 getTask 以根据预先指定的顺序获取任务。虽然我找不到太多关于如何做到这一点的文档。

其他信息:

  1. setGetTask 上有一小段,在redis documentation 中有一个例子。

    getTask <- function ( queue , job_id , ...)
    {
    
      key <- sprintf("
      redisEval("local x=redis.call('hkeys',KEYS[1])[1];
                   if x==nil then return nil end;
                   local ans=redis.call('hget',KEYS[1],x);
                   redis.call('hdel',KEYS[1],x);i
                   return ans",key)
    }
    
    setGetTask(getTask)
    

    虽然我认为文档中的代码在语法上不正确(缺少恕我直言“和右括号”)“)。我认为这在 CRAN 上是不可能的,因为文档的代码是在提交时执行的。

  2. 更改 getTask 函数不会改变任何关于工人获取任务的任何内容(即使在 redisEval 中引入明显的无意义,例如将其更改为 redisEval("dddddddddd((("))

  3. 从源代码安装包后,我只能访问 setGetTask 函数(我从 official CRAN package page of version 1.1.1 下载的(恕我直言,这与直接从 CRAN 安装没有区别)

数据:要执行的任务的数据框如下所示:

taskName;taskQueuePosition;parameter1;paramterN
taskT;1;val1;10
taskK;2;val2;8
taskP;3;val3;7
taskA;4;val4;7

我想用'taskQueuePosition'来控制顺序,编号小的任务应该先执行。

问题:

  1. 是否有人知道我可以从哪里获得有关使用 doRedis 或 setGetTask 执行此操作的更多信息的任何来源?
  2. 有人知道我需要如何更改 getTask 以实现上述目的吗?
  3. 还有其他聪明的想法来控制 foreach 循环中的执行顺序吗?最好是这样我可以在某些时候使用 doRedis 作为并行后端(更改这将意味着由于复杂的技术基础架构原因而对处理进行重大更改)。

代码(便于复制):

以下假设redis-server在本地机器上启动。

Redis 数据库填充:

library(doRedis)
library(foreach)

options('redis:num'=TRUE) # needed for proper execution

REDIS_JOB_QUEUE = "jobs"
registerDoRedis(REDIS_JOB_QUEUE)

# filling up the data frame
taskDF = data.frame(taskName=c("taskT","taskK","taskP","taskA"),
           taskQueuePosition=c(1,2,3,4),
           parameter1=c("val1","val2","val3","val4"),
           parameterN=c(10,8,7,7))

foreach(currTask=iter(taskDF, by='row'), 
        .verbose = T
) %dopar% {
  print(paste("Executing task: ",currTask$taskName))
  Sys.sleep(currTask$parameterN)
}

removeQueue(REDIS_JOB_QUEUE)

工人:

library(doRedis)
REDIS_JOB_QUEUE = "jobs"

startLocalWorkers(n=1, queue=REDIS_JOB_QUEUE)

【问题讨论】:

    标签: r foreach parallel-processing r-doredis


    【解决方案1】:

    我可以解决问题,现在可以控制任务执行的顺序了。

    其他信息:

    1. 文档中似乎存在拼写错误,导致 getTask 示例无法正常工作。通过考虑包中文件 task.R 中的 default_getTask 函数的形式,它应该看起来可能类似于:

    getTaskDefault <- function ( queue , job_id , ...)
    {
      key <- sprintf("%s:%s",queue, job_id)
      return(redisEval("local x=redis.call('hkeys',KEYS[1])[1];
                       if x==nil then return nil end;
                       local ans=redis.call('hget',KEYS[1],x);
                       redis.call('set', KEYS[1] .. '.start.' .. x, x);
                       redis.call('hdel',KEYS[1],x);
                       return ans",key))
    }
    

    函数第一行第一个百分号后面的字母似乎丢失了。这可以解释奇数的括号和引号。

    2) setGetTask 对我仍然没有任何效果。当我在填充数据库时通过.option 设置getTask 函数时(如vignette of the package 中所述),它被成功调用。

    3) 2)的信息表示我不需要getTask函数,所以我可以使用CRAN的包。

    -- 问题 -----

    1) doRedis 小插图描述了如何成功设置自定义 getTask。

    2 和 3) 当 getTask 函数中的 LUA 脚本进行如下修改时,任务以提交的方式从数据库中提取。这不正是我所要求的,但由于时间限制以及我(或更好地)不是关于 LUA 脚本的第一个想法的事实,它是一个令人满意的解决方案来控制 taskQueuePosition 列的提交顺序。

    getTaskInOrder <- function ( queue , job_id , ...)
    {
    
      key <- sprintf("%s:%s",queue, job_id)
      return(redisEval("
    
            local tasks=redis.call('hkeys',KEYS[1]); -- get all tasks
    
            local x=tasks[1];           -- get first task available task
            if x==nil then              -- if there are no tasks left, stop processing
              return nil 
            end;  
    
            local xMin = 65535;         -- if we have more tasks than 65535, getting the 
            -- task with the lowest taskID is not guaranteed to be the first one
            local i = 1;
            -- local iMinFound = -1;
            while (x ~= nil) do         -- search the array until there are no tasks left
            -- print('x: ',x)
            local xNum = tonumber(x);
            if(xNum<xMin) then
              xMin = xNum;
              -- iMinFound = i;
            end
            i=i+1;
            -- print('i is now: ',i);
            x=tasks[i];
            end
            -- print('Minimum is task number',xMin,' found at i ', iMinFound)
            x=tostring(xMin)            -- convert it back to a string (maybe it would 
                                        -- be better to keep the original string somewhere, 
                                        -- in case we loose some information whilst converting to number)
    
            -- print('x is now:',x);
            -- print(KEYS[1] .. '.start.' .. x, x);
            -- print('');
            local ans=redis.call('hget',KEYS[1],x);
            redis.call('set', KEYS[1] .. '.start.' .. x, x);
            redis.call('hdel',KEYS[1],x);
            return ans",key))
    }
    

    重要提示:我注意到,如果一个任务被中止,订单就会搞砸并且重新提交的任务(即使任务编号保持不变),将在最初提交的任务之后执行。这对我来说没问题。

    ----- 代码(便于复制):------

    这导致以下代码示例(任务数据框中有 12 个条目,而不是原来的 4 个):

    Redis 数据库填充:

    library(doRedis)
    library(foreach)
    
    options('redis:num'=TRUE) # needed for proper execution
    
    REDIS_JOB_QUEUE = "jobs"
    
    getTaskInOrder <- function ( queue , job_id , ...)
    {
      ...like above
    }
    
    registerDoRedis(REDIS_JOB_QUEUE)
    
    # filling up the data frame already in order of tasks to be executed
    # otherwise the dataframe has to be sorted by taskQueuePosition
    taskDF = data.frame(taskName=c("taskA","taskB","taskC","taskD","taskE","taskF","taskG","taskH","taskI","taskJ","taskK","taskL"),
           taskQueuePosition=c(1,2,3,4,5,6,7,8,9,10,11,12),
           parameter1=c("val1","val2","val3","val4","val1","val2","val3","val4","val1","val2","val3","val4"),
           parameterN=c(5,5,5,4,4,4,4,3,3,3,2,2))
    
    foreach(currTask=iter(taskDF, by='row'), 
            .verbose = T,
            .options.redis = list(getTask = getTaskInOrder
    ) %dopar% {
      print(paste("Executing task: ",currTask$taskName))
      Sys.sleep(currTask$parameterN)
    }
    
    removeQueue(REDIS_JOB_QUEUE)
    

    工人:

    library(doRedis)
    REDIS_JOB_QUEUE = "jobs"
    
    startLocalWorkers(n=1, queue=REDIS_JOB_QUEUE)
    

    另一个注意事项:以防万一您像我一样处理长时间的工作,请注意a bug in redis 1.1.1(CRAN 上的当前版本),这会导致重新提交任务(由于超时) 尽管工人们仍在为它们工作。

    【讨论】:

      猜你喜欢
      • 2020-08-14
      • 2016-01-03
      • 1970-01-01
      • 2023-03-29
      • 1970-01-01
      • 2017-12-06
      • 1970-01-01
      • 2011-10-15
      • 1970-01-01
      相关资源
      最近更新 更多