Выходные данные Azure Stream Analytics в Azure Cosmos DB

Вывод Azure Cosmos DB в Azure Stream Analytics записывает результаты обработки потоков в виде JSON-документации в контейнер Azure Cosmos DB. Он поддерживает архивирование данных и запросы с низкой задержкой к неструктурированным JSON-данным. Понимание того, как ведёт себя этот выход, помогает настроить его под пропускную способность, согласованность и разбиение, необходимые для вашего сценария.

Основы использования Azure Cosmos DB как целевого выходного пункта

Выход Azure Cosmos DB в Stream Analytics записывает результаты потоковой обработки в формате JSON в контейнеры Azure Cosmos DB. Если вы не знакомы с Azure Cosmos DB, для начала работы обратитесь к документации Azure Cosmos DB.

Stream Analytics подключается к Azure Cosmos DB только через SQL API. Другие API Azure Cosmos DB пока не поддерживаются. Если указать модулю Stream Analytics учетные записи Azure Cosmos DB, созданные при помощи других API, это может привести к неправильному сохранению данных. Когда вы используете Azure Cosmos DB в качестве выхода, установите задачу на уровень совместимости 1.2.

Stream Analytics не создает контейнеры в базе данных. Их необходимо создать заранее. После этого вы можете контролировать расходы на выставление счетов за контейнеры Azure Cosmos DB. Кроме того, можно настроить производительность, согласованность и емкость контейнеров напрямую через API Azure Cosmos DB. В следующих разделах подробно описаны некоторые параметры контейнера Azure Cosmos DB.

Настройка согласованности, доступности и задержки

Чтобы соответствовать требованиям вашего приложения, тонко настройте базу данных и контейнеры в Azure Cosmos DB и делайте компромиссы между согласованностью, доступностью, задержкой и пропускной способностью.

В зависимости от того, какой уровень согласованности чтения нужен вашему сценарию по сравнению с задержкой чтения и записи, выберите уровень согласованности в вашей базе данных. Чтобы повысить пропускную способность, увеличьте число единиц запросов (RU) для контейнера. Кроме того, Azure Cosmos DB по умолчанию активирует синхронное индексирование для каждой операции CRUD в вашем контейнере. Эта опция — ещё один полезный способ контролировать производительность чтения и записи в Azure Cosmos DB. Дополнительные сведения можно найти в статье о том, как изменить уровни согласованности базы данных и запросов.

Вставка и обновление Upsert в Stream Analytics

Используя интеграцию Stream Analytics с Azure Cosmos DB, вы можете вставлять или обновлять записи в контейнере на основе заданного столбца Document ID. Эта операция также называется upsert. Stream Analytics использует оптимистичный подход upsert. Обновления происходят только в случае ошибки при вставке с конфликтом идентификатора документа.

Используя уровень совместимости 1.0, Stream Analytics выполняет это обновление в виде PATCH-операции, поэтому поддерживает частичные обновления документа. Stream Analytics добавляет новые свойства или постепенно заменяет существующее свойство. Тем не менее изменение значений свойств массива в документе JSON приводит к перезаписи всего массива. То есть массивы не объединяются.

При использовании уровня совместимости 1.2 поведение операции upsert изменяется на вставку или замену документа. Подробно это поведение описано ниже, в разделе об уровне совместимости 1.2.

Если входящий JSON-документ имеет уже существующее поле ID, Azure Cosmos DB автоматически использует это поле как столбец Document ID. Stream Analytics обрабатывает любые последующие записи таким образом, что приводит к одной из следующих ситуаций:

  • Уникальные идентификаторы приводят к вставке.
  • При совпадении идентификаторов и значении Идентификатор документа, выставленном по ID, выполняется upsert.
  • В случае дублирования идентификаторов и незаданного Идентификатора документа после первого документа возникает ошибка.

Если вы хотите сохранить все документы, включая те, которые имеют дублирующийся идентификатор, переименуйте поле идентификатора в запросе (с помощью ключевого слова AS). Позвольте Azure Cosmos DB создать поле идентификатора или замените идентификатор другим значением столбца (с помощью ключевого слова AS или параметра Идентификатор документа).

