You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何通过Python AWS Lambda终止Apache Druid Coordinator Leader及识别重平衡卡顿

解决Apache Druid Segment重平衡卡顿相关问题

1. 通过AWS Lambda终止Druid Coordinator Leader节点

Druid官方未提供直接终止Coordinator Leader的API,需结合AWS资源管理工具实现,具体步骤如下:

  • 步骤1:关联Leader节点与AWS资源标识
    调用你已掌握的GET /druid/coordinator/v1/leader接口,返回结果包含Leader节点的私有IP或主机名。将该标识映射到对应AWS EC2实例ID,可通过实例的私有IP、私有DNS名称或自定义标签(如druid_role=coordinator)做匹配。

  • 步骤2:用boto3实现实例终止/重启
    在AWS Lambda中使用Python的boto3库操作EC2实例。若Druid节点属于Auto Scaling Group,终止实例后ASG会自动拉起新节点,满足你通过ZooKeeper重启Leader的需求。

    示例代码:

    import boto3
    import requests
    
    def lambda_handler(event, context):
        # 获取Coordinator Leader信息
        druid_coordinator_endpoint = "http://<你的Druid Coordinator基础地址>"
        leader_resp = requests.get(f"{druid_coordinator_endpoint}/druid/coordinator/v1/leader")
        leader_info = leader_resp.json()
        leader_host = leader_info["host"]  # 假设返回的host为私有IP
    
        # 匹配对应EC2实例
        ec2 = boto3.client('ec2')
        instances = ec2.describe_instances(
            Filters=[
                {'Name': 'private-ip-address', 'Values': [leader_host]},
                {'Name': 'tag:druid_role', 'Values': ['coordinator']}
            ]
        )
    
        if not instances['Reservations']:
            return {"status": "error", "msg": "未找到匹配的Coordinator实例"}
    
        instance_id = instances['Reservations'][0]['Instances'][0]['InstanceId']
    
        # 终止实例(ASG会自动重启新节点;若需仅重启,可改用stop_instances后调用start_instances)
        ec2.terminate_instances(InstanceIds=[instance_id])
    
        return {"status": "success", "terminated_instance_id": instance_id}
    

    注意:需为Lambda角色配置ec2:DescribeInstances和ec2:TerminateInstances(或ec2:StopInstances/ec2:StartInstances)权限。

2. 通过AWS Lambda识别Segment重平衡卡顿

利用Druid Coordinator的API接口获取Segment状态,结合时间阈值判断是否卡顿:

  • 核心API接口:

    • GET /druid/coordinator/v1/loadstatus:返回所有Segment的加载状态(LOADING/LOADED/FAILED等);
    • GET /druid/coordinator/v1/balance/status:返回当前平衡任务的进度、待处理移动任务数;
    • GET /druid/coordinator/v1/metadata/segments?full:返回Segment完整元数据,包含最后更新时间。
  • 卡顿判断逻辑:

    • 标记处于LOADING/MOVING状态且超过指定阈值(如30分钟)的Segment;
    • 检查平衡任务的最后更新时间,若超过阈值且仍有待处理任务,则判定为任务卡顿;
    • 对比Segment的lastUpdated时间与当前时间,超出阈值且状态未更新则视为卡顿。

    示例代码:

    import requests
    from datetime import datetime, timezone
    
    def lambda_handler(event, context):
        druid_coordinator_endpoint = "http://<你的Druid Coordinator基础地址>"
        threshold_minutes = 30  # 自定义卡顿时间阈值
        stuck_items = []
        current_time = datetime.now(timezone.utc)
    
        # 检查Segment加载状态
        load_status_resp = requests.get(f"{druid_coordinator_endpoint}/druid/coordinator/v1/loadstatus")
        load_status = load_status_resp.json()
    
        for seg_id, status in load_status.items():
            if status["status"] in ["LOADING", "MOVING"]:
                # 获取Segment元数据中的最后更新时间
                seg_meta_resp = requests.get(f"{druid_coordinator_endpoint}/druid/coordinator/v1/metadata/segments/{seg_id}")
                seg_meta = seg_meta_resp.json()
                last_updated = datetime.fromisoformat(seg_meta["lastUpdated"].replace('Z', '+00:00'))
                time_diff = (current_time - last_updated).total_seconds() / 60
    
                if time_diff > threshold_minutes:
                    stuck_items.append({
                        "type": "segment",
                        "segment_id": seg_id,
                        "status": status["status"],
                        "last_updated": seg_meta["lastUpdated"],
                        "stuck_duration": f"{round(time_diff)}分钟"
                    })
    
        # 检查平衡任务状态
        balance_status_resp = requests.get(f"{druid_coordinator_endpoint}/druid/coordinator/v1/balance/status")
        balance_status = balance_status_resp.json()
        pending_moves = balance_status.get("pendingMoves", 0)
        last_progress_time = balance_status.get("lastProgressTime")
    
        if last_progress_time and pending_moves > 0:
            progress_time = datetime.fromisoformat(last_progress_time.replace('Z', '+00:00'))
            progress_diff = (current_time - progress_time).total_seconds() / 60
            if progress_diff > threshold_minutes:
                stuck_items.append({
                    "type": "balance_task",
                    "pending_moves": pending_moves,
                    "last_progress_time": last_progress_time,
                    "stuck_duration": f"{round(progress_diff)}分钟"
                })
    
        return {
            "status": "success",
            "stuck_items": stuck_items,
            "total_stuck_count": len(stuck_items)
        }
    

内容的提问来源于stack exchange,提问作者Anil

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.05 12:45:25