Работая единственным QA в команде, в обязанности которого активно входит мониторинг, я столкнулась с проблемами тайм-менеджмента: очень много времени уходит на анализ упавших пайплайнов, заведение однотипных задач, поднятие тревоги, отмена тревоги, потому что проблема уже исправляется и прочие прелести тестерской работы. В один прекрасный осенний денек, сидя на окне и думая о нем (мониторинге, разумеется), я решила - хватит это терпеть. Посидела, погуглила, вспомнила родительские слова: “Хочешь сделать что-то хорошо - сделай это сама”. Что ж таков путь.
Задача: При падении шага пайплайна создавать задачу в Youtrack, со всеми необходимыми данными. Обязательно логируем, чтобы при желании можно было отправлять логи в телегу или еще куда удобно. Но в статье не будем рассматривать пересылку, ибо там все просто - постим логи на нужный url, да и все, грубо говоря.
Что понадобится:
блокнот (ну или IDE, если вы любитель всяких примочек)
Airflow (+ тестовый даг)
Youtrack
Плагин, как ни странно, нужно приютить в директории plugins, которая должна быть в корне проекта:
airflow / ├── .env ├── docker-compose.yaml ├── dags/ └── plugins/ └── my_plugin/ ├── __init__.py └── my_plugin.py
Подготавливаем переменные и указываем их в .env:
YOUTRACK_URI = ‘https://<you_address>/api/issues’
YOUTRACK_TOKEN - получаем в профиле YouTrack: Profile -> Account Security
YOUTRACK_PROJECT_ID - буквенный идентификатор проекта
В docker-compose.yaml прокидываем эти переменные:
environment: - YOUTRACK_URI=${YOUTRACK_URI} - YOUTRACK_TOKEN=${YOUTRACK_TOKEN} - YOUTRACK_PROJECT_ID=${YOUTRACK_PROJECT_ID}
Почему я сделала именно так, а не через Variable - помимо того, что это плагин у нас крутится на продуктовом Airflow, мы его еще используем на регрессии. Во время регрессии может случиться что угодно, поэтому надежнее всего хранить переменные в файле, а не в БД, как этого требует Variable.
Создаем файл с именем listener.py.
Сразу же импорты:
import requests import logging import json import os from airflow.listeners import hookimpl from airflow.models import DagRun, TaskInstance from airflow.utils.state import TaskInstanceState
requests - отправка HTTP-запросов (YouTrack);
logging - логирование json - работа с JSON данными;
json - работа с JSON-объектами;
os - доступ к переменным окружения;
hookimpl - декоратор, который сообщает Airflow: “эта функция — слушатель событий”;
DagRun - объекты, которые Airflow передаст слушателю при событии;
TaskInstanceState - набор статусов для понимания, когда слушателю работать;
Подгружаем переменные окружения:
YOUTRACK_URI = os.environ.get("YOUTRACK_URI") YOUTRACK_TOKEN = os.environ.get("YOUTRACK_TOKEN") YOUTRACK_PROJECT_ID = os.environ.get("YOUTRACK_PROJECT_ID")
Здесь же сразу укажем поля, которые будем заполнять в задаче, они не изменяемые, поэтому кладем их рядом с остальными переменными:
custom_fields = { "Priority": "Critical", "Type": "Task", "Kind": "Defect", "Team": "DataPipeline Core" }
Объявляем на весь Airflow, что функция является слушателем. Когда задача завершается с ошибкой, Airflow автоматически вызовет ее:
@hookimpl def on_task_instance_failed(previous_state: TaskInstanceState, task_instance: TaskInstance):
Как Airflow определяет, когда какую функцию использовать:
Ищется декоратор hookimpl (объявлятор слушателей);
Смотрится название функции, которая описывает событие. Можно посмотреть в документации какие еще бывают события и их названия.
В функцию передается таска и ее параметры, то есть объект TaskInstance, статус передается, чтобы убедиться, что задача зафэйлилась.
Теперь опишем нашего слушателя. Когда задача завершается с ошибкой, Airflow автоматически вызовет его. Из упавшего дага собираем данные, которые будем потом отправлять в задачу:
state: TaskInstanceState = task_instance.state task_name: str = task_instance.task_id start_date = task_instance.start_date end_date = task_instance.end_date dagrun = task_instance.dag_run dagrun_status = dagrun.state dagrun_id = dagrun.dag_id dag = task_instance.task.dag dag_name = dag.dag_id if dag else None log_url = task_instance.log_url
Проверяем сколько раз было падение таски. Нередко бывает, что случается краткосрочный сбой сети, где-то БД из-за большого объема не успела отработать, ну или черти в аду раскалили зеленый кубик докрасна, но потом позеленили его обратно. У нас было принято, что если после 5 попыток таска так и не взлетела, то все - крах, прах и мрак. Если таска все же полегла смертью храбрых, то идем проверять не существует ли уже такая задача:
try_number = task_instance.try_number if try_number == 5: find_task(dag_name, dagrun_id, task_name, start_date, end_date, log_url)
Теперь разберемся как нам найти существующую задачу. Мы в команде условились, что задачи будут иметь тему формата: [BUG][{dag_name}] Падение {task_name}. Ищем задачу с этой темой, если задача есть - логируем, что проблема известна, если нет - подготавливаем необходимые данные
Функция будет принимать все необходимые данные дага:
def find_task(dag_name, dagrun_id, task_name, start_date, end_date, log_url):
Логируем начало работы, это важный этап, т.к. с первого раза все равно ничего не получится, и чтобы долго не думать, где проблема, пихаем логи во все значимые места, не стесняемся:
logging.info(f’Start find task on Youtrack’)
Указываем данные, которые нам нужны для запроса: тема задачи, заголовки, параметры:
issue_theme = f"[BUG][{dag_name}] Падение {task_name} headers = { 'Authorization': f'Bearer {youtrack_token}', 'Content-Type': 'application/json', 'Accept': 'application/json' } params = { "query": f'project: {youtrack_project} AND summary: "{issue_theme}" AND State: -Complete' }
issue_theme - паттерн темы задачи, о котором мы договорились с коллегами;
headers - заголовки для общения с YouTrack: авторизация - по токену, общение - через json;
params - здесь кроется основная логика поиска:
ключ “query” содержит строку с поисковым запросом для YouTrack,
project - проект, для которого создаются задачи,
summary - тема задачи, в нашем случае паттерн,
State: -Complete - ищем задачу со статусом, отличным от “Закрыто”, на всякий пожарный поясню, что если задача в работе, то проблема считается известной, а для известных проблем новые задачи нам не требуются.
Теперь делаем простой get-запрос со всем подготовленными данными и логируем:
response = requests.get(youtrack_url, headers=headers, params=params) logging.info(f"Answer: {response.status_code}: {response.text}")
Далее принимаем решение создавать задачу или нет. Первым делом проверяем статус-код, если 200, работаем дальше:
if response.status_code == 200:
Обрабатываем тело запроса. Получаем json:
issues = response.json()
Если он не пустой - значит задача найдена, на этом, конечно, можно остановиться, но удобно видеть что это за задача, поэтому напишем еще буквально пару строк. На этапе внедрения плагина у нас может оказаться несколько задач на одну и ту же проблему, т.к. задачи заводились руками и никто не отменял (к сожалению) человеческий фактор, поэтому перебираем найденные задачи и логируем их id:
if issues: for issue in issues: logging.info(f'Task is exist: {issue['id']}')
Если же в теле ответа нет задач, логируем этот момент и подготавливаем данные для создания задачи:
else: logging.info("Task is not exist") summary = issue_theme description = f"Dag failed: {dag_name} \n \ Dag ID: {dagrun_id}\n \ Task id: {task_name} \n \ Start time: {start_date} \n \ Fail time: {end_date}\n \ Logs: {log_url}\n"
Далее создаем задачу. Можно вынести, конечно, это в другое место, но у нас плагин достаточно маленький, так что тут над архитектурой я решила не заморачиваться. Вызываем функцию создания задачи и передаем в нее необходимые данные:
created_issue = create_task(youtrack_url, youtrack_token, youtrack_project, summary, description, custom_fields)
Логируем удачность/неудачность создания:
if created_issue: logging.info(f"Issue created successfully: {created_issue['id']}") else: logging.info(f"Failed to create issue.")
Чуть не забыли залогировать ошибку, когда статус не 200
else: logging.info(f"Ошибка: {response.status_code}: {response.text}")
Теперь разберемся, как создавать задачи в YouTrack: Необходимо получить: тему, описание проблемы, спец.поля для заполнения: приоритет, тип задачи, тип проблемы, команда.
def create_task(summary, description, fields={}):
Не забываем логировать важные моменты:
logging.info('Start create task on Youtrack')
Указываем данные для запроса создания задачи:
headers = { "Authorization": f"Bearer {youtrack_token}", "Content-Type": "application/json", 'Accept': 'application/json' } issue_data = { "project": {"id": youtrack_project}, "summary": summary, "description": description, **fields }
Когда все заранее подготовлено, пытаемся отослать простой post-запрос, data принимает json, т.к. мы указали это в заголовке, поэтому не забываем преобразовать словарь
try: response = requests.post(youtrack_url, headers=headers, data=json.dumps(issue_data))
Запускаем проверку статуса:
response.raise_for_status()
Если все прошло успешно возвращаем созданную задачу:
created_issue = response.json() return created_issue
Если во время выполнения запроса произошла ошибка - логируем ее и возвращаем None:
except requests.exceptions.RequestException as e: logging.info(f"Error creating YouTrack issue: {e}") return None
Если ответ пришел не в json-формате, так же логируем и возвращаем None:
except json.JSONDecodeError as e: logging.info(f"Error decoding YouTrack response: {e}") logging.info(f"Response text: {response.text}") return None
Теперь наш плагин нужно инициализировать.
В папке с нашим плагином создаем файл init.py и импортируем менеджер плагинов и модуль логики, скажем так:
from airflow.plugins_manager import AirflowPlugin import listener
Создаем класс плагина, даем имя и в списке слушателей указываем наш свеженаписанный:
class FailDagToYoutrackPlugin(AirflowPlugin): name = 'FailDagToYoutrackPlagin' listeners = [listener]
Собственно это все, сохраняем, пересобираем Airflow, запускаем сломанный даг и наслаждаемся прелестями автоматизации.
