一、pushgateway简介

1.pushgateway的概念
pushgateway 是采用被动推送的方式,而不是类似于 prometheus server 主动连接 exporter 获取监控数据。
pushgateway 可以单独运行在一个节点,然后需要自定义监控脚本把需要监控的主动推送给 pushgateway的 API 接口, 然后 pushgateway 再等待 prometheus server 抓取数据
2.pushgateway的特点
pushgateway 本身没有任何抓取监控数据的功能
目前 pushgateway 只是被动的等待数据从客户端推送过来。

二、部署

1. 解压安装包

 unzip pushgateway-1.4.2.linux-amd64.zip

2.启动

sh start.sh 
#启动内容
DEPLOY_DIR=$(dirname "$0")
nohup $DEPLOY_DIR/pushgateway --web.listen-address :9091 --persistence.file="" > pushgateway.log 2>&1 &
  • 检查服务是否启动:
   netstat -tnlp | grep 9091

三、自定义脚本编写

  • 实现yarn Scheduler 调度器指标监控
#!/bin/bash
curl -X DELETE "http://pushgatewayIP:9091/metrics/job/分组标签/instance/主机名"
# 配置信息
YARN_API="http://yarnIP:8088/ws/v1/cluster/scheduler"
PUSHGATEWAY_URL="http://pushgatewayIP:9091"
JOB_NAME="分组标签"
INSTANCE_NAME=$(hostname -f)

# 临时文件
METRICS_FILE="/tmp/fair_scheduler_metrics.prom.$$"

# 日志函数
log() {
    echo "[$(date '+%Y-%m-%d %H:%M:%S')] $1" >&2
}

error() {
    echo "[$(date '+%Y-%m-%d %H:%M:%S')] ERROR: $1" >&2
}

# 清理函数
cleanup() {
    rm -f "$METRICS_FILE"
    log "临时文件已清理: $METRICS_FILE"
}

# 注册退出时的清理
trap cleanup EXIT

