一、背景

数函通相关实时任务已经在dinky开发后,通过dinky提交yarn集群采用yarn-application模式执行。

yarn集群部分节点偶尔会有重启、宕机的问题出现,这对于yarn这种高可用分布式平台没有影响,集群会经过短时间的failover自动恢复。

yarn job恢复后可能涉及以下变更:

  1. yarn job重启造成job id变更。
  2. yarn job的job id不变更,但下面的attemptId变更。
  3. yarn job整体不变更,但flink jobid变更。

dinky目前是通过job manager地址和flink jobid来监控flink job的状态,如果上述三部分变更发生,就会造成dinky无法获取相应flink job状态。

二、问题

dinky目前原生告警项都是启用的,包括unknown告警。

如果hadoop集群(yarn)的job重启,会产生如下问题:

  1. hadoop不会主动通知dinky,导致dinky依然使用旧的连接数据监控flink任务状态,那么dinky就无法监控到flink的任务状态,进而产生unknown告警。
  2. 如此产生的unknown告警不会自动修复,必须人工主动去hadoop任务主页查看job manager地址和flink jobid来重新配置dinky任务信息。
  3. 人工运维这些信息工作量庞大且均为重复劳动,动作滞后。
  4. 由于我们现在主要依靠dinky告警来监控flink任务状态,获取不到flink任务状态也会导致其他告警失效。

三、思路

我们解决这个问题的思路见下图。

四、交付物

脚本文件flink任务与yarn状态同步监控和配置脚本dw/realtime/cluster_monitor/yarn_monitor_flink_change.py
任务调度flink任务与yarn状态同步监控和配置调度realtime->实时运维工作流->dinky和yarn关于flink任务状态同步

dw/realtime/cluster_monitor/yarn_monitor_flink_change.py

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
import datetime
import random
import time

#  ** 语    言:Python
#  *********************************************************************************************************************
#  ** 所属主题:
#  ** 功能描述:获取yarn中任务的状态更新dinky
#  ** 创 建 者:
#  ** 创建日期:2025-04-02
#  ** 所需参数:
#  ** 修改日志:
#  ** 修改日期(yyyymmdd) 修改人(name) 修改内容(comment)
#  ** 20250402 lilianqiang init
#  *********************************************************************************************************************

# ##################################################### 核心参数 ########################################################

# ##################################################### 初始化 ##########################################################
import os
import sys
work_space = os.path.split(os.path.realpath(__file__))[0].split("dw")[0]
sys.path.append(work_space)

from dw.utils import date_util
from dw.utils.logger_util import logger_util
from dw.utils.database_util import mysql_dinky_execute
import datetime,requests
from bs4 import BeautifulSoup

# 准备程序所需参数
config_file_dir, file_name =  work_space + "dw/config/config.ini", os.path.split(os.path.realpath(__file__))[1].split(".")[0]
env = sys.argv[1]

logger = logger_util(file_name).logger
start_time = date_util.current_time()
logger.info("开始时间:" + str(start_time))
# ##################################################### 任务开始 ########################################################
# yarn-URL
resource_manager_url = "http://192.168.12.181:8088" if env == "Prod" else "http://192.168.12.211:8088"
# 指定队列名称
queue_name = "root.prod" if env == 'Prod' else 'root.default'

def get_yarn_jobs(queue_name, rm_url):
    api_url = f"{rm_url}/ws/v1/cluster/apps?queue={queue_name}&state=RUNNING&applicationTypes=Dinky%20Flink"
    try:
        response = requests.get(api_url)
        response.raise_for_status()  # 检查请求是否成功
        data = response.json()
        apps = data.get('apps', {}).get('app', [])
        return apps
    except requests.exceptions.RequestException as e:
        logger.error(f"请求失败: {e}")
        return []

def get_yarn_job_by_name(job_name, rm_url):
    jobs = get_yarn_jobs(queue_name, rm_url)
    matching_jobs = [job for job in jobs if job.get('name') == job_name]
    return matching_jobs

def get_yarn_attempts(job_id, rm_url):
    api_url = f"{rm_url}/ws/v1/cluster/apps/{job_id}/appattempts"
    try:
        response = requests.get(api_url)
        response.raise_for_status()
        data = response.json()
        attempts = data.get('appAttempts', {}).get('appAttempt', [])
        return attempts
    except requests.exceptions.RequestException as e:
        logger.error(f"请求失败: {e}")
        return []

def get_yarn_attempt_info(job_id, attempt_id, rm_url):
    api_url = f"{rm_url}/cluster/appattempt/{attempt_id}"
    try:
        response = requests.get(api_url)
        response.raise_for_status()
        return response.content.replace(b'\n', b'').replace(b' ', b'')
    except requests.exceptions.RequestException as e:
        logger.error(f"请求失败: {e}")
        return ''

def get_dinky_flink_job():
    try:
        flink_tasks = mysql_dinky_execute(config_file_dir, env, f"""
            select
                a.id,
                a.name,
                b.id as job_instance_id,
                b.jid,
                c.id as cluster_id,
                c.name as yarn_job_id,
                c.job_manager_host
            from
            (
                select
                    id,
                    name,
                    job_instance_id
                from dinky_task
                where enabled = 1
                    and step = 2
                    and type = 'yarn-application'
                    and dialect = 'FlinkSql'
                   -- and name = 'dwd-doris-dwd-risk-bankruptcy-reorganization-realtime'
            ) a
            inner join
            (
                select
                    id,
                    jid,
                    cluster_id
                from dinky_job_instance
                where status = 'UNKNOWN'
            ) b
                on a.job_instance_id = b.id
            inner join
            (
                select
                    id,
                    name,
                    type,
                    job_manager_host
                from dinky_cluster
                where enabled = 1
            ) c
                on b.cluster_id = c.id
            ;
        """)
        return flink_tasks
    except Exception as e:
        logger.error(f"{e}\n获取flink任务失败", exc_info=1)
        return

