• Главная
  • DataScience
  • CloudServicesEngineer
  • Поиск
Data Science

FS. ПР. Сокращатель ссылок

ПР. Сокращатель ссылок

Кратко:

  • Создание сокращателя ссылок с использованием сервисов Yandex Cloud.
  • Создание сервисного аккаунта, добавление ролей и создание бекета в Object Storage.
  • Создание бессерверной базы данных YDB и создание таблицы с помощью SQL-скрипта.
  • Создание функции для обработки ссылок и вставка кода в файл index.py.
  • Настройка Yandex API Gateway с использованием спецификации for-serverless-shortener.yml.
  • Запуск и проверка работоспособности созданного приложения.
  • Возможность дальнейшего развития и расширения функциональности приложения.

Практическая работа. Сокращатель ссылок

В рамках этого курса вы изучили несколько ключевых сервисов Yandex Cloud, относящихся к группе Serverless. Давайте объединим их для решения ещё одной практической задачи и создадим сервис, который конвертирует длинные ссылки в короткие.

Шаг 1. Сервисный аккаунт

Создание аккаунта

Создайте сервисный аккаунт с именем serverless-shortener:
 export SERVICE_ACCOUNT_SHORTENER_ID=$(yc iam service-account create --name serverless-shortener \
  --description "service account for serverless" \
  --format json | jq -r .) 
Проверьте текущий список сервисных аккаунтов:
yc iam service-account list
После проверки запишите идентификатор созданного сервисного аккаунта в переменную SERVICE_ACCOUNT_SHORTENER_ID:
echo "export SERVICE_ACCOUNT_SHORTENER_ID=<идентификатор сервисного аккаунта>" >> ~/.bashrc && . ~/.bashrc
echo $SERVICE_ACCOUNT_SHORTENER_ID

Назначение ролей

Добавьте созданному сервисному аккаунту роли editor, storage.viewer и ydb.admin:
echo "export FOLDER_ID=$(yc config get folder-id)" >> ~/.bashrc && . ~/.bashrc
echo $FOLDER_ID

echo "export OAUTH_TOKEN=$(yc config get token)" >> ~/.bashrc && . ~/.bashrc
echo $OAUTH_TOKEN

echo "export CLOUD_ID=$(yc config get cloud-id)" >> ~/.bashrc && . ~/.bashrc
echo $CLOUD_ID

yc resource-manager folder add-access-binding $FOLDER_ID \
  --subject serviceAccount:$SERVICE_ACCOUNT_SHORTENER_ID \
  --role editor

yc resource-manager folder add-access-binding $FOLDER_ID \
  --subject serviceAccount:$SERVICE_ACCOUNT_SHORTENER_ID \
  --role ydb.admin

yc resource-manager folder add-access-binding $FOLDER_ID \
  --subject serviceAccount:$SERVICE_ACCOUNT_SHORTENER_ID \
  --role storage.viewer

Шаг 2. Создание бакета в Object Storage

Сделаем для нашего сервиса веб-интерфейс. Поскольку это будет статическая веб-страница, разместим её в объектном хранилище.
В консоли управления в вашем рабочем каталоге выберите сервис Object Storage. Нажмите кнопку Создать бакет.
На странице создания бакета:
  1. Введите имя бакета. В нашем примере это будет storage-for-serverless-shortener.
  2. Ограничьте максимальный размер бакета (например 1 ГБ).
  3. Выберите тип доступа Публичный во всех случаях.
  4. Выберите класс хранилища Стандартное.
Нажмите кнопку Создать бакет для завершения операции.
image
Создайте файл index.html и загрузите его в созданный бакет — это будет стартовая страничка для нашего сокращателя:
<!DOCTYPE html>
<html lang="en">

<head>
    <meta charset="UTF-8">
    <title>Сокращатель URL</title>
    <!-- предостережет от лишнего GET запроса на адрес /favicon.ico -->
    <link rel="icon" href="data:;base64,iVBORw0KGgo=">
</head>

<body>
<h1>Добро пожаловать</h1>
<form action="javascript:shorten()">
    <label for="url">Введите ссылку:</label><br>
    <input id="url" name="url" type="text"><br>
    <input type="submit" value="Сократить">
</form>
<p id="shortened"></p>
</body>

<script>
    function shorten() {
        const link = document.getElementById("url").value
        fetch("/shorten", {
            method: 'POST',
            headers: {
                'Content-Type': 'application/json'
            },
            body: link
        })
            .then(response => response.json())
            .then(data => {
                const url = data.url
                document.getElementById("shortened").innerHTML = `<a href=${url}>${url}</a>`
            })
            .catch(error => {
                document.getElementById("shortened").innerHTML = `<p>Произошла ошибка ${error}, попробуйте еще раз</p>`
            })
    }
</script>

</html>

Шаг 3. Создание базы данных

Создадим бессерверную базу данных YDB с именем for-serverless-shortener. Чтобы не переключаться из терминала, снова воспользуемся CLI. Обязательно укажите флаг --serverless для выбора типа создаваемой базы данных.
yc ydb database create for-serverless-shortener \
  --serverless \
  --folder-id $FOLDER_ID

yc ydb database list
Далее выполните команду:
yc ydb database get --name for-serverless-shortener
В выводе вы увидите значение endpoint. Оно состоит из двух частей: собственно эндпоинта (обычно это ydb.serverless.yandexcloud.net:2135) и пути базы данных (он указывается после ключевого слова database и начинается с символа /, например /ru-central1/...).
Сохраним адрес эндпоинта в переменную YDB_ENDPOINT, а путь базы данных — в переменную YDB_DATABASE. Они пригодятся нам для подключения функции.
echo "export YDB_ENDPOINT=<YDB_ENDPOINT>" >> ~/.bashrc && . ~/.bashrc
echo $YDB_ENDPOINT

echo "export YDB_DATABASE=<YDB_DATABASE>" >> ~/.bashrc && . ~/.bashrc
echo $YDB_DATABASE
Для дальнейшей работы нам понадобится утилита интерфейса командной строки YDB CLI:
curl https://storage.yandexcloud.net/yandexcloud-ydb/install.sh | bash
С помощью CLI создадим авторизованный ключ сервисного аккаунта serverless-shortener:
yc iam key create \
--service-account-name serverless-shortener \
--output serverless-shortener.sa
Сохраним путь к файлу с ключом в переменную окружения:
echo "export SA_KEY_FILE=$PWD/serverless-shortener.sa" >> ~/.bashrc && . ~/.bashrc
echo $SA_KEY_FILE
Проверим работоспособность с помощью команды:
ydb \
  --endpoint $YDB_ENDPOINT \
  --database $YDB_DATABASE \
  --sa-key-file $SA_KEY_FILE \
  discovery whoami \
  --groups
  1. Сохраним в файл links.yql SQL-скрипт для создания таблицы:
    CREATE TABLE links
    (
        id Utf8,
        link Utf8,
        PRIMARY KEY (id)
    );
    COMMIT;
 
Запустите создание таблицы, а затем проверьте результат:
ydb \
  --endpoint $YDB_ENDPOINT \
  --database $YDB_DATABASE \
  --sa-key-file $SA_KEY_FILE \
  scripting yql --file links.yql

ydb \
  --endpoint $YDB_ENDPOINT \
  --database $YDB_DATABASE \
  --sa-key-file $SA_KEY_FILE \
  scheme describe links

Шаг 4. Создание функции

В рабочем каталоге создайте файл index.py:
import ydb
import urllib.parse
import hashlib
import base64
import json
import os


def decode(event, body):
    # тело запроса может быть закодировано
    is_base64_encoded = event.get('isBase64Encoded')
    if is_base64_encoded:
        body = str(base64.b64decode(body), 'utf-8')
    return body


def response(statusCode, headers, isBase64Encoded, body):
    return {
        'statusCode': statusCode,
        'headers': headers,
        'isBase64Encoded': isBase64Encoded,
        'body': body,
    }


def get_config():
    endpoint = os.getenv("endpoint")
    database = os.getenv("database")
    if endpoint is None or database is None:
        raise AssertionError("Нужно указать обе переменные окружения")
    credentials = ydb.construct_credentials_from_environ()
    return ydb.DriverConfig(endpoint, database, credentials=credentials)


