ПР. Сокращатель ссылок
Кратко:
- Создание сокращателя ссылок с использованием сервисов 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. Нажмите кнопку Создать бакет.
На странице создания бакета:
- Введите имя бакета. В нашем примере это будет
storage-for-serverless-shortener. - Ограничьте максимальный размер бакета (например 1 ГБ).
- Выберите тип доступа
Публичныйво всех случаях. - Выберите класс хранилища
Стандартное.
Нажмите кнопку Создать бакет для завершения операции.

Создайте файл
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
- Сохраним в файл
links.yqlSQL-скрипт для создания таблицы: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.

Итак, вы создали приложение с использованием Cloud Functions, API Gateway, Object Storage и YDB. Конечно, вы можете развивать его и дальше, расширяя функциональность.
Вводный курс по serverless-разработке на этом завершён. Осталось пройти лишь заключительный тест, который проверит ваши знания по всем рассмотренным в курсе сервисам.
- Категория: Cloud Services Engineer
- Просмотров: 1208
ПР. Однократная отправка сообщений
Кратко:
- Реализация проекта по конвертированию видеофайлов в 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. Сервисный аккаунт и 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"
}

Воспользуемся полученным идентификатором задачи
task_id для получения статуса из базы данных. Для этого внесите в вводные данные JSON следующие изменения:
{"action":"get_task_status", "task_id":"<идентификатор задачи>"}
Нажмите кнопку Запустить тест. Так как мы ещё не обрабатывали задачи в очереди, результат очевиден:
{
"ready": false
}

Шаг 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.

Удаление триггера
По завершении работы не забудьте удалить созданный триггер
ffmpeg, иначе он будет продолжать работать:
yc serverless trigger delete ffmpeg
Не забудьте также удалить или остановить все созданные вами ресурсы.
- Категория: Cloud Services Engineer
- Просмотров: 928
ПР. Проверка доступности веб-ресурсов
Кратко:
- Создание системы проверки доступности веб-ресурсов с использованием Yandex Cloud Functions, API Gateway и Yandex Message Queue.
- Добавление возможности ставить задачи по проверке доступности других веб-ресурсов.
- Использование библиотеки boto3 для работы с YMQ.
- Создание очереди Yandex Message Queue и функции для проверки доступности URL.
- Обновление спецификации API Gateway для предоставления доступа к функции.
- Создание функции для чтения из очереди и проверка ее работы.
- Создание триггера для вызова функции обработки сообщений из очереди один раз в минуту.
- Удаление триггера-таймера после завершения практической работы.
- Не забудьте удалить или остановить все созданные вами ресурсы.
Практическая работа. Проверка доступности веб-ресурсов
В этом уроке вы доработаете систему проверки доступности веб-ресурсов, которую создали на предыдущих практических занятиях. В текущем варианте она проверяет только доступность сайта ya.ru. Теперь давайте добавим в неё возможность ставить задачи по проверке доступности других веб-ресурсов.
Общая архитектура системы
У системы есть два метода:
CheckUrl— ставит задачу на проверку указанного URL.GetResult— считывает результаты проверки.
Метод
CheckUrl обрабатывается функцией, которая будет складывать все запросы в очередь. Функция-обработчик будет вызываться раз в секунду, считывать URL из очереди, проверять его доступность и записывать результат в базу данных. Оттуда этот результат можно будет получить с помощью метода GetResult.
Мы не будем менять уже созданные функции и таблицу в 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 и нажмите кнопку Создать очередь.

В настройках создаваемой очереди задайте имя очереди
my-first-queue, затем выберите тип очереди Стандартная и нажмите кнопку Создать.
Очередь создана.

Теперь зайдите в настройки очереди, чтобы посмотреть параметры подключения к ней. Нам потребуется значение URL.

Шаг 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/"
},
Нажмите кнопку Запустить тест.

Если вы всё сделали правильно, то увидите код статуса
200. При этом в очереди увеличится количество сообщений.
Шаг 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/
После каждого запроса количество сообщений в очереди будет увеличиваться на одно.

Шаг 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
Протестируйте функцию.

После её выполнения количество сообщений в очереди уменьшится, а в базе данных появится новая таблица с результатами тестирования доступности функции.


Шаг 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-выражения можно прочитать в документации.
Теперь у нас есть функция, которая раз в минуту будет пробовать взять из очереди 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
Знакомство с Yandex Message Queue
Кратко:
- Очереди в Yandex Message Queue состоят из тела и метаданных сообщения.
- Тело обрабатывается приложением, метаданные используются для подтверждения обработки, оценки времени и проверки передачи сообщения.
- Сервис YMQ организует очередь между отправителями и получателями сообщений.
- Существуют стандартные и FIFO очереди, отличающиеся порядком обработки и гарантией доставки сообщений.
- Параметры очередей включают стандартный таймаут видимости, срок хранения сообщений, максимальный размер сообщения и задержку доставки.
- Dead Letter Queue используется для обработки недоставленных сообщений.
- Работа с очередями осуществляется через API или AWS CLI.
- Тарификация зависит от типа очереди и исходящего трафика.
Знакомство с Yandex Message Queue
На прошлом уроке вы узнали, что такое очереди и зачем они нужны. На этом уроке вы разберётесь со спецификой реализации очередей в Yandex Message Queue и их параметрами.
Что такое сообщение
Основным понятием для очередей в Yandex Message Queue (YMQ) является сообщение. Сообщение состоит из тела — ваших данных — и метаданных, его дополнительных атрибутов.
Тело сообщения обрабатывается вашим приложением, а метаданные удобно использовать для других самых разнообразных целей: например, для подтверждения обработки сообщения, оценки времени его обработки или чтобы удостовериться, что оно было передано без изменений.

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

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