Разделение данных в Azure Cosmos DB

Azure Cosmos DB автоматически масштабирует секции в зависимости от рабочей нагрузки. Используйте неограниченное количество контейнеров для разделения данных. При записи в неограниченные контейнеры Stream Analytics использует параллельные записывающие процессы в таком же количестве, как на предыдущем шаге запроса или в схеме секционирования входных данных.

Примечание.

Azure Stream Analytics поддерживает только контейнеры без ограничений с ключами раздела на верхнем уровне. Например, /region поддерживается. Вложенные ключи разделов (например, /region/name) не поддерживаются.

В зависимости от выбранного ключа секции может появиться следующее предупреждение:

CosmosDB Output contains multiple rows and just one row per partition key. If the output latency is higher than expected, consider choosing a partition key that contains at least several hundred records per partition key.

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

Размер хранилища для документов, относящихся к одному значению ключа секции, ограничен 20 ГБ (а ограничение размера физической секции составляет 50 ГБ). Идеальный ключ раздела — это тот, который часто появляется как фильтр в ваших запросах и обладает достаточной мощностью для масштабируемости решения.

Ключи секций, используемые для запросов Stream Analytics и Azure Cosmos DB, не должны быть идентичными. Для полностью параллельных топологий используйте PartitionId, , в качестве ключа раздела в запросе Stream Analytics, но этот вариант может быть не лучшим выбором в качестве ключа раздела для контейнера Azure Cosmos DB.

Ключ секции также является границей для транзакций в хранимых процедурах и триггерах для Azure Cosmos DB. Выберите ключ раздела так, чтобы документы, встречающиеся вместе в транзакциях, имели одно и то же значение ключа раздела. Статья Секционирование в Azure Cosmos DB содержит дополнительные сведения о выборе ключа секции.

Для контейнеров Azure Cosmos DB с фиксированной пропускной способностью Stream Analytics не предоставляет возможности ни вертикального, ни горизонтального масштабирования после исчерпания их емкости. Они имеют верхний предел в 10 ГБ и 10 000 ЕЗ/с для пропускной способности. Чтобы перенести данные из фиксированного контейнера в контейнер с неограниченными возможностями (например, с ключом раздела и пропускной способностью не менее 1000 RU/с), используйте средство миграции данных или библиотеку канала изменений.

Возможность записи в несколько фиксированных контейнеров прекращается. Не используйте его для масштабирования своей работы в Stream Analytics.

Улучшенная пропускная способность с уровнем совместимости 1.2

При использовании уровня совместимости 1.2 Stream Analytics поддерживает встроенную интеграцию для массовой записи в Azure Cosmos DB. При использовании этой интеграции Stream Analytics эффективно записывает данные в Azure Cosmos DB, обеспечивая максимальную пропускную способность и эффективно обрабатывая запросы на регулирование скорости.

Улучшенный механизм записи доступен на новом уровне совместимости из-за разницы в поведении upsert. Используя уровни до версии 1.2, upsert вставляет или объединяет документ. При использовании 1.2 поведение операции upsert изменяется: документ вставляется или заменяется.

Используя уровни до версии 1.2, Stream Analytics использует кастомную хранящую процедуру для массового размещения документов по ключу раздела в Azure Cosmos DB. Там Stream Analytics записывает пакет в виде транзакции. Даже если одна запись сталкивается с временной ошибкой (throttling), сервису Stream Analytics приходится повторно обрабатывать весь пакет. Такое поведение делает сценарии даже с разумным троттлингом медленными.

В следующем примере показаны два идентичных задания Stream Analytics, считывающие одни и те же входные данные с концентраторов событий Azure. Оба задания Stream Analytics полностью секционированы с помощью простого запроса и записывают данные в идентичные контейнеры Azure Cosmos DB. Метрики слева относятся к заданию с уровнем совместимости 1.0. Метрики справа взяты из задания, настроенного с параметром 1.2. Ключ партиционирования контейнера Azure Cosmos DB — это уникальный идентификатор GUID, который поступает из входящего события.

