分布式环境下基于 Whenever 的动态扩缩容配置方案

在分布式环境中实现 Ruby 定时任务的动态扩缩容,需解决两个核心问题:

  1. 任务去重:防止多个节点重复执行相同任务
  2. 动态感知:实时响应节点数量变化
解决方案架构
graph TD
    A[主节点] -->|广播| B[节点1]
    A -->|广播| C[节点2]
    A -->|广播| D[节点N]
    E[Redis] -->|存储| F[节点状态]
    G[Kubernetes] -->|扩缩容事件| H[Whenever 配置更新]

关键技术实现
  1. 动态配置层(Whenever 扩展)
# config/dynamic_schedule.rb
DynamicSchedule = lambda do
  # 从中央配置服务获取当前节点数
  node_count = ConfigService.fetch_node_count
  
  every(1.day, at: '02:00') do
    runner "DataCleanJob.perform_later", 
           shard_id: ENV['NODE_INDEX'] % node_count
  end
end

Whenever.recipes << DynamicSchedule

  1. 节点协调服务
# lib/shard_allocator.rb
class ShardAllocator
  REDIS = Redis.new(url: ENV['REDIS_URL'])
  
  def self.assign_shard
    node_id = SecureRandom.uuid
    REDIS.sadd("active_nodes", node_id)
    
    # 计算当前节点分片ID
    active_nodes = REDIS.smembers("active_nodes").sort
    active_nodes.index(node_id)
  end
end

  1. Kubernetes 集成(HPA 触发时)
# 扩缩容事件处理脚本
kubectl get pods -l app=worker | wc -l > current_node_count
rails runner 'ConfigService.update_node_count(File.read("current_node_count"))'

  1. 任务执行器(带分片校验)
# app/jobs/data_clean_job.rb
class DataCleanJob < ApplicationJob
  def perform
    current_shard = ENV['NODE_SHARD_ID'].to_i
    return unless current_shard == today_shard_id
    
    # 实际业务逻辑
    User.inactive.delete_all
  end

  private
  
  def today_shard_id
    Date.today.jd % ConfigService.node_count
  end
end

扩缩容工作流程
  1. 扩容时

    • Kubernetes 创建新 Pod
    • Pod 启动时通过 ShardAllocator 注册节点
    • Whenever 重新计算任务分片
  2. 缩容时

    • Kubernetes 发送 TERM 信号
    • 节点从 Redis 移除注册信息
    • 存活节点重新分配分片
关键配置项
# .env.production
REDIS_URL=redis://cluster-redis:6379/1
NODE_INDEX=<由初始化脚本自动注入>
NODE_SHARD_ID=<由ShardAllocator生成>

性能优化建议
  1. 使用 Redis 管道批量处理节点状态更新
  2. 为分片计算添加本地缓存(TTL 5分钟)
  3. 设置任务执行超时:
    every 1.hour do
      runner "TimeoutJob.perform_now", timeout: 55.minutes
    end
    

容错机制
# 节点心跳检测
every 5.minutes do
  runner "NodeRegistry.heartbeat"
end

# lib/node_registry.rb
def self.heartbeat
  REDIS.expire("node:#{ENV['NODE_ID']}", 10.minutes)
  REDIS.zadd("active_nodes", Time.now.to_i, ENV['NODE_ID'])
end

此方案特点:

  • 无中心节点:通过分片算法实现自协调
  • 零配置变更:扩缩容自动生效
  • 优雅退化:Redis 故障时自动切换为全节点执行
  • 资源优化:空闲节点自动跳过任务执行

:实际部署时需配合 Kubernetes Readiness Probe 确保节点注册完成后再接收任务。

Logo

开源鸿蒙跨平台开发社区汇聚开发者与厂商,共建“一次开发,多端部署”的开源生态,致力于降低跨端开发门槛,推动万物智联创新。

更多推荐