В Базовых параметрах помимо имени и типа очереди вы можете указать:
- Стандартный таймаут видимости. Это время, на которое сообщение скрывается из очереди после чтения получателем. Пока сообщение скрыто, другие получатели не могут получить сообщение из очереди. Минимальный таймаут видимости — 30 секунд, максимальный — 12 часов.
- Срок хранения сообщений. Вы можете указать, как долго каждое сообщение может храниться в очереди в ожидании чтения получателем. Это значение должно быть в промежутке от 60 секунд до 14 дней.
- Максимальный размер сообщения. Может составлять от 1 до 256 КБ.
- Задержка доставки. Иногда нужно, чтобы обработчик получил сообщение не сразу, а позже. Здесь вы можете указать время, в течение которого новое сообщение нельзя получить из очереди. Значение должно быть в промежутке от 0 секунд до 15 минут.
- Время ожидания при получении сообщения. В течение этого времени получатель будет ожидать поступления сообщений. Если в очереди появятся сообщения, вызов будет сделан раньше, чем указано в этой настройке. Если же по истечении этого времени сообщения не появились, будет возвращен пустой список.

В блоке Настройки очередей недоставленных сообщений вы можете настроить работу с так называемой 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,514 ₽
Более подробную информацию вы найдёте в документации.
На следующем уроке вы расширите созданное ранее приложение для проверки доступности yandex.ru, добавив в него возможность ставить задачи по проверке доступности других веб-ресурсов при помощи очередей.
- Категория: Cloud Services Engineer
- Просмотров: 898
Что такое очереди
Кратко:
- Очереди используются для организации независимой работы поставщиков и потребителей информации в системах реального времени.
- Они служат буфером между поставщиком и потребителем, позволяя развязать два параллельных процесса.
- Очереди используются внутри программ при взаимодействии потоков и в организации больших систем.
- Они помогают решать проблемы интеграции приложений, их отказоустойчивости и масштабирования.
- Существует два распространённых типа очередей: FIFO и стандартная.
- У очередей FIFO сравнительно небольшое количество вызовов API в секунду.
- В свою очередь из стандартных очередей сообщения по возможности считываются последовательно.
- Преимуществом этого подхода является более высокая пропускная способность.
Что такое очереди
Зачем нужны очереди
Предположим, вы создаете поисковую систему. У вас есть роботы (поставщики), которые собирают ссылки и передают их обработчикам (потребителям) для разбора страниц и записи результата в базу данных. Самый простой способ передавать ссылки от роботов к обработчикам — вызывать обработчики напрямую из роботов. Однако, такая реализация имеет ряд недочётов, например:
- Как правило, роботы работают быстрее обработчиков: данные будут накапливаться внутри роботов, а значит, их память или диск будут быстро переполняться.
- Роботам нужно знать рабочий интерфейс обработчиков, который может меняться по мере совершенствования системы, а значит, будет нужно адаптировать и код роботов.
- Нет гарантии, что обработчик разберет ссылку целиком.
Практически все эти проблемы решаются при помощи очередей, т. е. последовательности некоторой информации. Главная задача очередей — организовать независимую работу поставщиков и потребителей информации в системах реального времени. Вот что это означает в нашем примере.
Робот пишет сообщение в очередь: «У меня есть новая ссылка, вот она». Обработчик ссылок читает записанное сообщение из очереди и забирает себе ссылку на разбор. Чтобы другой обработчик не начал обрабатывать ту же самую ссылку, очередь скрывает это сообщение на некоторый промежуток времени (таймаут видимости), достаточный для его обработки потребителем.
Если сообщение обработано успешно, оно удаляется из очереди, а обработчик забирает себе следующее. Если обработчик вышел из строя или произошёл обрыв соединения, по истечении таймаута видимости сообщение снова становится доступным в очереди, и его может взять другой обработчик.
Таким образом, очередь служит буфером между поставщиком и потребителем, позволяя развязать два параллельных процесса. Если потребитель начинает медленнее вычитывать элементы, то очередь будет их накапливать. А шанс догнать процесс и обработать данные побыстрее даст потом.
Очереди используются внутри программ при взаимодействии потоков и в организации больших систем, где взаимодействуют несколько программ или сервисов. В таких системах очереди помогают решать проблемы интеграции приложений, их отказоустойчивости и масштабирования, например:
- надёжной передачи данных и команд между компонентами;
- обработки данных от множества устройств IoT;
- обработки событий, по которым должна быть вызвана функция.
Стоит учитывать, что у очередей есть ряд ограничений. Главным из них является неопределённое время на обработку принимающей стороной. Сложно предсказать, когда потребитель получит и обработает элемент.
Виды очередей
Существует два распространённых типа очередей: FIFO и стандартная.
FIFO означает «first in, first out», т. е. порядок размещения данных в очереди совпадает с порядком, в котором данные из нее считываются. Очереди FIFO используются в тех случаях, когда нужно обеспечить строгий порядок доставки и однократную обработку сообщений. Например, соблюдение исходного порядка важно для финансовых транзакций. Если представить, что у клиента на счету было 2 000 рублей, и он сначала положил туда 10 000, а затем потратил 5 000, то очевидно, что порядок выполнения транзакций имеет значение. Стоит отметить, что у очередей FIFO сравнительно небольшое количество вызовов API в секунду.

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

В Yandex Cloud сервисом очередей является Yandex Message Queue. Он относится к PaaS-слою и объединён с другими serverless-сервисами в группу Бессерверные вычисления. В следующем уроке мы поговорим о его специфике.
- Категория: Cloud Services Engineer
- Просмотров: 597