Снимок экрана: сравнение метрик Stream Analytics.

Скорость поступления событий в Event Hubs в два раза превышает объём, на приём которого настроены контейнеры Azure Cosmos DB (20 000 RU), поэтому в Azure Cosmos DB можно ожидать троттлинг. Однако задание с уровнем 1.2 постоянно выполняет запись с более высокой пропускной способностью (выходные события в минуту), а также с меньшим средним использованием единиц потоковой передачи SU%. В вашей среде эта разница зависит от нескольких факторов. Эти факторы включают в себя формат события, размер входного события или сообщения, ключи секций и запросы.

Снимок экрана: сравнение метрик Azure Cosmos DB.

При использовании версии 1.2 Stream Analytics эффективнее задействует 100 % доступной пропускной способности Azure Cosmos DB, практически без повторных отправок из-за троттлинга или ограничения скорости. Это поведение обеспечивает лучший интерфейс для других рабочих нагрузок, таких как запросы, выполняемые в контейнере одновременно. Если вы хотите увидеть, как Stream Analytics масштабируется с Azure Cosmos DB в качестве приемника для 1000–10 000 сообщений в секунду, попробуйте этот пример проекта Azure.

Пропускная способность вывода Azure Cosmos DB одинакова при использовании версий 1.0 и 1.1. Мы настоятельно рекомендуем использовать уровень совместимости 1.2 в Stream Analytics при работе с Azure Cosmos DB.

Параметры Azure Cosmos DB для выходных данных JSON

Когда вы настраиваете Azure Cosmos DB как выход в Stream Analytics, следующие свойства определяют результат.

Снимок экрана: поля сведений для выходного потока Azure Cosmos DB.

Поле Описание
Алиас результата Псевдоним для ссылки на эти выходные данные в запросе Stream Analytics.
Подписка Подписка Azure.
Код счета Имя или URI конечной точки учетной записи Azure Cosmos DB.
Ключ учетной записи Общедоступный ключ доступа к учетной записи Azure Cosmos DB.
База данных Имя базы данных Azure Cosmos DB.
Имя контейнера Имя контейнера, например MyContainer. Должен существовать один контейнер с именем MyContainer.
Код документа Необязательно. Название столбца в выходных событиях служит уникальным ключом для операций вставки или обновления. Если оставить его пустым, Stream Analytics вставляет все события без опции обновления.

После настройки вывода данных в Azure Cosmos DB можно использовать его в запросе в качестве целевого объекта предложения INTO. Когда вы используете вывод Azure Cosmos DB таким образом, нужно специально установить ключ раздела.

Выходная запись должна содержать столбец с учетом регистра, названный так же, как ключ раздела в Azure Cosmos DB. Для достижения большей параллелизации в выражении может потребоваться предложение PARTITION BY, использующее тот же столбец.

Вот пример запроса:

    SELECT TollBoothId, PartitionId
    INTO CosmosDBOutput
    FROM Input1 PARTITION BY PartitionId

Обработка ошибок и повторные попытки

Если происходит временная ошибка, недоступность службы или ограничение скорости, когда Stream Analytics отправляет события в Azure Cosmos DB, Stream Analytics повторяет попытки до успешного завершения операции. Но он не выполняет повторные попытки при ошибках «Неавторизовано» (код ошибки HTTP 401), «Не найдено» (код ошибки HTTP 404), «Доступ запрещён» (код ошибки HTTP 403) или «Неверный запрос» (код ошибки HTTP 400).

Распространённые проблемы, приводящие к сбоям вывода Azure Cosmos DB

Ряд условий может привести к сбою вывода Azure Cosmos DB. Выходные данные из Stream Analytics могут нарушать уникальное ограничение индекса контейнера, столбец PartitionKey может отсутствовать, или столбец Id может отсутствовать. Для получения дополнительной информации об уникальных индексных ограничениях см. раздел «Уникальные ключи в Azure Cosmos DB».