Skip to contents

RushWorker inherits all methods from Rush. Upon initialization, the worker registers itself in the Redis database as a running worker. This class is usually not constructed directly by the user.

In addition to the inherited methods, the worker provides methods that require a worker identity:

  • $pop_task(): Pop a task from the queue and mark it as running.

  • $push_running_tasks(xss): Create running tasks evaluated by the worker.

  • $finish_tasks(keys, yss): Save the output of tasks and mark them as finished.

  • $fail_tasks(keys, conditions): Mark tasks as failed and optionally save the condition objects.

  • $n_queued_available_tasks: Number of queued tasks the worker can pop.

Value

Object of class R6::R6Class and RushWorker.

Super class

Rush -> RushWorker

Public fields

worker_id

(character(1))
Identifier of the worker.

profile

(character(1))
Name of the mirai compute profile the worker runs on. NULL if the worker runs on the default compute profile.

heartbeat

(callr::r_bg)
Background process for the heartbeat.

Active bindings

terminated

(logical(1))
Whether to shutdown the worker. Used in the worker loop to determine whether to continue.

n_queued_available_tasks

(integer(1))
Number of queued tasks the worker can pop, i.e. the tasks in the shared queue and in the queue of the compute profile the worker runs on. Tasks queued for other compute profiles are not counted.

Methods

Inherited methods


RushWorker$new()

Creates a new instance of this R6 class.

Usage

RushWorker$new(
  network_id,
  config = NULL,
  worker_id = NULL,
  profile = NULL,
  heartbeat_period = NULL,
  heartbeat_expire = NULL
)

Arguments

network_id

(character(1))
Identifier of the rush network. Manager and workers must have the same id. Keys in Redis are prefixed with the instance id.

config

(redux::redis_config)
Redis configuration options. If NULL, configuration set by rush_plan() is used. If rush_plan() has not been called, the REDIS_URL environment variable is parsed. If REDIS_URL is not set, a default configuration is used. See redux::redis_config for details.

worker_id

(character(1))
Identifier of the worker. Keys in redis specific to the worker are prefixed with the worker id.

profile

(character(1))
Name of the mirai compute profile the worker runs on. If NULL, the worker runs on the default compute profile.

heartbeat_period

(integer(1))
Period of the heartbeat in seconds. The heartbeat is updated every heartbeat_period seconds. Must be at least 1 second.

heartbeat_expire

(integer(1))
Time to live of the heartbeat in seconds. The heartbeat key is set to expire after heartbeat_expire seconds. Must be at least heartbeat_period, otherwise a live worker is reaped as lost between two heartbeats. Set it larger than the longest pause a worker may experience, for example from garbage collection or swapping, because a live worker wrongly declared lost can leave a task in an inconsistent state.


RushWorker$pop_task()

Pop a task from the queue and mark it as running. Returns NULL if no task is available.

A worker running on a compute profile takes tasks from the queue of its profile first and falls back to the shared queue, so that tasks pushed without a profile are processed by any worker. A worker running on the default compute profile only takes tasks from the shared queue.

Usage

RushWorker$pop_task(timeout = 1, fields = "xs")

Arguments

timeout

(numeric(1))
Time to wait for task in seconds.

fields

(character())
Fields to be returned.


RushWorker$push_running_tasks()

Create running tasks.

Usage

RushWorker$push_running_tasks(xss, xss_extra = NULL, extra = NULL)

Arguments

xss

(list of named list())
Lists of arguments for the function e.g. list(list(x1 = 1, x2 = 2), list(x1 = 3, x2 = 4)). If xss is empty, no tasks are created and the method returns an empty character().

xss_extra

(list of named list())
List of additional information stored along with the task e.g. list(list(timestamp_xs = Sys.time()), list(timestamp_xs = Sys.time())).

extra

(list)
Deprecated argument for additional information stored along with the task. Use xss_extra instead.

Returns

(character())
Keys of the tasks.


RushWorker$finish_tasks()

Save the output of tasks and mark them as finished.

Usage

RushWorker$finish_tasks(keys, yss, yss_extra = NULL, extra = NULL)

Arguments

keys

(character(1))
Keys of the associated tasks.

yss

(list of named list())
Lists of results for the function e.g. list(list(y1 = 1, y2 = 2), list(y1 = 3, y2 = 4)).

yss_extra

(list of named list())
List of additional information stored along with the results e.g. list(list(timestamp_ys = Sys.time()), list(timestamp_ys = Sys.time())).

extra

(named list())
Deprecated argument for additional information stored along with the results. Use yss_extra instead.

Returns

(RushWorker)
Invisible self.


RushWorker$fail_tasks()

Move running tasks to failed and optionally save the condition objects.

Usage

RushWorker$fail_tasks(keys, conditions = NULL)

Arguments

keys

(character())
Keys of the running tasks to be moved.

conditions

(list())
List conditions e.g. list(simpleError("Error"), simpleError("Error")). Defaults to list(message = "Task failed").

Returns

(RushWorker)
Invisible self.


RushWorker$set_terminated()

Mark the worker as terminated. Last step in the worker loop before the worker terminates.

Usage

RushWorker$set_terminated()

Returns

(RushWorker)
Invisible self.


RushWorker$clone()

The objects of this class are cloneable with this method.

Usage

RushWorker$clone(deep = FALSE)

Arguments

deep

Whether to make a deep clone.