# 生成Prometheus指标
generate_metrics() {
    local json_data="$1"
    
    # 清空指标文件
    > "$METRICS_FILE"
    
    # 首先写入所有指标的 HELP 和 TYPE 定义
    cat >> "$METRICS_FILE" << 'EOF'
# HELP yarn_queue_used_memory_tb Fair scheduler queue used memory in TB
# TYPE yarn_queue_used_memory_tb gauge
# HELP yarn_queue_used_vcores Fair scheduler queue used vCores
# TYPE yarn_queue_used_vcores gauge
# HELP yarn_queue_min_memory_tb Fair scheduler queue min memory in TB
# TYPE yarn_queue_min_memory_tb gauge
# HELP yarn_queue_min_vcores Fair scheduler queue min vCores
# TYPE yarn_queue_min_vcores gauge
# HELP yarn_queue_max_memory_tb Fair scheduler queue max memory in TB
# TYPE yarn_queue_max_memory_tb gauge
# HELP yarn_queue_max_vcores Fair scheduler queue max vCores
# TYPE yarn_queue_max_vcores gauge
# HELP yarn_queue_steady_share_memory Fair scheduler queue steady fair share memory
# TYPE yarn_queue_steady_share_memory gauge
# HELP yarn_queue_steady_share_vcores Fair scheduler queue steady fair share vCores
# TYPE yarn_queue_steady_share_vcores gauge
# HELP yarn_queue_instant_share_memory Fair scheduler queue instantaneous fair share memory
# TYPE yarn_queue_instant_share_memory gauge
# HELP yarn_queue_instant_share_vcores Fair scheduler queue instantaneous fair share vCores
# TYPE yarn_queue_instant_share_vcores gauge
# HELP yarn_queue_active_applications Fair scheduler queue active applications
# TYPE yarn_queue_active_applications gauge
# HELP yarn_queue_pending_applications Fair scheduler queue pending applications
# TYPE yarn_queue_pending_applications gauge
# HELP yarn_queue_memory_usage_percent Fair scheduler queue memory usage percentage
# TYPE yarn_queue_memory_usage_percent gauge
# HELP yarn_queue_vcores_usage_percent Fair scheduler queue vCores usage percentage
# TYPE yarn_queue_vcores_usage_percent gauge
# HELP yarn_queue_memory_min_usage_percent Fair scheduler queue memory usage vs min resources percentage
# TYPE yarn_queue_memory_min_usage_percent gauge
# HELP yarn_queue_vcores_min_usage_percent Fair scheduler queue vCores usage vs min resources percentage
# TYPE yarn_queue_vcores_min_usage_percent gauge
EOF

    # 获取所有队列
    local queue_list_file="/tmp/fair_queue_list.$$"
    echo "$json_data" | jq -r '.. | .queueName? // empty' | grep -v "^root$" > "$queue_list_file"
    
    # 检查是否有队列
    if [ ! -s "$queue_list_file" ]; then
        error "未找到任何队列"
        return 1
    fi
    
    log "找到队列: $(cat "$queue_list_file" | tr '\n' ' ')"
    
    # 处理每个队列
    while IFS= read -r queue_name; do
        [ -z "$queue_name" ] && continue
        
        log "处理队列: $queue_name"
        
        # 提取队列数据
        local queue_data
        queue_data=$(echo "$json_data" | jq -c --arg q "$queue_name" '.. | select(.queueName? == $q)')
        
        # 基础资源指标
        local used_memory=$(echo "$queue_data" | jq -r '.usedResources.memory // 0')
        local used_memory_tb=$(echo "scale=4; $used_memory / 1024 / 1024" | bc)
        local used_vcores=$(echo "$queue_data" | jq -r '.usedResources.vCores // 0')
        local min_memory=$(echo "$queue_data" | jq -r '.minResources.memory // 0')
        local min_memory_tb=$(echo "scale=4; $min_memory / 1024 / 1024" | bc)
        local min_vcores=$(echo "$queue_data" | jq -r '.minResources.vCores // 0')
        local max_memory=$(echo "$queue_data" | jq -r '.maxResources.memory // 0')
        local max_memory_tb=$(echo "scale=4; $max_memory / 1024 / 1024" | bc)
        local max_vcores=$(echo "$queue_data" | jq -r '.maxResources.vCores // 0')
        local steady_memory=$(echo "$queue_data" | jq -r '.steadyFairResources.memory // 0')
        local steady_vcores=$(echo "$queue_data" | jq -r '.steadyFairResources.vCores // 0')
        local instant_memory=$(echo "$queue_data" | jq -r '.fairResources.memory // 0')
        local instant_vcores=$(echo "$queue_data" | jq -r '.fairResources.vCores // 0')
        
        # 应用数量指标
        local active_apps=$(echo "$queue_data" | jq -r '.numActiveApplications // 0')
        local pending_apps=$(echo "$queue_data" | jq -r '.numPendingApplications // 0')
        
        # 计算使用率
        local memory_usage=0
        local vcores_usage=0
        local memory_min_usage=0
        local vcores_min_usage=0
        
        # 计算相对于最大资源的使用率
        if [ "$max_memory" -gt 0 ]; then
            memory_usage=$(echo "scale=2; $used_memory * 100 / $max_memory" | bc)
        fi
        
        if [ "$max_vcores" -gt 0 ]; then
            vcores_usage=$(echo "scale=2; $used_vcores * 100 / $max_vcores" | bc)
        fi
        
        # 计算相对于最小资源的使用率
        if [ "$min_memory" -gt 0 ]; then
            memory_min_usage=$(echo "scale=2; $used_memory * 100 / $min_memory" | bc)
        fi
        
        if [ "$min_vcores" -gt 0 ]; then
            vcores_min_usage=$(echo "scale=2; $used_vcores * 100 / $min_vcores" | bc)
        fi
        
        # 写入队列指标数据
        cat >> "$METRICS_FILE" << EOF
yarn_queue_used_memory_tb{instance="$INSTANCE_NAME",queue="$queue_name"} $used_memory_tb
yarn_queue_used_vcores{instance="$INSTANCE_NAME",queue="$queue_name"} $used_vcores
yarn_queue_min_memory_tb{instance="$INSTANCE_NAME",queue="$queue_name"} $min_memory_tb
yarn_queue_min_vcores{instance="$INSTANCE_NAME",queue="$queue_name"} $min_vcores
yarn_queue_max_memory_tb{instance="$INSTANCE_NAME",queue="$queue_name"} $max_memory_tb
yarn_queue_max_vcores{instance="$INSTANCE_NAME",queue="$queue_name"} $max_vcores
yarn_queue_steady_share_memory{instance="$INSTANCE_NAME",queue="$queue_name"} $steady_memory
yarn_queue_steady_share_vcores{instance="$INSTANCE_NAME",queue="$queue_name"} $steady_vcores
yarn_queue_instant_share_memory{instance="$INSTANCE_NAME",queue="$queue_name"} $instant_memory
yarn_queue_instant_share_vcores{instance="$INSTANCE_NAME",queue="$queue_name"} $instant_vcores
yarn_queue_active_applications{instance="$INSTANCE_NAME",queue="$queue_name"} $active_apps
yarn_queue_pending_applications{instance="$INSTANCE_NAME",queue="$queue_name"} $pending_apps
yarn_queue_memory_usage_percent{instance="$INSTANCE_NAME",queue="$queue_name"} $memory_usage
yarn_queue_vcores_usage_percent{instance="$INSTANCE_NAME",queue="$queue_name"} $vcores_usage
yarn_queue_memory_min_usage_percent{instance="$INSTANCE_NAME",queue="$queue_name"} $memory_min_usage
yarn_queue_vcores_min_usage_percent{instance="$INSTANCE_NAME",queue="$queue_name"} $vcores_min_usage
EOF
        
    done < "$queue_list_file"
    
    # 清理临时文件
    rm -f "$queue_list_file"
    
    echo "$METRICS_FILE"
}

