如何通过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
相关产品推荐
相关产品推荐

