Работая единственным 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 определяет, когда какую функцию использовать:

  1. Ищется декоратор hookimpl (объявлятор слушателей);

  2. Смотрится название функции, которая описывает событие. Можно посмотреть в документации какие еще бывают события и их названия.

  3. В функцию передается таска и ее параметры, то есть объект 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, запускаем сломанный даг и наслаждаемся прелестями автоматизации.