ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Hydra RQ Launcher 插件完全指南:基于 Redis Queue 的分布式任务调度与多节点并行执行

Hydra RQ Launcher 插件完全指南:基于 Redis Queue 的分布式任务调度与多节点并行执行 Hydra RQ Launcher 插件完全指南基于 Redis Queue 的分布式任务调度与多节点并行执行【免费下载链接】hydraHydra is a framework for elegantly configuring complex applications项目地址: https://gitcode.com/GitHub_Trending/hyd/hydra本指南围绕 Hydra 官方 RQ Launcher 插件hydra-rq-launcher展开系统讲解如何借助 Redis QueueRQ实现跨多节点、基于队列的多任务并行执行从安装、hydra/launcherrq的启用方式到全部配置参数入队策略、Redis 连接、轮询间隔的环境变量与 YAML 覆盖方法再到 worker 的启动、串行化原理与源码级执行链路。读完本文你将能够在自己的 Hydra 应用中把 multirun 任务批量投递到 Redis 队列并由任意多台机器上的 worker 分布式消费同时掌握常见配置陷阱如 cloudpickle 序列化器不一致的排查方法。一、插件定位什么场景需要 RQ LauncherRQ Launcher 是 Hydra 官方维护的一个启动器Launcher插件为 Hydra 提供基于Redis QueueRQ的分布式执行与任务排队能力。它把 Hydra 的 multirun多任务批量运行模式从单机串行/进程池并行扩展为投递到 Redis 队列、由独立 worker 进程消费的分布式模型因此可以跨多个计算节点并行执行任务。从源码结构看该插件由以下几个核心部分组成仓库路径plugins/hydra_rq_launcher/hydra_plugins/hydra_rq_launcher/rq_launcher.pyRQLauncher类实现 Hydra 的Launcher接口hydra_plugins/hydra_rq_launcher/_core.pylaunch 核心逻辑负责建连、入队、轮询与结果回收hydra_plugins/hydra_rq_launcher/config.py结构化配置RQLauncherConf、EnqueueConf、RedisConf并通过ConfigStore注册为hydra/launcherrqhydra_plugins/hydra_rq_launcher/serializer.py基于 cloudpickle 的 RQ 序列化器example/可直接运行的示例应用tests/test_rq_launcher.pyLauncher 标准测试套件与序列化相关回归测试。使用前提与限制原文档明确指出使用本插件必须有一个可访问的 Redis 服务器如果你的并行需求仅限于单机多进程那么 Joblib Launcher 更合适——它不需要任何外部数据库RQ 本身不支持 Windows其官方文档中的 Limitations 一节有说明因此该插件主要面向 Linux / macOS 环境。setup.py中声明的支持系统为 Operating System :: MacOS 与 Operating System :: POSIX :: LinuxPython 版本要求3.10。二、安装插件使用 pip 直接安装并升级到最新版本pip install hydra-rq-launcher --upgrade安装完成后插件会以命名空间包hydra_plugins.hydra_rq_launcher的形式注册到 Hydra 的插件发现机制中。setup.py中的find_namespace_packages(include[hydra_plugins.*])保证了这一点测试test_discovery()也专门验证了RQLauncher能被 Hydra 的插件子系统Plugins.instance().discover(Launcher)正确发现。从依赖声明见plugins/hydra_rq_launcher/setup.py看运行时依赖包括rq2.10.0,3Redis Queue 库本身cloudpickle任务序列化fakeredis2.36.2,3无 Redis 服务器时的同步模拟模式redis.mocktrue依赖hydra-core1.1.0.dev7Hydra 核心。三、快速开始启用 RQ Launcher方式一命令行参数启用在运行任意 Hydra 应用时追加参数python your_app.py hydra/launcherrq方式二配置文件默认启用在应用配置或独立的配置文件中通过 defaults 覆盖hydra/launcherdefaults: - override hydra/launcher: rq仓库自带示例plugins/hydra_rq_launcher/example/config.yaml就是这种方式defaults: - override hydra/launcher: rq task: 1 hydra: launcher: enqueue: job_timeout: 1d对应示例应用plugins/hydra_rq_launcher/example/my_app.py是一个最简 Hydra 应用读取cfg.task打印当前进程 ID 并休眠 1 秒模拟任务执行可直接用于体验插件效果。四、默认配置详解每一个参数的作用与取值启用后可通过 Hydra 的--cfg hydra查看 launcher 的完整默认配置python your_app.py hydra/launcherrq --cfg hydra -p hydra.launcher输出对应以下结构化配置与plugins/hydra_rq_launcher/hydra_plugins/hydra_rq_launcher/config.py中的定义一一对应# package hydra.launcher _target_: hydra_plugins.hydra_rq_launcher.rq_launcher.RQLauncher enqueue: job_timeout: null ttl: null result_ttl: null failure_ttl: null at_front: false job_id: null description: null queue: default redis: host: ${oc.env:REDIS_HOST,localhost} port: ${oc.env:REDIS_PORT,6379} db: ${oc.env:REDIS_DB,0} password: ${oc.env:REDIS_PASSWORD,null} ssl: ${oc.env:REDIS_SSL,False} ssl_ca_certs: ${oc.env:REDIS_SSL_CA_CERTS,null} mock: ${oc.env:REDIS_MOCK,False} stop_after_enqueue: false wait_polling: 1.04.1enqueue入队行为控制对应EnqueueConfenqueue子配置对应 RQ 的Queue.enqueue()关键字参数逐项说明如下注释语义来自config.py参数默认值含义job_timeoutnull任务最长运行时间超时会被终止。支持d/h/m/s时间单位如1d表示 1 天。null表示不限制源码中在入队前会将None归一化为-1传给 RQ见_core.pyttlnull任务在队列中允许排队的最长时间超时任务会被丢弃。同样支持d/h/m/s单位result_ttlnull成功任务及其结果在 Redis 中的保留时间。null时在源码中被归一化为-1无限期保留failure_ttlnull失败任务在 Redis 中的保留时间。null同样归一化为-1at_frontfalse是否把任务放到队首而不是队尾优先执行job_idnull任务 ID。null时由插件自动生成一个 UUID见_core.py中str(uuid.uuid4())并会写入hydra.job.id同时作为输出目录名的一部分descriptionnull任务描述。null时自动使用该任务的 overrides 字符串 .join(filter_overrides(overrides))实测建议job_timeout: 1d这样的写法已经在示例配置example/config.yaml中出现可作为时间单位语法的最直接参照。4.2queue队列名称默认队列名为default。所有任务都会投递到这个队列worker 侧需要消费同名队列。_core.py中通过Queue(namerq_cfg.queue, connectionconnection, ...)创建 RQ 队列对象运行时日志也会打印enqueuing N job(s) in queue : queue_name以方便核对。4.3redisRedis 连接配置对应RedisConf所有连接参数均通过oc.env插值从环境变量读取未设置时使用默认值配置项环境变量默认值含义hostREDIS_HOSTlocalhostRedis 服务器地址portREDIS_PORT6379Redis 端口dbREDIS_DB0Redis 数据库编号passwordREDIS_PASSWORDnull无密码连接密码sslREDIS_SSLFalse是否启用 SSL 连接ssl_ca_certsREDIS_SSL_CA_CERTSnullSSL 使用的 CA 证书路径mockREDIS_MOCKFalse是否以同步模拟模式运行无需真实 Redis仅用于测试_core.py中的连接逻辑印证了这些参数的使用方式当redis.mock为假默认异步模式时创建真实的redis.Redis(...)连接传入 host/port/db/password/ssl/ssl_ca_certs当redis.mock为真时改用FakeStrictRedis()以单线程同步方式运行日志会打印 Running in synchronous mode。测试套件tests/conftest.py正是通过monkeypatch.setenv(REDIS_MOCK, True)在无 Redis 服务的情况下跑完整套集成测试。4.4stop_after_enqueue与wait_pollingstop_after_enqueue: false默认不停止。若设为true插件在完成全部任务入队后抛出StopAfterEnqueue自定义异常立即退出_core.py中raise StopAfterEnqueue适用于只投递、不等待结果的场景例如把任务交给其他调度系统接管wait_polling: 1.0入队完成后主进程每隔多少秒轮询一次任务状态日志中会打印Polling job statuses every wait_polling sec。轮询逻辑见_core.py的while True循环不断检查所有 job 的get_status()直到全部处于finished或failed状态。五、通过环境变量配置 Redis 连接插件约定使用一组环境变量存放 Redis 连接信息。在bash/zsh等 shell 中可这样设置export REDIS_HOSTlocalhost export REDIS_PORT6379 export REDIS_DB0 export REDIS_PASSWORD其中REDIS_HOST、REDIS_PORT、REDIS_DB、REDIS_PASSWORD分别对应服务器的主机地址、端口、数据库编号与密码。启用 SSL 连接当 Redis 服务端开启 TLS 时设置以下环境变量export REDIS_SSLtrue export REDIS_SSL_CA_CERTS/etc/ssl/certs/ca-certificates.crtREDIS_SSL控制是否启用 SSLREDIS_SSL_CA_CERTS指向自定义 CA 证书路径对应 redis-py 的ssl_ca_certs参数。注意这些环境变量之所以生效是因为config.py中每个字段都通过II(oc.env:REDIS_HOST,localhost)这类 OmegaConf 插值绑定到环境变量——即使不改配置文件环境变量也会自动覆盖默认值。六、启动 RQ Worker 消费任务配置好环境变量后需要启动 RQ worker 才能真正消费队列中的任务。worker 通过 URL 方式连接 Redisrq worker --url redis://:$REDIS_PASSWORD$REDIS_HOST:$REDIS_PORT/$REDIS_DB几点必须注意的实践要点原文档明确强调worker 的 Python 环境中必须安装任务所需的全部依赖包括 Hydra 与hydra-rq-launcher本身任务序列化使用cloudpickle。worker 必须使用与入队侧一致的序列化器否则反序列化会失败。serializer.py中CloudpickleSerializer只是把cloudpickle.dumps/loads直接暴露给 RQ结合_core.py顶部的_SERIALIZER_ERROR提示如果 worker 未以 Hydra 的 cloudpickle 序列化器启动插件会抛出CompactHydraException错误信息中建议用如下方式显式指定序列化器启动 workerrq worker -S hydra_plugins.hydra_rq_launcher.serializer.CloudpickleSerializer相关回归测试见tests/test_rq_launcher.py中的test_unserializable_rq_result_error与test_unserializable_rq_result_stream_error它们分别模拟了 RQ 旧版rq:job:*hash 结果与新版本rq:results:*结果流两种反序列化失败路径均断言抛出包含 cloudpickle serializer 提示的CompactHydraException。七、运行示例一次完整的 multirun 分布式执行仓库在plugins/hydra_rq_launcher/example/下提供了配套示例。假设你已安装插件、启动 Redis 并拉起 worker运行python my_app.py --multirun task1,2,3,4,5会将 5 个任务投递到default队列交由 worker 实例处理。终端输出大致如下$ python my_app.py --multirun task1,2,3,4,5 [HYDRA] RQ Launcher is enqueuing 5 job(s) in queue : default [HYDRA] Sweep output dir : multirun/2020-06-15/18-00-00 [HYDRA] Enqueued 13b3da4e-03f7-4d16-9ca8-cfb3c48afeae [HYDRA] #1 : task1 [HYDRA] Enqueued 00c6a32d-e5a4-432c-a0f3-b9d4ef0dd585 [HYDRA] #2 : task2 [HYDRA] Enqueued 63b90f27-0711-4c95-8f63-70164fd850df [HYDRA] #3 : task3 [HYDRA] Enqueued b1d49825-8b28-4516-90ca-8106477e1eb1 [HYDRA] #4 : task4 [HYDRA] Enqueued ed96bdaa-087d-4c7f-9ecb-56daf948d5e2 [HYDRA] #5 : task5 [HYDRA] Finished enqueuing [HYDRA] Polling job statuses every 1.0 sec这段日志对应的源码逻辑位于_core.py的launch()先创建 sweep 输出目录hydra.sweep.dir根据redis.mock选择真实 Redis 连接或FakeStrictRedis并构造Queue(name..., connection..., is_async..., serializerCloudpickleSerializer)遍历每个 job override用hydra_context.config_loader.load_sweep_config()解析出该任务的独立配置再把自动生成的job_id与序号写入sweep_config.hydra.job.id / hydra.job.num调用queue.enqueue(execute_job, hydra_context..., sweep_config..., task_function..., singleton_state..., **enqueue_keywords)投递任务worker 侧实际执行的是execute_job函数全部入队后按wait_polling间隔轮询直至所有 job 结束再逐个取回JobReturn结果。execute_job_core.py完整复现了 Hydra 任务的标准执行环境setup_globals()初始化全局上下文、Singleton.set_state(singleton_state)恢复入队时捕获的单例状态、HydraConfig.instance().set_config(sweep_config)注入任务配置最后调用run_job(...)并把结果写入hydra.sweep.dir/hydra.sweep.subdir对应的输出目录。这意味着每个 worker 任务都拥有与本地运行一致的输出布局。八、测试验证与可靠性保障plugins/hydra_rq_launcher/tests/test_rq_launcher.py从多个层面验证了插件的正确性test_discovery验证插件可被 Hydra 插件系统发现TestRQLauncher(LauncherTestSuite)运行 Hydra 提供的Launcher 通用测试套件参数化launcher_namerq覆盖所有 Launcher 必须满足的行为契约TestRQLauncherIntegration(IntegrationTestSuite)通过-m hydra/launcherrq走完整集成测试test_example_app直接跑示例应用断言task1,2,3,4四个任务全部返回、且每个返回的 overrides 与预期一致还额外参数化覆盖了hydra.launcher.redis.ssltrue的配置分支test_cloudpickle_serializer验证 cloudpickle 序列化函数序列化再反序列化后调用value(41) 42两个test_unserializable_rq_result_*验证序列化器不匹配时的错误提示。这些测试默认在REDIS_MOCKTrue的模拟 Redis 上运行见tests/conftest.py说明即使没有 Redis 服务你也能在本地通过mock模式完整走通插件流程——这非常适合快速验证任务配置是否正确。九、延伸监控、常见模式与配置技巧任务监控RQ 提供基于控制台与 Web 界面的监控手段官方文档 job monitoring 与 RQ Dashboard可实时查看队列深度、任务状态与失败详情。多节点并行时worker 的启动数量、机器分布完全由你自行规划插件只负责入队与结果回收更多配置模式插件的配置覆盖除了上面的环境变量方式也支持在 YAML 中直接写hydra.launcher.enqueue.job_timeout: 1d见示例example/config.yaml或在命令行用hydra.launcher.queuemyqueue等方式覆盖任意字段。关于 Hydra 插件配置的标准做法可进一步参考 configuring_plugins.md输出目录相对路径提醒_core.py会在 sweep 目录为相对路径时打印警告——因为 worker 启动的工作目录可能与入队方不同生产环境建议使用绝对路径配置hydra.sweep.dir避免结果落盘位置不一致结果序列化一致性的本质CloudpickleSerializer的作用范围包括函数、局部闭包等标准 pickle 无法处理的 Python 对象测试中直接序列化了一个嵌套函数value这也是分布式任务必须选择 cloudpickle 而非默认 pickle 的原因。若 worker 与入队方版本不一致导致结果无法反序列化插件会给出明确的CompactHydraException指引而非静默失败。十、小结RQ Launcher 是 Hydra 生态中把 multirun 任务真正推上队列、跨节点分发的官方插件hydra/launcherrq一行即可切换启动器enqueue、queue、redis、wait_polling、stop_after_enqueue五组配置完整覆盖入队策略、目标队列、Redis 连接与等待行为配合环境变量、YAML 与命令行三种覆盖方式可灵活适配各种部署环境。理解其_core.py的入队—轮询—回收链路以及 cloudpickle 序列化约束是让分布式任务稳定运行的关键。若你的并行需求仅限单机优先考虑无需外部依赖的 Joblib Launcher一旦需要多节点队列调度RQ Launcher 便是开箱即用的选择。【免费下载链接】hydraHydra is a framework for elegantly configuring complex applications项目地址: https://gitcode.com/GitHub_Trending/hyd/hydra创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表