yarn任务重启造成dinky1.1.0的flink任务uknown的问题
·
一、背景
数函通相关实时任务已经在dinky开发后,通过dinky提交yarn集群采用yarn-application模式执行。
yarn集群部分节点偶尔会有重启、宕机的问题出现,这对于yarn这种高可用分布式平台没有影响,集群会经过短时间的failover自动恢复。
yarn job恢复后可能涉及以下变更:
- yarn job重启造成job id变更。
- yarn job的job id不变更,但下面的attemptId变更。
- yarn job整体不变更,但flink jobid变更。
dinky目前是通过job manager地址和flink jobid来监控flink job的状态,如果上述三部分变更发生,就会造成dinky无法获取相应flink job状态。
二、问题
dinky目前原生告警项都是启用的,包括unknown告警。
如果hadoop集群(yarn)的job重启,会产生如下问题:
- hadoop不会主动通知dinky,导致dinky依然使用旧的连接数据监控flink任务状态,那么dinky就无法监控到flink的任务状态,进而产生unknown告警。
- 如此产生的unknown告警不会自动修复,必须人工主动去hadoop任务主页查看job manager地址和flink jobid来重新配置dinky任务信息。
- 人工运维这些信息工作量庞大且均为重复劳动,动作滞后。
- 由于我们现在主要依靠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()}")
更多推荐
所有评论(0)