def get_flink_jobs(jobmanager_url):
    api_url = f"{jobmanager_url}/jobs/overview"
    try:
        response = requests.get(api_url)
        response.raise_for_status()  # 检查请求是否成功
        data = response.json()
        jobs = data.get('jobs', [])
        return jobs
    except requests.exceptions.RequestException as e:
        print(f"请求失败: {e}")
        return []

def get_flink_job_details(job_id, jobmanager_url):
    api_url = f"{jobmanager_url}/jobs/{job_id}"
    try:
        response = requests.get(api_url)
        response.raise_for_status()  # 检查请求是否成功
        data = response.json()
        return data
    except requests.exceptions.RequestException as e:
        print(f"请求失败: {e}")
        return None

def get_flink_job_by_name(job_name, jobmanager_url):
    jobs = get_flink_jobs(jobmanager_url)
    matching_jobs = [job for job in jobs if job.get('name') == job_name]
    return matching_jobs

flink_jobs = get_dinky_flink_job()
if flink_jobs:
    for flink_task in flink_jobs:
        flink_task_name = flink_task[1]
        job_instance_id = flink_task[2]
        flink_job_id = flink_task[3]
        cluster_id = flink_task[4]
        yarn_job_id = flink_task[5]
        yarn_jm_host = flink_task[6]
        logger.info(f"flink任务名: {flink_task_name}, flink instance id:{job_instance_id}, yarn job id:{flink_job_id}, yarn jm host:{cluster_id}, yarn job id:{yarn_job_id}, flink jobmanager host:{yarn_jm_host}")

        ### 检测dinky任务的yarn集群状态信息是否一致
        yarn_jobs = get_yarn_job_by_name(flink_task_name, resource_manager_url)
        if len(yarn_jobs) > 0:
            yarn_job = yarn_jobs[0]
            trackingUrl = yarn_job['trackingUrl']
            if yarn_job:
                logger.info(f"作业 ID: {yarn_job['id']}, 名称: {yarn_job['name']}, 状态: {yarn_job['state']}, 应用类型: {yarn_job['applicationType']}")
                attempts = get_yarn_attempts(yarn_job['id'], resource_manager_url)
                jm_address = ''
                for attempt in attempts:
                    logger.info(f"==尝试 ID: {attempt['id']}, 状态: {attempt['appAttemptId']}, am地址: {attempt['nodeId']}")
                    resp_content = get_yarn_attempt_info(yarn_job['id'], attempt['appAttemptId'], resource_manager_url)
                    soup = BeautifulSoup(resp_content, 'html.parser')
                    jm_address = soup.find_all(lambda tag:tag.string == 'Node:')[0].find_next().contents[0]
                    if jm_address and jm_address != 'N/A':
                        break
                if jm_address:
                    logger.info(f"==jm_address:{jm_address},yarn_jm_host:{yarn_jm_host}")
                    if jm_address != yarn_jm_host:
                        update_jm_sql = f"update dinky_cluster set job_manager_host = '{jm_address}' where id = {cluster_id}"
                        try:
                            mysql_dinky_execute(config_file_dir, env, update_jm_sql)
                            logger.info(f"#######################更新jm_address成功:\nsql:{update_jm_sql}")
                        except Exception as e:
                            logger.error(f"{e}\n更新jm_address失败:\nsql:{update_jm_sql}", exc_info=1)
                            exit(255)

            ### 检测dinky任务的flink集群状态信息是否一致
            flink_job = get_flink_job_by_name(flink_task_name, trackingUrl) # 这里必须按名字去查,因为当前mysql里的jid已经失效了
            if flink_job and len(flink_job) > 0:
                if flink_job[0]['jid']:
                    jid = flink_job[0]['jid']
                    logger.info(f"==jid:{jid},flink_job_id:{flink_job_id}")
                    if jid != flink_job_id:
                        update_jid_sql = f"update dinky_job_instance set jid = '{jid}', status = 'RUNNING' where id = {job_instance_id}"
                        try:
                            mysql_dinky_execute(config_file_dir, env, update_jid_sql)
                            logger.info(f"#######################更新jid成功:\nsql:{update_jid_sql}")
                        except Exception as e:
                            logger.error(f"{e}\n更新jid失败:\nsql:{update_jid_sql}", exc_info=1)
                            exit(255)
                else:
                    logger.error(f"flink集群中不存在jid是{flink_job_id}的flink作业。")
            else:
                logger.error(f"flink集群中不存在name是{flink_task_name}的flink作业。")
        else:
            logger.info(f"yarn集群中不存在{flink_task_name}的{yarn_job_id}作业,可能该作业在flink中已经被停止。")
else:
    logger.info(f"flink集群中不存在unkown的作业。")
# ##################################################### 任务结束 ########################################################
end_time = date_util.current_time()
delta_time = end_time - start_time
logger.info(f"开始时间:{str(start_time)}),结束时间:{str(end_time)},用时(秒):{delta_time.total_seconds()}")

更多推荐