# 验证指标文件格式
validate_metrics_format() {
    local metrics_file="$1"
    
    log "验证指标文件格式..."
    
    if [ ! -f "$metrics_file" ] || [ ! -s "$metrics_file" ]; then
        error "指标文件不存在或为空"
        return 1
    fi
    
    if ! grep -q "HELP" "$metrics_file" || ! grep -q "TYPE" "$metrics_file"; then
        error "指标文件缺少 HELP 或 TYPE 定义"
        return 1
    fi
    
    local data_lines=$(grep -v "^#" "$metrics_file" | grep -v "^$" | wc -l)
    if [ "$data_lines" -eq 0 ]; then
        error "指标文件没有数据行"
        return 1
    fi
    
    log "指标格式验证通过,数据行数: $data_lines"
    return 0
}

# 推送指标到Pushgateway
push_metrics() {
    local metrics_file="$1"
    
    if ! validate_metrics_format "$metrics_file"; then
        return 1
    fi
    
    log "推送指标到 Pushgateway..."
    
    # 先删除旧数据
    curl -X DELETE "http://pushgatewayIP:9091/metrics/job/分组标签/instance/主机名  >/dev/null 2>&1
    sleep 1
    
    # 推送新指标
    local response
    response=$(curl -s -o /dev/null -w "%{http_code}" \
        --connect-timeout 10 \
        --max-time 30 \
        -X POST \
        --data-binary "@$metrics_file" \
        "$PUSHGATEWAY_URL/metrics/job/$JOB_NAME/instance/$INSTANCE_NAME")
    
    if [ "$response" = "200" ]; then
        log "指标推送成功 (HTTP $response)"
        return 0
    else
        error "指标推送失败 (HTTP $response)"
        return 1
    fi
}

