Примечание.
Для доступа к этой странице требуется авторизация. Вы можете попробовать войти или изменить каталоги.
Для доступа к этой странице требуется авторизация. Вы можете попробовать изменить каталоги.
В этой статье показано, как использовать повторное секционирование для масштабирования запроса Azure Stream Analytics в сценариях, в которых невозможна полная параллелизация.
Использование параллелизации может оказаться невозможным, если:
- вы не контролируете ключ раздела для входного потока.
- ваш источник "распыляет" входные данные по нескольким секциям, которые позднее необходимо объединить.
Повторное секционирование или перегруппировка требуются при обработке данных в потоке, не сегментированном в соответствии с естественной схемой ввода, например PartitionId для Центров событий. При повторном секционировании каждый сегмент можно обрабатывать независимо, что позволяет линейно масштабировать конвейер потоковой передачи.
Как выполнить повторное разделение
Можно перераспределить входные данные двумя способами:
- Используйте отдельное задание Stream Analytics, которое выполняет перераспределение.
- использовать одно задание, но сначала выполнить повторное секционирование перед пользовательской логикой аналитики.
Создание отдельного задания Stream Analytics для повторного разбиения входных данных на разделы
Можно создать задание, которое считывает входные данные и записывает в концентратор событий с использованием ключа партиции. Затем этот концентратор событий может служить входными данными для другого задания Stream Analytics, в котором реализована логика аналитики. При настройке выходных данных концентратора событий в задании необходимо указать ключ раздела, с помощью которого Stream Analytics будет повторно разделять данные.
-- For compat level 1.2 or higher
SELECT *
INTO output
FROM input
--For compat level 1.1 or lower
SELECT *
INTO output
FROM input PARTITION BY PartitionId
Переразбиение входных данных в одном задании Stream Analytics
Вы также можете включить шаг в ваш запрос, который сначала перераспределяет входные данные, что затем могут использовать другие шаги в запросе. Например, если вы хотите повторно секционировать входные данные на основе DeviceId, запрос будет выглядеть следующим образом:
WITH RepartitionedInput AS
(
SELECT *
FROM input PARTITION BY DeviceID
)
SELECT DeviceID, AVG(Reading) as AvgNormalReading
INTO output
FROM RepartitionedInput
GROUP BY DeviceId, TumblingWindow(minute, 1)
Следующий пример запроса объединяет два потока перераспределенных данных. При присоединении двух потоков повторно разделенных данных потоки должны иметь тот же ключ раздела и одинаковое количество разделов. Результатом будет поток, имеющий единую схему секционирования.
WITH step1 AS
(
SELECT * FROM input1
PARTITION BY DeviceID
),
step2 AS
(
SELECT * FROM input2
PARTITION BY DeviceID
)
SELECT * INTO output
FROM step1 PARTITION BY DeviceID
UNION step2 PARTITION BY DeviceID
Выходная схема должна соответствовать ключу секции потока и счетчику секций, чтобы каждый подпоток можно было очистить независимо. Поток можно также объединить и повторно секционировать по другой схеме перед очисткой, но следует избегать этого метода, так как он увеличивает общую задержку обработки и нагрузку на ресурсы.
Единицы потоковой передачи для перераспределения
Экспериментируйте и следите за использованием ресурсов задания, чтобы определить точное необходимое количество секций. Количество единиц потоковой передачи (SU) должно быть скорректировано в соответствии с физическими ресурсами, необходимыми для каждой секции. Как правило, для каждого раздела требуется шесть SUs. При нехватке ресурсов, выделенных для задания, система применит повторное секционирование только в том случае, если это повысит производительность задания.
Повторное секционирование для выходных данных SQL
Если ваша работа использует базу данных SQL для вывода, используйте явное переподеление, чтобы соответствовать оптимальному количеству разделов и увеличить пропускную способность. Поскольку SQL лучше всего работает с восемью процессами записи, изменение числа потоков до восьми перед очисткой или где-то выше по потоку может повысить производительность заданий.
Если существует более восьми входных секций, наследование схемы секционирования входных данных может не быть подходящим вариантом. Рекомендуется использовать в запросе INTO, чтобы указать точное количество записывающих модулей.
Следующий пример считывает данные из входного потока, независимо от того, выполняется ли его естественное секционирование, и выполняет повторное секционирование потока на десять частей в соответствии с измерением DeviceID, а затем сбрасывает данные на вывод.
SELECT * INTO [output]
FROM [input]
PARTITION BY DeviceID INTO 10
Дополнительные сведения см. в статье Вывод данных Azure Stream Analytics в базу данных SQL Azure.