def execute(config, query, params):
    with ydb.Driver(config) as driver:
        try:
            driver.wait(timeout=5)
        except TimeoutError:
            print("Connect failed to YDB")
            print("Last reported errors by discovery:")
            print(driver.discovery_debug_details())
            return None

        session = driver.table_client.session().create()
        prepared_query = session.prepare(query)

        return session.transaction(ydb.SerializableReadWrite()).execute(
            prepared_query,
            params,
            commit_tx=True
        )


def insert_link(id, link):
    config = get_config()
    query = """
        DECLARE $id AS Utf8;
        DECLARE $link AS Utf8;

        UPSERT INTO links (id, link) VALUES ($id, $link);
        """
    params = {'$id': id, '$link': link}
    execute(config, query, params)


def find_link(id):
    print(id)
    config = get_config()
    query = """
        DECLARE $id AS Utf8;

        SELECT link FROM links where id=$id;
        """
    params = {'$id': id}
    result_set = execute(config, query, params)
    if not result_set or not result_set[0].rows:
        return None

    return result_set[0].rows[0].link


def shorten(event):
    body = event.get('body')

    if body:
        body = decode(event, body)
        original_host = event.get('headers').get('Origin')
        link_id = hashlib.sha256(body.encode('utf8')).hexdigest()[:6]
        # в ссылке могут быть закодированные символы, например, %. это помешает работе api-gateway при редиректе,
        # поэтому следует избавиться от них вызовом urllib.parse.unquote
        insert_link(link_id, urllib.parse.unquote(body))
        return response(200, {'Content-Type': 'application/json'}, False, json.dumps({'url': f'{original_host}/r/{link_id}'}))

    return response(400, {}, False, 'В теле запроса отсутствует параметр url')


def redirect(event):
    link_id = event.get('pathParams').get('id')
    redirect_to = find_link(link_id)

    if redirect_to:
        return response(302, {'Location': redirect_to}, False, '')

    return response(404, {}, False, 'Данной ссылки не существует')


# эти проверки нужны, поскольку функция у нас одна
# в идеале сделать по функции на каждый путь в api-gw
def get_result(url, event):
    if url == "/shorten":
        return shorten(event)
    if url.startswith("/r/"):
        return redirect(event)

    return response(404, {}, False, 'Данного пути не существует')


def handler(event, context):
    url = event.get('url')
    if url:
        # из API-gateway url может прийти со знаком вопроса на конце
        if url[-1] == '?':
            url = url[:-1]
        return get_result(url, event)

    return response(404, {}, False, 'Эту функцию следует вызывать при помощи api-gateway')
Создайте файл requirements.txt со следующим содержимым:
ydb==2.13.3
Находясь в директории с исходными файлами, упакуйте их в zip-архив:
zip src.zip index.py requirements.txt
Создадим нашу функцию for-serverless-shortener, задав все необходимые переменные. В переменные окружения функции необходимо добавить:
  • endpoint — нужно указать протокол grpcs:// и добавить значение Эндпоинт из секции YDB эндпоинт, обычно получается grpcs://ydb.serverless.yandexcloud.net:2135;
  • database — это значение поля База данных из секции YDB эндпоинт (начинается с /ru-central1/....);
  • USE_METADATA_CREDENTIALS — выставите значение переменной в 1.
Также сразу сделаем функцию публичной.
yc serverless function create \
  --name for-serverless-shortener \
  --description "function for serverless-shortener"

yc serverless function version create \
  --function-name for-serverless-shortener \
  --memory=256m \
  --execution-timeout=5s \
  --runtime=python39 \
  --entrypoint=index.handler \
  --service-account-id $SERVICE_ACCOUNT_SHORTENER_ID \
  --environment USE_METADATA_CREDENTIALS=1 \
  --environment endpoint=grpcs://ydb.serverless.yandexcloud.net:2135 \
  --environment database=$YDB_DATABASE \
  --source-path src.zip

yc serverless function allow-unauthenticated-invoke for-serverless-shortener

Шаг 5. Конфигурирование Yandex API Gateway

Создадим спецификацию for-serverless-shortener.yml со следующим содержанием:
openapi: 3.0.0
info:
  title: for-serverless-shortener
  version: 1.0.0
paths:
  /:
    get:
      x-yc-apigateway-integration:
        type: object_storage
        bucket:             <bucket_name>        # <-- имя бакета
        object:             <html_file>          # <-- имя html-файла
        presigned_redirect: false
        service_account:    <service_account_id> # <-- идентификатор сервисного аккаунта
      operationId: static
  /shorten:
    post:
      x-yc-apigateway-integration:
        type: cloud_functions
        function_id:  <function_id>               # <-- идентификатор функции
      operationId: shorten
  /r/{id}:
    get:
      x-yc-apigateway-integration:
        type: cloud_functions
        function_id:  <function_id>               # <-- идентификатор функции
      operationId: redirect
      parameters:
        - description: id of the url
          explode: false
          in: path
          name: id
          required: true
          schema:
            type: string
          style: simple
Не забудьте подставить в спецификацию актуальные для вас значения переменных.
Используем спецификацию для инициализации:
yc serverless api-gateway create \
  --name for-serverless-shortener \
  --spec=for-serverless-shortener.yml \
  --description "for serverless shortener"