# 主函数
main() {
    log "开始收集 Fair Scheduler 指标..."
    
    # 检查依赖
    if ! command -v jq &> /dev/null; then
        error "jq 命令未找到,请先安装"
        exit 1
    fi
    
    if ! command -v curl &> /dev/null; then
        error "curl 命令未找到"
        exit 1
    fi
    
    if ! command -v bc &> /dev/null; then
        error "bc 命令未找到,请先安装: yum install bcapt-get install bc"
        exit 1
    fi
    
    # 获取YARN API数据
    log "从 YARN API 获取数据: $YARN_API"
    local json_data
    json_data=$(curl -s --connect-timeout 30 --max-time 60 "$YARN_API")
    
    if [ $? -ne 0 ] || [ -z "$json_data" ]; then
        error "无法从 YARN API 获取数据"
        exit 1
    fi
    
    # 验证JSON数据
    if ! echo "$json_data" | jq -e . >/dev/null 2>&1; then
        error "YARN API 返回无效的JSON数据"
        exit 1
    fi
    
    # 生成指标
    log "生成指标文件..."
    local metrics_file
    metrics_file=$(generate_metrics "$json_data")
    
    # 检查指标文件是否正常生成
    if [[ "$metrics_file" =~ fair_scheduler_metrics\.prom\.[0-9]+$ ]] && [ -f "$metrics_file" ]; then
        # 显示前10行用于调试
        log "指标文件前10行:"
        head -10 "$metrics_file"
        
        # 推送指标
        if push_metrics "$metrics_file"; then
            log "Fair Scheduler 指标收集完成"
            exit 0
        else
            error "Fair Scheduler 指标收集失败"
            exit 1
        fi
    else
        error "指标文件生成失败"
        exit 1
    fi
}

# 执行主函数
main "$@"
  • 自定义presto指标监控
#!/bin/bash
# 脚本名称:monitor_presto.sh
# 定时任务:*/5 * * * * /root/monitor_presto.sh 2>&1 >> /var/log/monitor_presto.log
                                                                                                                  
# 增加错误处理
set -e
                                                                                                                  
# 配置变量
INSTANCE_NAME=`hostname -f | cut -d'.' -f1`  #获取本机名,用于后面的的标签
PUSHGATEWAY_URL="http://pushgatewayIP:9091/metrics/job/pushgateway/instance/$INSTANCE_NAME"  # Pushgateway地址
METRICS=""  #初始化指标变量
COORDINATOR="linking-presto-"  # presto地址和端口


# 1.Coordinator节点数量(固定为1)
coord_count=$(curl -s "http://$COORDINATOR/ui/api/stats" | jq -r '.activeCoordinators')
METRICS+="presto_coordinator_num $coord_count\n"

# 2. Worker节点数量(从/v1/node获取)
worker_count=$(curl -s "http://$COORDINATOR/ui/api/stats" | jq -r '.activeWorkers')
METRICS+="presto_worker_num $worker_count\n"


# 阻塞数

presto_blocked_queries=$(curl -s "http://$COORDINATOR/v1/query?state=QUEUED" | jq -r 'length')
METRICS+="presto_blocked_queries ${presto_blocked_queries:-0}\n"

# 排队等待数查询
queued_queries=$(curl -s "http://$COORDINATOR/ui/api/stats" | jq -r '.queuedQueries')
METRICS+="presto_queued_queries $queued_queries\n"


#运行的驱动器数量
running_drivers=$(curl -s "http://$COORDINATOR/ui/api/stats" | jq -r '.runningDrivers')
METRICS+="presto_running_drivers $running_drivers\n"


#集群中可用的处理器核心总数
total_processors=$(curl -s "http://$COORDINATOR/ui/api/stats" | jq -r '.totalAvailableProcessors')
METRICS+="presto_total_processors $total_processors\n"

#节点运行时长

uptime_count=$(curl -s "http://$COORDINATOR/v1/node" | \
  jq -r '.[] | "presto_uptime_count{uri=\"\(.uri)\"} \(.age | sub("d$"; "") | tonumber)"')

METRICS+="# HELP presto_uptime_count Total recent requests\n# TYPE presto_uptime_count  gauge\n$uptime_count\n"


#总请求数量

recent_requests=$(curl -s "http://$COORDINATOR/v1/node" | \
  jq -r '.[] | "presto_recent_requests{uri=\"\(.uri)\"} \(.recentRequests | floor)"')
METRICS+="# HELP presto_recent_requests Total recent requests\n# TYPE presto_recent_requests gauge\n$recent_requests\n"