В результате успешного создания API-шлюза получим значение параметра domain:
yc serverless api-gateway list
yc serverless api-gateway get --name for-serverless-shortener
Чтобы проверить работоспособность API-шлюза и созданного приложения целиком, скопируйте служебный домен (вида https://<идентификатор API Gateway>.apigw.yandexcloud.net/) и вставьте адрес в браузер.
Добавляйте адреса сайтов в форму, они будут сохранятся в базу данных. А вам будет доступна ссылка, за которой будет скрываться оригинальный адрес. Ваше приложение полностью работоспособно. Теперь вы умеете использовать serverless-стек технологий Yandex Cloud.
image
Итак, вы создали приложение с использованием Cloud Functions, API Gateway, Object Storage и YDB. Конечно, вы можете развивать его и дальше, расширяя функциональность.
Вводный курс по serverless-разработке на этом завершён. Осталось пройти лишь заключительный тест, который проверит ваши знания по всем рассмотренным в курсе сервисам.

 

Категория: Cloud Services Engineer
Просмотров: 1208

YMQ. ПР. Однократная отправка сообщений

ПР. Однократная отправка сообщений

Кратко:

  • Реализация проекта по конвертированию видеофайлов в GIF с использованием Yandex Cloud Functions.
  • Использование Yandex Message Queue для передачи задач в очередь.
  • Создание очереди с помощью утилиты aws или Terraform.
  • Создание базы данных в Yandex Database Service (YDB) и создание документной таблицы.
  • Создание объекта хранения в Yandex Object Storage (Object Storage) и получение URL для доступа к нему.
  • Создание функций для обработки видео и получения URL для результатов.
  • Создание триггера для вызова функции обработчика и отправки сообщений в очередь.
  • Протестирование системы с использованием идентификатора задачи для получения статуса из базы данных.
  • Удаление триггера после завершения работы

Практическая работа. Однократная отправка сообщений

В этой практической работе мы реализуем проект, который позволит пользователям конвертировать видеофайлы в GIF. Такая задача хорошо подходит для Cloud Functions, потому что конвертирование отнимает немало ресурсов процессора, и чем качественнее видео, тем больше ресурсов требуется на его обработку.

Почему для решения этой задачи нужны очереди?

Представим, что мы попытались решить эту задачу «в лоб». Пользователь заходит на страницу и вводит ссылку на видеофайл. Сервис скачивает его, конвертирует и отдает ссылку на GIF. Возникают две серьёзные проблемы:
  1. Синхронное соединение не всегда стабильно. Чем дольше вы его держите, тем выше вероятность, что оно разорвётся. В этом случае всё придётся сделать заново. А если соединение нестабильно, то пользователь может и не дождаться результата.
  2. Задача ресурсоёмкая: если сервисом одновременно воспользуются много пользователей с большими видеороликами, мощностей может не хватить.
Чтобы избежать этих проблем, в архитектуру сервиса необходимо встроить очередь.

Шаг 1. Сервисный аккаунт и Lockbox

Создание сервисного аккаунта
Создайте сервисный аккаунт с именем ffmpeg-account-for-cf:
export SERVICE_ACCOUNT=$(yc iam service-account create --name ffmpeg-account-for-cf \
  --description "service account for serverless" \
  --format json | jq -r .)
Проверьте текущий список сервисных аккаунтов:
yc iam service-account list
После проверки запишите ID созданного сервисного аккаунта в переменную SERVICE_ACCOUNT_ID:
echo "export SERVICE_ACCOUNT_FFMPEG_ID=<ID>" >> ~/.bashrc && . ~/.bashrc
echo $SERVICE_ACCOUNT_FFMPEG_ID

 

Назначение роли сервисному аккаунту

Добавим вновь созданному сервисному аккаунту роли storage.viewer, storage.uploader, ymq.reader, ymq.writer, ydb.admin, serverless.functions.invoker, и lockbox.payloadViewer:
echo "export FOLDER_ID=$(yc config get folder-id)" >> ~/.bashrc && . ~/.bashrc
echo $FOLDER_ID

yc resource-manager folder add-access-binding $FOLDER_ID \
  --subject serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --role storage.viewer

yc resource-manager folder add-access-binding $FOLDER_ID \
  --subject serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --role storage.uploader

yc resource-manager folder add-access-binding $FOLDER_ID \
  --subject serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --role ymq.reader

yc resource-manager folder add-access-binding $FOLDER_ID \
  --subject serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --role ymq.writer

yc resource-manager folder add-access-binding $FOLDER_ID \
  --subject serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --role ydb.admin

yc resource-manager folder add-access-binding $FOLDER_ID \
  --subject serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --role serverless.functions.invoker

yc resource-manager folder add-access-binding $FOLDER_ID \
  --subject serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --role lockbox.payloadViewer

yc resource-manager folder add-access-binding $FOLDER_ID \
  --subject serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --role editor
Вы можете назначить несколько ролей и с помощью команды set-access-binding. Но эта команда полностью перезаписывает права доступа к ресурсу и все текущие роли на него будут удалены! Поэтому сначала убедитесь, что ресурсу не назначены роли, которые вы не хотите потерять:
yc resource-manager folder list-access-bindings $FOLDER_ID

yc resource-manager folder set-access-bindings $FOLDER_ID \
  --access-binding role=storage.viewer,subject=serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --access-binding role=storage.uploader,subject=serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --access-binding role=ymq.reader,subject=serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --access-binding role=ymq.writer,subject=serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --access-binding role=ydb.admin,subject=serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --access-binding role=serverless.functions.invoker,subject=serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --access-binding role=lockbox.payloadViewer,subject=serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID \
  --access-binding role=editor,subject=serviceAccount:$SERVICE_ACCOUNT_FFMPEG_ID

 

Создание ключа доступа для сервисного аккаунта
Этот этап нужен для получения идентификатора ключа доступа и секретного ключа, которые будут использованы для загрузки файлов в Object Storage, работы с Yandex Message Queue и т. д. Для создания ключа доступа необходимо вызвать следующую команду:
yc iam access-key create --service-account-name ffmpeg-account-for-cf
В результате вы получите примерно следующее:
    access_key:
        id: ajefraollq5puj2tir1o
        service_account_id: ajetdv28pl0a1a8r41f0
        created_at: "2021-08-23T21:13:05.677319393Z"
        key_id: BTPNvWthv0ZX2xVmlPIU
    secret: cWLQ0HrTM0k_qAac43cwMNJA8VV_rfTg_kd4xVPi 
Здесь key_id — это идентификатор ключа доступа ACCESS_KEY_ID. А secret — это секретный ключ SECRET_ACCESS_KEY. Переменные ACCESS_KEY_ID и SECRET_ACCESS_KEY могут быть использованы для задания соответствующих значений aws_access_key_id и aws_secret_access_key при использовании библиотеки boto3.
 
Создание элемента в сервисе Lockbox
В сервисе Lockbox (находится на стадии Preview) создайте ваш первый секрет, состоящий из набора версий, в которых хранятся ваши данные. Версия содержит наборы ключей и значений:
  • Ключ — несекретное название для значения, по которому вы будете его идентифицировать.
  • Значение — это секретные данные.
Версия не изменяется. Для любого изменения количества пар ключей-значений или их содержимого необходимо создать новую версию. Создадим секрет с именем ffmpeg-sa-key и парой ключей ACCESS_KEY_ID и SECRET_ACCESS_KEY:
yc lockbox secret create --name ffmpeg-sa-key \
  --folder-id $FOLDER_ID \
  --description "keys for serverless" \
  --payload '[{"key": "ACCESS_KEY_ID", "text_value": <ACCESS_KEY_ID>}, {"key": "SECRET_ACCESS_KEY", "text_value": "<SECRET_ACCESS_KEY>"}]'
Получим и запишем значение SECRET_ID, оно нам потребуется при создании функции:
yc lockbox secret list

yc lockbox secret get --name ffmpeg-sa-key

echo "export SECRET_ID=<SECRET_ID>" >> ~/.bashrc && . ~/.bashrc
echo $SECRET_ID

 

Шаг 2. Создание очереди Yandex Message Queue

Для создания очереди Yandex Message Queue вы можете использовать три разных способа:
  • консоль управления;
  • консольная утилита aws;
  • Terraform.
Создание очереди с помощью утилиты aws
Воспользуемся AWS CLI. Для начала задайте конфигурацию с помощью команды aws configure. При этом от вас потребуется ввести:
  • AWS Access Key ID — идентификатор ключа доступа key_id сервисного аккаунта, полученный на предыдущем шаге.
  • AWS Secret Access Key — секретный ключ secret сервисного аккаунта, полученный на предыдущем шаге.
  • Default region name — используйте значение ru-central1.
По завершению конфигурации вы сможете создать очередь:
aws configure
aws sqs create-queue --queue-name ffmpeg --endpoint https://message-queue.api.cloud.yandex.net/
В результате успешного выполнения предыдущей команды в ответ вы получите URL:
    {
        "QueueUrl": "https://message-queue.api.cloud.yandex.net/b1ga4gj7agij03ln6aov/dj6000000003kv2t02b3/ffmpeg"
    } 
Запишем значения URL в переменную YMQ_QUEUE_URL. Она потребуется нам при создании функции:
echo "export YMQ_QUEUE_URL=<YMQ_QUEUE_URL>" >> ~/.bashrc && . ~/.bashrc
echo $YMQ_QUEUE_URL
Ещё вам потребует значение атрибута QueueArn, получим его:
aws sqs get-queue-attributes \
  --endpoint https://message-queue.api.cloud.yandex.net \
  --queue-url $YMQ_QUEUE_URL \
  --attribute-names QueueArn
В результате вы получите ответ вида:
    {
        "Attributes": {
            "QueueArn": "yrn:yc:ymq:ru-central1:b1gl21bkgss4msekt08i:ffmpeg"
        }
    } 
Сохраним значение QueueArn в переменную YMQ_QUEUE_ARN:
echo "export YMQ_QUEUE_ARN=<YMQ_QUEUE_ARN>" >> ~/.bashrc && . ~/.bashrc
echo $YMQ_QUEUE_ARN

Шаг 3. Создание базы данных в сервисе YDB

Создадим базу данных YDB с именем ffmpeg и типом serverless, используя для этого флаг --serverless:
yc ydb database create ffmpeg \
  --serverless \
  --folder-id $FOLDER_ID

yc ydb database list
Сразу получим и сохраним document_api_endpoint в значение переменной DOCAPI_ENDPOINT:
yc ydb database get --name ffmpeg
echo "export DOCAPI_ENDPOINT=<DOCAPI_ENDPOINT>" >> ~/.bashrc && . ~/.bashrc
echo $DOCAPI_ENDPOINT
Как только база данных создана, воспользуемся ранее использованной утилитой AWS CLI для создания документной таблицы в этой базе данных. Всю конфигурацию возьмем из файла tasks.json:
{
  "AttributeDefinitions": [
    {
      "AttributeName": "task_id",
      "AttributeType": "S"
    }
  ],
  "KeySchema": [
    {
      "AttributeName": "task_id",
      "KeyType": "HASH"
    }
  ],
  "TableName": "tasks"
}
Находясь в одном каталоге с файлом tasks.json, вызовите следующую команду для создания таблицы:
aws dynamodb create-table \
  --cli-input-json file://tasks.json \
  --endpoint-url $DOCAPI_ENDPOINT \
  --region ru-central1
В консоли управления убедитесь, что БД ffmpeg создана, и в ней есть пустая таблица tasks.

Шаг 4. Создание бакета в сервисе Object Storage

Самый простой способ создания бакета в Object Storage — это использование консоли управления.
В консоли управления в вашем рабочем каталоге выберите сервис Object Storage. Нажмите кнопку Создать бакет. На странице создания бакета введите имя, в нашем примере это будет storage-for-ffmpeg, остальные параметры не меняйте.
Нажмите кнопку Создать бакет для завершения операции. Далее вы всегда сможете поменять класс хранилища, его размер и настройки доступа.
Сохраним название бакета для дальнейшего использования:
echo "export S3_BUCKET=<имя бакета>" >> ~/.bashrc && . ~/.bashrc
echo $S3_BUCKET

Шаг 5. Создание функций

При создании функций нам потребуется ряд переменных:
  • SECRET_ID — идентификатор секрета (можно получить из таблицы со списком секретов);
  • YMQ_QUEUE_URL — URL очереди (можно получить на странице обзора);
  • DOCAPI_ENDPOINT — его можно получить на странице обзора БД, нужен именно Document API;
  • S3_BUCKET — имя бакета, в нашем случае это storage-for-ffmpeg.
Проверим заданные ранее переменные:
echo $SERVICE_ACCOUNT_FFMPEG_ID
echo $SECRET_ID
echo $YMQ_QUEUE_URL
echo $DOCAPI_ENDPOINT
echo $S3_BUCKET
Для обработки видео понадобится утилита FFmpeg. Скачайте статический релизный бинарный файл для Linux amd64 на сайте ffmpeg.org (обычно он находится в разделе FFmpeg Static Builds и называется примерно так: ffmpeg-release-amd64-static.tar.xz). Извлеките из архива файл ffmpeg. Обратите внимание, что у этого файла должны быть заданы права доступа на выполнение (установить нужный флаг можно с помощью команды chmod a+x ffmpeg).
Поскольку через консоль управления можно прикладывать файлы размером не более 3,5 МБ, загрузим код функций и файл ffmpeg в объектное хранилище (Object Storage).
Исходный код в файле index.py содержит обе необходимые нам функции:
import json
import os
import subprocess
import uuid
from urllib.parse import urlencode

import boto3
import requests
import yandexcloud
from yandex.cloud.lockbox.v1.payload_service_pb2 import GetPayloadRequest
from yandex.cloud.lockbox.v1.payload_service_pb2_grpc import PayloadServiceStub

boto_session = None
storage_client = None
docapi_table = None
ymq_queue = None


def get_boto_session():
    global boto_session
    if boto_session is not None:
        return boto_session

    # initialize lockbox and read secret value
    yc_sdk = yandexcloud.SDK()
    channel = yc_sdk._channels.channel("lockbox-payload")
    lockbox = PayloadServiceStub(channel)
    response = lockbox.Get(GetPayloadRequest(secret_id=os.environ['SECRET_ID']))

    # extract values from secret
    access_key = None
    secret_key = None
    for entry in response.entries:
        if entry.key == 'ACCESS_KEY_ID':
            access_key = entry.text_value
        elif entry.key == 'SECRET_ACCESS_KEY':
            secret_key = entry.text_value
    if access_key is None or secret_key is None:
        raise Exception("secrets required")
    print("Key id: " + access_key)

    # initialize boto session
    boto_session = boto3.session.Session(
        aws_access_key_id=access_key,
        aws_secret_access_key=secret_key
    )
    return boto_session


def get_ymq_queue():
    global ymq_queue
    if ymq_queue is not None:
        return ymq_queue

    ymq_queue = get_boto_session().resource(
        service_name='sqs',
        endpoint_url='https://message-queue.api.cloud.yandex.net',
        region_name='ru-central1'
    ).Queue(os.environ['YMQ_QUEUE_URL'])
    return ymq_queue


def get_docapi_table():
    global docapi_table
    if docapi_table is not None:
        return docapi_table

    docapi_table = get_boto_session().resource(
        'dynamodb',
        endpoint_url=os.environ['DOCAPI_ENDPOINT'],
        region_name='ru-central1'
    ).Table('tasks')
    return docapi_table


def get_storage_client():
    global storage_client
    if storage_client is not None:
        return storage_client

    storage_client = get_boto_session().client(
        service_name='s3',
        endpoint_url='https://storage.yandexcloud.net',
        region_name='ru-central1'
    )
    return storage_client

# API handler

def create_task(src_url):
    task_id = str(uuid.uuid4())
    get_docapi_table().put_item(Item={
        'task_id': task_id,
        'ready': False
    })
    get_ymq_queue().send_message(MessageBody=json.dumps({'task_id': task_id, "src": src_url}))
    return {
        'task_id': task_id
    }


def get_task_status(task_id):
    task = get_docapi_table().get_item(Key={
        "task_id": task_id
    })
    if task['Item']['ready']:
        return {
            'ready': True,
            'gif_url': task['Item']['gif_url']
        }
    return {'ready': False}


def handle_api(event, context):
    action = event['action']
    if action == 'convert':
        return create_task(event['src_url'])
    elif action == 'get_task_status':
        return get_task_status(event['task_id'])
    else:
        return {"error": "unknown action: " + action}

# Converter handler

def download_from_ya_disk(public_key, dst):
    api_call_url = 'https://cloud-api.yandex.net/v1/disk/public/resources/download?' + \
                   urlencode(dict(public_key=public_key))
    response = requests.get(api_call_url)
    download_url = response.json()['href']
    download_response = requests.get(download_url)
    with open(dst, 'wb') as video_file:
        video_file.write(download_response.content)


def upload_and_presign(file_path, object_name):
    client = get_storage_client()
    bucket = os.environ['S3_BUCKET']
    client.upload_file(file_path, bucket, object_name)
    return client.generate_presigned_url('get_object', Params={'Bucket': bucket, 'Key': object_name}, ExpiresIn=3600)


def handle_process_event(event, context):
    for message in event['messages']:
        task_json = json.loads(message['details']['message']['body'])
        task_id = task_json['task_id']
        # Download video
        download_from_ya_disk(task_json['src'], '/tmp/video.mp4')
        # Convert with ffmpeg
        subprocess.run(['ffmpeg', '-i', '/tmp/video.mp4', '-r', '10', '-s', '320x240', '/tmp/result.gif'])
        result_object = task_id + ".gif"
        # Upload to Object Storage and generate presigned url
        result_download_url = upload_and_presign('/tmp/result.gif', result_object)
        # Update task status in DocAPI
        get_docapi_table().update_item(
            Key={'task_id': task_id},
            AttributeUpdates={
                'ready': {'Value': True, 'Action': 'PUT'},
                'gif_url': {'Value': result_download_url, 'Action': 'PUT'},
            }
        )
    return "OK"
Сгенерируйте файл requirements.txt:
pipreqs $PWD --force
Находясь в директории с исходными файлами, упакуем все нужные файлы в ZIP-архив.
zip src.zip index.py requirements.txt ffmpeg
В Object Storage для простоты используем тот же бакет, куда далее будем складывать видео. На вкладке Объекты, вверху справа нажмите кнопку Загрузить и выберите созданный архив.
Создадим функции ffmpeg-api и ffmpeg-converter, при этом сразу зададим все необходимые переменные и сервисный аккаунт:
yc serverless function create \
  --name ffmpeg-api \
  --description "function for ffmpeg-api"

yc serverless function create \
  --name ffmpeg-converter \
  --description "function for ffmpeg-converter"

yc serverless function version create \
  --function-name ffmpeg-api \
  --memory=256m \
  --execution-timeout=5s \
  --runtime=python37 \
  --entrypoint=index.handle_api \
  --service-account-id $SERVICE_ACCOUNT_FFMPEG_ID \
  --environment SECRET_ID=$SECRET_ID \
  --environment YMQ_QUEUE_URL=$YMQ_QUEUE_URL \
  --environment DOCAPI_ENDPOINT=$DOCAPI_ENDPOINT \
  --package-bucket-name $S3_BUCKET \
  --package-object-name src.zip

yc serverless function version create \
  --function-name ffmpeg-converter \
  --memory=2048m \
  --execution-timeout=600s \
  --runtime=python37 \
  --entrypoint=index.handle_process_event \
  --service-account-id $SERVICE_ACCOUNT_FFMPEG_ID \
  --environment SECRET_ID=$SECRET_ID \
  --environment YMQ_QUEUE_URL=$YMQ_QUEUE_URL \
  --environment DOCAPI_ENDPOINT=$DOCAPI_ENDPOINT \
  --environment S3_BUCKET=$S3_BUCKET \
  --package-bucket-name $S3_BUCKET \
  --package-object-name src.zip

Тестирование функции

В консоли управления перейдите из рабочего каталога в раздел Cloud Functions и выберите ранее созданную функцию ffmpeg-api. Перейдите на вкладку Тестирование в боковом меню, выберите шаблон данных Без шаблона и добавьте во вводные данные JSON:
{"action":"convert", "src_url":"https://disk.yandex.ru/i/38RbVC0spb_jQQ"}
Нажмите кнопку Запустить тест. Этим самым мы загрузим файл в хранилище и создадим задачу в БД. Если всё сделано правильно, то вы увидите такой результат:
    {
        "task_id": "133e05c2-1b98-41cc-9aab-b816d71af343"
    } 
image
Воспользуемся полученным идентификатором задачи task_id для получения статуса из базы данных. Для этого внесите в вводные данные JSON следующие изменения:
{"action":"get_task_status", "task_id":"<идентификатор задачи>"}
Нажмите кнопку Запустить тест. Так как мы ещё не обрабатывали задачи в очереди, результат очевиден:
    {
        "ready": false
    } 
image

Шаг 6. Создание триггера

Теперь создайте триггер, который будет вызывать функцию обработки сообщений из очереди. После создания триггер начинает работать через пять минут. Он будет брать по одному сообщению и раз в 10 секунд отправлять в функцию:
yc serverless trigger create message-queue \
  --name ffmpeg \
  --queue $YMQ_QUEUE_ARN \
  --queue-service-account-id $SERVICE_ACCOUNT_FFMPEG_ID \
  --invoke-function-name ffmpeg-converter  \
  --invoke-function-service-account-id $SERVICE_ACCOUNT_FFMPEG_ID \
  --batch-size 1 \
  --batch-cutoff 10s
С этого момента очередь начнёт обрабатываться. Можно проверить, готова ли задача, и, если это так, запросить по URL результат обработки из Object Storage.
Теперь у нас есть функция, которая выполняет функцию API, через которую мы можем ставить задачи в очередь на исполнение. Триггер раз в 10 секунд берет по одному сообщению в очереди и передает функции обработчику. Функция-обработчик формирует результат и обновляет данные в базе данных. При этом мы получаем сконвертированные GIF-файлы из видео.
Протестируйте систему, используя полученный ранее идентификатор задачи task_id для получения статуса из базы данных. Для этого внесите изменения в вводные данные JSON:
{"action":"get_task_status", "task_id":"133e05c2-1b98-41cc-9aab-b816d71af343"}
Нажмите кнопку Запустить тест. Если задача уже успела обработаться, то вы получите URL.
image

Удаление триггера

По завершении работы не забудьте удалить созданный триггер ffmpeg, иначе он будет продолжать работать:
yc serverless trigger delete ffmpeg
Не забудьте также удалить или остановить все созданные вами ресурсы.
 
Проверьте себя
Очереди FIFO в Yandex Message Queue поддерживают:
 
Правильный ответ:
  • Любое количество поставщиков и потребителей
Стандартные очереди в Yandex Message Queue:
 
Правильный ответ:
  • Пытаются сохранять порядок полученных сообщений при передаче поставщикам, но не гарантируют его

 

Категория: Cloud Services Engineer
Просмотров: 928

YMQ. ПР. Проверка доступности веб-ресурсов

ПР. Проверка доступности веб-ресурсов

Кратко:

  • Создание системы проверки доступности веб-ресурсов с использованием Yandex Cloud Functions, API Gateway и Yandex Message Queue.
  • Добавление возможности ставить задачи по проверке доступности других веб-ресурсов.
  • Использование библиотеки boto3 для работы с YMQ.
  • Создание очереди Yandex Message Queue и функции для проверки доступности URL.
  • Обновление спецификации API Gateway для предоставления доступа к функции.
  • Создание функции для чтения из очереди и проверка ее работы.
  • Создание триггера для вызова функции обработки сообщений из очереди один раз в минуту.
  • Удаление триггера-таймера после завершения практической работы.
  • Не забудьте удалить или остановить все созданные вами ресурсы.

Практическая работа. Проверка доступности веб-ресурсов

В этом уроке вы доработаете систему проверки доступности веб-ресурсов, которую создали на предыдущих практических занятиях. В текущем варианте она проверяет только доступность сайта ya.ru. Теперь давайте добавим в неё возможность ставить задачи по проверке доступности других веб-ресурсов.

Общая архитектура системы

У системы есть два метода:
  1. CheckUrl — ставит задачу на проверку указанного URL.
  2. GetResult — считывает результаты проверки.
Метод CheckUrl обрабатывается функцией, которая будет складывать все запросы в очередь. Функция-обработчик будет вызываться раз в секунду, считывать URL из очереди,  проверять его доступность и записывать результат в базу данных. Оттуда этот результат можно будет получить с помощью метода GetResult.
image
Мы не будем менять уже созданные функции и таблицу в PostgreSQL, сделаем новые.
Работать с YMQ из функций мы будем с помощью библиотеки boto3. Чтобы её использовать, нужно создать сервисный аккаунт с секретным ключом доступа, а затем настроить зависимости функции. Сделаем это после того, как создадим очередь.

Шаг 1. Проверить наличие сервисного аккаунта

Если вы ранее создавали сервисный аккаунт с именем service-account-for-cf, добавляли вновь созданному сервисному аккаунту роли editor и другие, то вам остаётся только создать ключ доступа:
yc iam access-key create --service-account-name service-account-for-cf
В результате вы получите примерно следующее:
    access_key:
        id: ajefraollq5puj2tir1o
        service_account_id: ajetdv28pl0a1a8r41f0
        created_at: "2021-08-23T21:13:05.677319393Z"
        key_id: BTPNvWthv0ZX2xVmlPIU
    secret: cWLQ0HrTM0k_qAac43cwMNJA8VV_rfTg_kd4xVPi 
Здесь key_id — это идентификатор ключа доступа ACCESS_KEY. А secret — это секретный ключ SECRET_KEY. Переменные ACCESS_KEY и SECRET_KEY могут быть использованы для задания соответствующих значений aws_access_key_id и aws_secret_access_key при использовании библиотеки boto3.

Шаг 2. Создание очереди Yandex Message Queue

Вы можете создать очередь одним из трёх способов:
  • через консоль управления;
  • с помощью консольной утилиты aws;
  • с помощью Terraform.
В этом уроке мы будем использовать консоль управления. Откройте раздел Message Queue и нажмите кнопку Создать очередь.
image
В настройках создаваемой очереди задайте имя очереди my-first-queue, затем выберите тип очереди Стандартная и нажмите кнопку Создать.
image
Очередь создана.
image
Теперь зайдите в настройки очереди, чтобы посмотреть параметры подключения к ней. Нам потребуется значение URL.
image

Шаг 3. Создание функции

Для создания функции зададим ряд переменных:
  • VERBOSE_LOG — определяет, пишет ли функция подробности своего выполнения в журнал.
  • AWS_ACCESS_KEY_ID — значение «Идентификатор ключа» из сервисного аккаунта, который мы сделали ранее.
  • AWS_SECRET_ACCESS_KEY — значение «Секретный ключ» из того же сервисного аккаунта.
  • QUEUE_URL — URL на очередь, его можно получить на обзорной странице созданной ранее очереди.
Чтобы задать переменные, выполните в консоли следующие команды:
echo "export VERBOSE_LOG=True" >> ~/.bashrc && . ~/.bashrc
echo "export AWS_ACCESS_KEY_ID=<AWS_ACCESS_KEY_ID>" >> ~/.bashrc && . ~/.bashrc
echo "export AWS_SECRET_ACCESS_KEY=<AWS_SECRET_ACCESS_KEY>" >> ~/.bashrc && . ~/.bashrc
echo "export QUEUE_URL=<QUEUE_URL>" >> ~/.bashrc && . ~/.bashrc
Создайте файл my-url-receiver-function.py со следующим содержанием:
import logging
import os
import boto3

logger = logging.getLogger()
logger.setLevel(logging.INFO)

verboseLogging = eval(os.environ['VERBOSE_LOG'])  ## Convert to bool
queue_url = os.environ['QUEUE_URL']

def log(logString):
    if verboseLogging:
        logger.info(logString)

def handler(event, context):

    # Get url
    try:
        url = event['queryStringParameters']['url']
    except Exception as error:
        logger.error(error)
        statusCode = 400
        return {
            'statusCode': statusCode
        }

    # Create client
    client = boto3.client(
        service_name='sqs',
        endpoint_url='https://message-queue.api.cloud.yandex.net',
        region_name='ru-central1'
    )

    # Send message to queue
    client.send_message(
        QueueUrl=queue_url,
        MessageBody=url
    )
    log('Successfully sent test message to queue')

    statusCode = 200

    return {
        'statusCode': statusCode
    }
Затем воспользуйтесь командой pipreqs $PWD --force для формирования файла requirements.txt и упакуйте файлы с функцией и требованиями в ZIP-архив.
zip my-url-receiver-function my-url-receiver-function.py requirements.txt
Создайте функцию и её версию:
yc serverless function create \
  --name  my-url-receiver-function \
  --description "function for url"

yc serverless function version create \
  --function-name=my-url-receiver-function \
  --memory=256m \
  --execution-timeout=5s \
  --runtime=python312 \
  --entrypoint=my-url-receiver-function.handler \
  --service-account-id $SERVICE_ACCOUNT_ID \
  --environment VERBOSE_LOG=$VERBOSE_LOG \
  --environment AWS_ACCESS_KEY_ID=$AWS_ACCESS_KEY_ID \
  --environment AWS_SECRET_ACCESS_KEY=$AWS_SECRET_ACCESS_KEY \
  --environment QUEUE_URL=$QUEUE_URL \
  --source-path my-url-receiver-function.zip

Тестирование функции

Перейдите в раздел Cloud Functions консоли управления облаком и выберите созданную функцию my-url-receiver-function. На вкладке Тестирование в боковом меню выберите шаблон HTTPS-вызов и замените раздел queryStringParameters:
    "queryStringParameters": {
        "a": "2",
        "b": "1",
    }, 
на аналогичный, но с параметром url с любым сайтом. Важно указывать ссылку целиком.
    "queryStringParameters": {
        "url": "https://ya.ru/"
    }, 
Нажмите кнопку Запустить тест.
image
Если вы всё сделали правильно, то увидите код статуса 200. При этом в очереди увеличится количество сообщений.
image

Шаг 4. Обновление спецификации API Gateway

Функция готова, но по умолчанию она не является публичной. Предоставим доступ к ней с помощью API-шлюза. Для этого необходимо обновить ранее созданную спецификацию hello-world.yaml. Если у вас нет её под рукой, выгрузите её из облака:
yc serverless api-gateway get-spec \
  --name hello-world >> hello-world-new.yaml
Внесите изменения, добавив секцию о ранее созданной функции:
    /check:
        get:
            x-yc-apigateway-integration:
                type: cloud-functions
                function_id: <идентификатор функции>
                service_account_id: <идентификатор сервисного аккаунта>
            operationId: add-url 
Обновите конфигурацию:
yc serverless api-gateway update \
  --name hello-world \
  --spec=hello-world-new.yaml
Для тестирования выполните вызов функции в браузере:
https://<идентификатор API Gateway>.apigw.yandexcloud.net/check?url=https://ya.ru/
После каждого запроса количество сообщений в очереди будет увеличиваться на одно.
image

Шаг 5. Создание функции для чтения из очереди

В предыдущих работах мы создавали функцию, использующую подключение к БД. Здесь мы повторим этот опыт.
Проверим, что нам доступны переменные для инициации подключения: CONNECTION_ID, DB_USER, DB_HOST. Мы создавали их в предыдущей работе с помощью следующих команд:
echo "export CONNECTION_ID=<CONNECTION_ID>" >> ~/.bashrc && . ~/.bashrc
echo "export DB_USER=<DB_USER>" >> ~/.bashrc && . ~/.bashrc
echo "export DB_HOST=<DB_HOST>" >> ~/.bashrc && . ~/.bashrc
Также для работы с очередью нам потребуются переменные VERBOSE_LOG, AWS_ACCESS_KEY_ID, AWS_SECRET_ACCESS_KEY и QUEUE_URL, заданные на предыдущих шагах.
Создадим функцию function-for-url-from-mq.py и воспользуемся командой pipreqs $PWD --force, чтобы сформировать для нее файл requirements.txt.
import logging
import os
import boto3
import datetime
import requests

#Эти библиотеки нужны для работы с PostgreSQL
import psycopg2
import psycopg2.errors
import psycopg2.extras

CONNECTION_ID = os.getenv("CONNECTION_ID")
DB_USER = os.getenv("DB_USER")
DB_HOST = os.getenv("DB_HOST")
QUEUE_URL = os.environ['QUEUE_URL']

# Настраиваем функцию для записи информации в журнал функции
# Получаем стандартный логер языка Python
logger = logging.getLogger()
logger.setLevel(logging.INFO)
# Вычитываем переменную VERBOSE_LOG, которую мы указываем в переменных окружения 
verboseLogging = eval(os.environ['VERBOSE_LOG'])  ## Convert to bool

#Функция log, которая запишет текст в журнал выполнения функции, если в переменной окружения VERBOSE_LOG будет значение True
def log(logString):
    if verboseLogging:
        logger.info(logString)

#Получаем подключение
def getConnString(context):
    """
    Extract env variables to connect to DB and return a db string
    Raise an error if the env variables are not set
    :return: string
    """
    connection = psycopg2.connect(
        database=CONNECTION_ID, # Идентификатор подключения
        user=DB_USER, # Пользователь БД
        password=context.token["access_token"],
        host=DB_HOST, # Точка входа
        port=6432,
        sslmode="require")
    return connection

"""
    Create SQL query with table creation
"""
def makeCreateDataTableQuery(table_name):
    query = f"""CREATE TABLE public.{table_name} (
    url text,
    result integer,
    time float
    )"""
    return query

def makeInsertDataQuery(table_name, url, result, time):
    query = f"""INSERT INTO {table_name} 
    (url, result,time)
    VALUES('{url}', {result}, {time})
    """
    return query

def handler(event, context):

    # Create client
    client = boto3.client(
        service_name='sqs',
        endpoint_url='https://message-queue.api.cloud.yandex.net',
        region_name='ru-central1'
    )

    # Receive sent message
    messages = client.receive_message(
        QueueUrl=QUEUE_URL,
        MaxNumberOfMessages=1,
        VisibilityTimeout=60,
        WaitTimeSeconds=1
    ).get('Messages')

    if messages is None:
        return {
            'statusCode': 200
        }

    for msg in messages:
        log('Received message: "{}"'.format(msg.get('Body')))

    # Get url from message
    url = msg.get('Body');

    # Check url
    try:
        now = datetime.datetime.now()
        response = requests.get(url, timeout=(1.0000, 3.0000))
        timediff = datetime.datetime.now() - now
        result = response.status_code
    except requests.exceptions.ReadTimeout:
        result = 601
    except requests.exceptions.ConnectTimeout:
        result = 602
    except requests.exceptions.Timeout:
        result = 603
    log(f'Result: {result} Time: {timediff.total_seconds()}')
    
    connection = getConnString(context)
    log(f'Connecting: {connection}')    
    cursor = connection.cursor()

    table_name = 'custom_request_result'
    sql = makeInsertDataQuery(table_name, url, result, timediff.total_seconds())

    log(f'Exec: {sql}')
    try:
        cursor.execute(sql)
    except psycopg2.errors.UndefinedTable as error:
        log(f'Table not exist - create and repeate insert')
        connection.rollback()
        logger.error(error)
        createTable = makeCreateDataTableQuery(table_name)
        log(f'Exec: {createTable}')
        cursor.execute(createTable)
        connection.commit()
        log(f'Exec: {sql}')
        cursor.execute(sql)
    except Exception as error:
        logger.error( error)

    connection.commit()
    cursor.close()
    connection.close()

    # Delete processed messages
    for msg in messages:
        client.delete_message(
            QueueUrl=QUEUE_URL,
            ReceiptHandle=msg.get('ReceiptHandle')
        )
        print('Successfully deleted message by receipt handle "{}"'.format(msg.get('ReceiptHandle')))

    statusCode = 200

    return {
        'statusCode': statusCode
    }
При создании сразу задайте все необходимые переменные и сервисный аккаунт:
zip function-for-url-from-mq function-for-url-from-mq.py requirements.txt

yc serverless function create \
  --name function-for-url-from-mq \
  --description "function for url from mq"

yc serverless function version create \
  --function-name=function-for-url-from-mq \
  --memory=256m \
  --execution-timeout=5s \
  --runtime=python312 \
  --entrypoint=function-for-url-from-mq.handler \
  --service-account-id $SERVICE_ACCOUNT_ID \
  --environment VERBOSE_LOG=True \
  --environment CONNECTION_ID=$CONNECTION_ID \
  --environment DB_USER=$DB_USER \
  --environment DB_HOST=$DB_HOST \
  --environment AWS_ACCESS_KEY_ID=$AWS_ACCESS_KEY_ID \
  --environment AWS_SECRET_ACCESS_KEY=$AWS_SECRET_ACCESS_KEY \
  --environment QUEUE_URL=$QUEUE_URL \
  --source-path function-for-url-from-mq.zip
Протестируйте функцию.
image
После её выполнения количество сообщений в очереди уменьшится, а в базе данных появится новая таблица с результатами тестирования доступности функции.
image
image

Шаг 6. Создание триггера

Создадим триггер, который будет вызывать функцию обработки сообщений из очереди один раз в минуту. Он будет использовать cron-выражение:
yc serverless trigger create timer \
  --name trigger-for-mq \
  --invoke-function-name function-for-url-from-mq \
  --invoke-function-service-account-id $SERVICE_ACCOUNT_ID \
  --cron-expression '* * * * ? *'
Cron-выражение * * * * ? * означает вызов функции function-for-url-from-mq один раз в минуту. Подробнее про cron-выражения можно прочитать в документации.
image
Теперь у нас есть функция, которая раз в минуту будет пробовать взять из очереди URL и проверить его. Также есть метод REST API, который позволяет записывать URL в очередь независимо от работы обработчика. Мы можем вызывать созданный метод как угодно часто. Очередь будет просто накапливаться, а затем обработчик будет постепенно её разбирать.
В итоге вы получили асинхронную систему проверки доступности URL с доступом по REST API. Вы не создали ни одной виртуальной машины, но решили вопросы масштабирования и отказоустойчивости системы.

Удаление триггера-таймера

По завершении практической работы не забудьте удалить созданный вами триггер trigger-for-mq, иначе он будет работать, пока не исчерпает деньги на аккаунте:
yc serverless trigger delete trigger-for-mq
Не забудьте удалить или остановить все созданные вами ресурсы: триггеры, очереди YMQ и кластер базы данных.
Следующий практический урок завершает тему. Вы попробуете создать онлайн-сервис, конвертирующий произвольные видеофайлы в GIF-анимацию. Для этого вы объедините в одно решение сервисы Yandex Cloud Functions, Yandex Message Queue, YDB и Yandex Object Storage. А заодно закрепите использование консольных инструментов yc и aws.

 

Категория: Cloud Services Engineer
Просмотров: 890

YMQ. Знакомство с Yandex Message Queue

Знакомство с Yandex Message Queue

Кратко:

  • Очереди в Yandex Message Queue состоят из тела и метаданных сообщения.
  • Тело обрабатывается приложением, метаданные используются для подтверждения обработки, оценки времени и проверки передачи сообщения.
  • Сервис YMQ организует очередь между отправителями и получателями сообщений.
  • Существуют стандартные и FIFO очереди, отличающиеся порядком обработки и гарантией доставки сообщений.
  • Параметры очередей включают стандартный таймаут видимости, срок хранения сообщений, максимальный размер сообщения и задержку доставки.
  • Dead Letter Queue используется для обработки недоставленных сообщений.
  • Работа с очередями осуществляется через API или AWS CLI.
  • Тарификация зависит от типа очереди и исходящего трафика.

Знакомство с Yandex Message Queue

На прошлом уроке вы узнали, что такое очереди и зачем они нужны. На этом уроке вы разберётесь со спецификой реализации очередей в Yandex Message Queue и их параметрами.

Что такое сообщение

Основным понятием для очередей в Yandex Message Queue (YMQ) является сообщение. Сообщение состоит из тела — ваших данных — и метаданных, его дополнительных атрибутов.
Тело сообщения обрабатывается вашим приложением, а метаданные удобно использовать для других самых разнообразных целей: например, для подтверждения обработки сообщения, оценки времени его обработки или чтобы удостовериться, что оно было передано без изменений.
image
Представим, что мы считали это простое сообщение из очереди, давайте рассмотрим его содержимое:
  • Body — т. е. тело сообщения. Здесь хранится то, что собственно передаётся.
  • Attributes — набор атрибутов сообщения, указывающих время первого получения (время UNIX), количество попыток обработки этого сообщения, время отправки (время UNIX).
  • ReceiptHandle — идентификатор получения. Он указывает на факт получения сообщения и назначается системой при его считывании. Этот идентификатор используется для удаления полученного сообщения из очереди или изменения его таймаута видимости.
  • MD5OfBody — хеш-сумма тела сообщения, созданная 128-битным алгоритмом хеширования MD5.
  • MessageId — уникальный идентификатор сообщения. Он возвращается вам из YMQ, когда вы отправляете сообщение. Этот идентификатор удобно использовать для различной диагностики.

Как работает сервис YMQ

Задача YMQ — организовать очередь между приложениями-отправителями и приложениями получателями сообщений. Чтобы отправить и получить сообщение, отправители и получатели должны обращаться к сервису сами.
Каждое сообщение за время жизни проходит следующие этапы:
  • отправка в очередь;
  • хранение в очереди, пока оно не будет считано или пока не истечёт время хранения;
  • чтение сообщения потребителем и пометка сообщения на это время, как находящегося в обработке;
  • удаление из очереди, если сообщение было успешно обработано или перенесено в Dead Letter Queue:
image

Типы очередей в YMQ

Сервис YMQ поддерживает два типа очередей — стандартные и FIFO. На предыдущем уроке мы уже разобрали их основные отличия. Давайте остановимся теперь на них подробнее.
Стандартные очереди позволяют сохранять сообщения, которые затем читаются приложениями в произвольном порядке. Такой подход упрощает систему обработки сообщений. Приложения должны быть рассчитаны на ситуации приёма данных не в хронологическом порядке их поступления.
Стандартные очереди обеспечивают гарантию, что каждое сообщение будет доставлено до получателя хотя бы один раз. В исключительных случаях данные могут быть доставлены до считывающего приложения несколько раз. Обрабатывающие системы должны быть готовы к подобным ситуациям. Такие очереди лучше подходят для обработки не связанных между собой сообщений и обеспечивают более высокую пропускную способность, то есть работают быстрее.
Очереди FIFO позволяют обеспечить строгую очерёдность выдачи сообщений запрашивающей/обрабатывающей стороне и обеспечивают семантику строгой однократной гарантированной доставки сообщений. Такие очереди подходят для передачи связанных сообщений и работают медленнее из-за того, что сообщения должны быть обработаны по очереди.
FIFO очереди часто используются для обработки финансовых данных. Представьте, что в одну FIFO очередь отправляются действия с банковскими счетами разных пользователей. Данные из очереди обрабатывают несколько получателей, и было бы удобно, чтобы все действия одного пользователя попадали одному получателю. Для такой цели можно использовать группировку сообщений.
При помощи специальных идентификаторов группы можно обеспечить отправку сразу нескольких потоков упорядоченных сообщений для разных получателей сообщений в рамках одной очереди FIFO. Вместо одной очереди FIFO получается несколько очередей по количеству групп. В каждой группе запись сообщений и их считывание происходит по схеме FIFO.

Параметры очередей

Давайте разберём, какие параметры бывают у очередей в Yandex Message Queue.
image
В Базовых параметрах помимо имени и типа очереди вы можете указать:
  • Стандартный таймаут видимости. Это время, на которое сообщение скрывается из очереди после чтения получателем. Пока сообщение скрыто, другие получатели не могут получить сообщение из очереди. Минимальный таймаут видимости — 30 секунд, максимальный — 12 часов.
  • Срок хранения сообщений. Вы можете указать, как долго каждое сообщение может храниться в очереди в ожидании чтения получателем. Это значение должно быть в промежутке от 60 секунд до 14 дней.
  • Максимальный размер сообщения. Может составлять от 1 до 256 КБ.
  • Задержка доставки. Иногда нужно, чтобы обработчик получил сообщение не сразу, а позже. Здесь вы можете указать время, в течение которого новое сообщение нельзя получить из очереди. Значение должно быть в промежутке от 0 секунд до 15 минут.
  • Время ожидания при получении сообщения. В течение этого времени получатель будет ожидать поступления сообщений. Если в очереди появятся сообщения, вызов будет сделан раньше, чем указано в этой настройке. Если же по истечении этого времени сообщения не появились, будет возвращен пустой список.
image
В блоке Настройки очередей недоставленных сообщений вы можете настроить работу с так называемой Dead Letter Queue (DLQ, дословно — очередь невостребованных писем). Это специальная очередь, куда могут перенаправляться сообщения, которые получатели не смогли обработать в обычных очередях. Собирая такие сообщения в отдельной очереди, вы можете исследовать ошибки, возникающие при их обработке.
Чтобы воспользоваться этой функцией, вам придется сначала завести отдельную очередь того же типа, что и очередь, откуда перенаправляются сообщения, которые дошли до адресата, но не были обработаны. Включите функцию Перенаправлять недоставленные сообщения, выберите заранее созданную DLQ, а затем укажите количество попыток, после которых необработанное сообщение направляется в эту очередь.

Как можно работать с очередями в Yandex Message Queue

Помещение данных в очередь выполняется при помощи программных вызовов через специальный API, либо с использованием AWS CLI. Пример использования AWS CLI с YMQ вы можете найти в документации.
YMQ поддерживает API и другие подходы, которые используют в сервисе Amazon SQS, поэтому для работы с ними вы можете использовать уже существующие инструменты, например библиотеки boto3 для Python .

Тарификация

В рамках сервиса Message Queue тарифицируется количество запросов к стандартным очередям и очередям FIFO, а также исходящий трафик. Для целей тарификации каждые 64 КБ данных запроса считаются отдельным запросом.
Первые 100 000 запросов в месяц к очередям любого типа не оплачиваются. А далее их стоимость зависит от типа очереди: для стандартных очередей это 48,7600 ₽ за 1 миллион запросов, для FIFO — 61,1500 ₽.
Исходящий трафик до 10 ГБ не тарифицируется, а затем оплата за него  составляет 1,5300 ₽ за 1 ГБ.
Допустим, за месяц было сделано 2,75 млн запросов объемом 48 КБ к стандартной очереди. Это значит, что нам нужно вычесть нетарифицируемое количество запросов из общего и привести его к тарифу за миллион запросов, а затем приплюсовать исходящий трафик исходя из того, что каждое сообщение умещается в 64 КБ, приведя размер каждого сообщения к гигабайтам. Считаем:
48,7600×((2750000−100000)/1000000)+1,53×((2750000×48/1024/1024)−10)=129,214+177,30=306,51448,7600×((2750000−100000)/1000000)+1,53×((2750000×48/1024/1024)−10)=129,214+177,30=306,514 ₽
Более подробную информацию вы найдёте в документации.
На следующем уроке вы расширите созданное ранее приложение для проверки доступности yandex.ru, добавив в него возможность ставить задачи по проверке доступности других веб-ресурсов при помощи очередей.
 
Проверьте себя
Что происходит с сообщением, когда получатель начал обрабатывать сообщение?
 
Правильный ответ:
  • Остаётся в очереди, но скрывается от других
Какие типы очередей поддерживает сервис YMQ?
 
Правильный ответ:
  • Оба типа 

 

Категория: Cloud Services Engineer
Просмотров: 898

YMQ. Что такое очереди

Что такое очереди

Кратко:

  • Очереди используются для организации независимой работы поставщиков и потребителей информации в системах реального времени.
  • Они служат буфером между поставщиком и потребителем, позволяя развязать два параллельных процесса.
  • Очереди используются внутри программ при взаимодействии потоков и в организации больших систем.
  • Они помогают решать проблемы интеграции приложений, их отказоустойчивости и масштабирования.
  • Существует два распространённых типа очередей: FIFO и стандартная.
  • У очередей FIFO сравнительно небольшое количество вызовов API в секунду.
  • В свою очередь из стандартных очередей сообщения по возможности считываются последовательно.
  • Преимуществом этого подхода является более высокая пропускная способность.

Что такое очереди

Зачем нужны очереди

Предположим, вы создаете поисковую систему. У вас есть роботы (поставщики), которые собирают ссылки и передают их обработчикам (потребителям) для разбора страниц и записи результата в базу данных. Самый простой способ передавать ссылки от роботов к обработчикам — вызывать обработчики напрямую из роботов. Однако, такая реализация имеет ряд недочётов, например:
  • Как правило, роботы работают быстрее обработчиков: данные будут накапливаться внутри роботов, а значит, их память или диск будут быстро переполняться.
  • Роботам нужно знать рабочий интерфейс обработчиков, который может меняться по мере совершенствования системы, а значит, будет нужно адаптировать и код роботов.
  • Нет гарантии, что обработчик разберет ссылку целиком.
Практически все эти проблемы решаются при помощи очередей, т. е. последовательности некоторой информации. Главная задача очередей — организовать независимую работу поставщиков и потребителей информации в системах реального времени. Вот что это означает в нашем примере.
Робот пишет сообщение в очередь: «У меня есть новая ссылка, вот она». Обработчик ссылок читает записанное сообщение из очереди и забирает себе ссылку на разбор. Чтобы другой обработчик не начал обрабатывать ту же самую ссылку, очередь скрывает это сообщение на некоторый промежуток времени (таймаут видимости), достаточный для его обработки потребителем.
Если сообщение обработано успешно, оно удаляется из очереди, а обработчик забирает себе следующее. Если обработчик вышел из строя или произошёл обрыв соединения, по истечении таймаута видимости сообщение снова становится доступным в очереди, и его может взять другой обработчик.
Таким образом, очередь служит буфером между поставщиком и потребителем, позволяя развязать два параллельных процесса. Если потребитель начинает медленнее вычитывать элементы, то очередь будет их накапливать. А шанс догнать процесс и обработать данные побыстрее даст потом.
Очереди используются внутри программ при взаимодействии потоков и в организации больших систем, где взаимодействуют несколько программ или сервисов. В таких системах очереди помогают решать проблемы интеграции приложений, их отказоустойчивости и масштабирования, например:
  • надёжной передачи данных и команд между компонентами;
  • обработки данных от множества устройств IoT;
  • обработки событий, по которым должна быть вызвана функция.
Стоит учитывать, что у очередей есть ряд ограничений. Главным из них является неопределённое время на обработку принимающей стороной. Сложно предсказать, когда потребитель получит и обработает элемент.

Виды очередей

Существует два распространённых типа очередей: FIFO и стандартная.
FIFO означает «first in, first out», т. е. порядок размещения данных в очереди совпадает с порядком, в котором данные из нее считываются. Очереди FIFO используются в тех случаях, когда нужно обеспечить строгий порядок доставки и однократную обработку сообщений. Например, соблюдение исходного порядка важно для финансовых транзакций. Если представить, что у клиента на счету было 2 000 рублей, и он сначала положил туда 10 000, а затем потратил 5 000, то очевидно, что порядок выполнения транзакций имеет значение. Стоит отметить, что у очередей FIFO сравнительно небольшое количество вызовов API в секунду.
image
В свою очередь из стандартных очередей сообщения по возможности считываются последовательно, но соблюдение исходного порядка при доставке сообщений не гарантируется. Преимуществом этого подхода является более высокая в сравнении с FIFO пропускная способность: стандартные очереди поддерживают сравнительно большое количество вызовов API в секунду (отправка, принятие или удаление сообщения).
image
В Yandex Cloud сервисом очередей является Yandex Message Queue. Он относится к PaaS-слою и объединён с другими serverless-сервисами в группу Бессерверные вычисления. В следующем уроке мы поговорим о его специфике.
 
Проверьте себя
Как называется тип очередей, который сохраняет последовательность элементов?
 
Правильный ответ:
  • FIFO
В каких задачах полезны очереди?
 
Правильный ответ:
  • Передачи данных
Сколько поставщиков и потребителей может быть у очереди?
 
Правильный ответ:
  • В теории — бесконечное количество. На практике — зависит от конкретной очереди

 

Категория: Cloud Services Engineer
Просмотров: 597
  1. SY. ПР. Запуск тестового приложения
  2. SY. ПР. Загрузка данных, выполнение запросов AWS CLI
  3. SY. Document API
  4. SY. Тарификация YDB в бессерверном режиме

Страница 5 из 19

  • 1
  • 2
  • 3
  • 4
  • 5
  • 6
  • 7
  • 8
  • 9
  • 10
© Gantry Framework 2016 - 2026
Developed by RocketTheme exclusively
for Gantry 5.
  • Главная
  • Начало
  • Карта
Back to top