#失败请求数

recent_failures=$(curl -s "http://$COORDINATOR/v1/node" | \
  jq -r '.[] | "presto_recent_failures{uri=\"\(.uri)\"} \(.recentFailures | floor)"')
METRICS+="# HELP presto_recent_failures Total recent requests\n# TYPE presto_recent_failures gauge\n$recent_failures\n"

##成功请求数

recent_success=$(curl -s "http://$COORDINATOR/v1/node" | \
  jq -r '.[] | "presto_recent_success{uri=\"\(.uri)\"} \(.recentSuccesses | floor)"')
METRICS+="# HELP presto_recent_success Total recent requests\n# TYPE presto_recent_success gauge\n$recent_success\n"

# 3. 查询成功数(FINISHED状态查询)
success_count=$(curl -s "http://$COORDINATOR/v1/query?state=FINISHED" | jq -r 'length' 2>/dev/null || echo 0)
METRICS+="presto_finish_job $success_count\n"
                                                                                                                  
# 4. 查询失败次数(FAILED状态查询)
fail_count=$(curl -s "http://$COORDINATOR/v1/query?state=FAILED" | jq -r 'length' 2>/dev/null || echo 0)
METRICS+="presto_failed_job $fail_count\n"
                                                                                                            
# 5. 当前活动查询数(RUNNING状态)
running_count=$(curl -s "http://$COORDINATOR/v1/query?state=RUNNING" | jq -r 'length' 2>/dev/null || echo 0)
METRICS+="presto_running_job $running_count\n"
                                                                                                               
# 6. Coordinator内存使用情况
mem_data=$(curl -s "http://$COORDINATOR/v1/status")
heap_used_mb=$(echo "$mem_data" | jq -r '.heapUsed / (1024*1024) | floor' 2>/dev/null || echo 0)
heap_used_kb=$(echo "$mem_data" | jq -r '.heapUsed  | floor' 2>/dev/null || echo 0)

heap_available_mb=$(echo "$mem_data" | jq -r '.heapAvailable / (1024*1024) | floor' 2>/dev/null || echo 1)
if [ $heap_available_mb -eq 0 ]; then
    heap_available_mb=1  # 避免除零错误
fi
mem_usage_pct=$(( 100 * heap_used_mb / heap_available_mb ))

METRICS+="presto_heap_used $heap_available_mb\n"
METRICS+="presto_heap_used_pct $mem_usage_pct\n"
METRICS+="presto_heap_used_kb $heap_used_kb\n"


WORKERS=("10.150.83.126" "10.150.83.188" "10.150.83.182")
PORT="8081"
METRICS="# HELP presto_worker_memory_max_bytes Max available memory in bytes\n"
METRICS+="# TYPE presto_worker_memory_max_bytes gauge\n"
METRICS+="# HELP presto_worker_memory_free_bytes Free memory in bytes\n"
METRICS+="# TYPE presto_worker_memory_free_bytes gauge\n"
METRICS+="# HELP presto_worker_memory_reserved_bytes Reserved memory in bytes\n"
METRICS+="# TYPE presto_worker_memory_reserved_bytes gauge\n"

for IP in "${WORKERS[@]}"; do
  DATA=$(curl -sf "http://$IP:$PORT/v1/memory/general")
  
  # 提取指标并追加
  METRICS+="presto_worker_memory_max_bytes{ip=\"$IP\"} $(echo "$DATA" | jq '.maxBytes')\n"
  METRICS+="presto_worker_memory_free_bytes{ip=\"$IP\"} $(echo "$DATA" | jq '.freeBytes')\n"
  METRICS+="presto_worker_memory_reserved_bytes{ip=\"$IP\"} $(echo "$DATA" | jq '.reservedBytes')\n"
done





# 最后推送到 Pushgateway
echo -e "$METRICS" | curl --data-binary @- $PUSHGATEWAY_URL
echo "监控数据已推送到 Pushgateway"

四、Prometheus配置pushgateway

- job_name: pushgateway
  static_configs:
    - targets: ['10.127.0.1:9091']
  honor_labels: true 
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