Моделирование высокочастотной торговли с помощью Stream Analytics

Azure Stream Analytics поддерживает расширенную аналитику за счет сочетания языка SQL, пользовательских функций JavaScript (UDF) и пользовательских агрегатов (UDA). Расширенная аналитика включает онлайн-обучение моделей машинного обучения и вычисление оценок, а также имитационное моделирование процессов с сохранением состояния. В этой статье описывается, как выполнять линейную регрессию в Azure Stream Analytics задании, которое выполняет непрерывную подготовку и оценку в сценарии с высокой частотой торговли.

Необходимые условия

Высокочастотный рабочий процесс торговли

Логический поток высокочастотной торговли:

  1. Получение котировок в режиме реального времени с биржи ценных бумаг.
  2. Построение прогнозной модели на основе котировок для прогнозирования изменения цен.
  3. Размещение заказов на покупку или продажу, чтобы заработать деньги от успешного прогнозирования движения цен.

Для этого сценария требуется:

  • Поток котировок в реальном времени.
  • Предиктивная модель, которая может работать с котировками в реальном времени.
  • Симуляция торговли, демонстрирующая прибыль или потерю алгоритма торговли.

Поток котировок в реальном времени

Important

WebSocket API для торговли IEX (iextrading.com), на который ссылаются в этом разделе, выведен из эксплуатации. IEX Cloud теперь предоставляет рыночные данные через IEX Cloud с различными проверками подлинности и конечными точками. Измените URL-адрес и проверку подлинности в реализации соответствующим образом.

Important

Пакеты NuGet SocketIoClientDotNet и WindowsAzure.ServiceBus, используемые в этом примере, считаются устаревшими. Для новых проектов используйте актуальную клиентскую библиотеку Socket.IO и пакет Azure.Messaging.EventHubs с EventHubProducerClient вместо устаревшей EventHubClient.

Investors Exchange (IEX) ранее предлагала бесплатные котировки спроса и предложения в реальном времени через socket.io. Вы можете написать простую консольную программу для получения котировок в реальном времени и отправки их в Центры событий Azure, используя его в качестве источника данных. Следующий код является скелетом программы. Код исключает обработку ошибок для краткости. Кроме того, необходимо добавить в проект пакеты NuGet SocketIoClientDotNet и WindowsAzure.ServiceBus.

using Quobject.SocketIoClientDotNet.Client;
using Microsoft.ServiceBus.Messaging;
var symbols = "msft,fb,amzn,goog";
var eventHubClient = EventHubClient.CreateFromConnectionString(connectionString, eventHubName);
var socket = IO.Socket("https://ws-api.iextrading.com/1.0/tops");
socket.On(Socket.EVENT_MESSAGE, (message) =>
{
    eventHubClient.Send(new EventData(Encoding.UTF8.GetBytes((string)message)));
});
socket.On(Socket.EVENT_CONNECT, () =>
{
    socket.Emit("subscribe", symbols);
});

Предостережение

Этот пример кода предназначен только для иллюстрации. Конечная точка API IEX WebSocket и используемые здесь пакеты NuGet больше недоступны. Не используйте этот код в рабочей среде. Сведения о текущих альтернативах см. в ВАЖНЫХ примечаниях выше в этом разделе.

Ниже приведены некоторые созданные примеры событий:

{"symbol":"MSFT","marketPercent":0.03246,"bidSize":100,"bidPrice":74.8,"askSize":300,"askPrice":74.83,"volume":70572,"lastSalePrice":74.825,"lastSaleSize":100,"lastSaleTime":1506953355123,"lastUpdated":1506953357170,"sector":"softwareservices","securityType":"commonstock"}
{"symbol":"GOOG","marketPercent":0.04825,"bidSize":114,"bidPrice":870,"askSize":0,"askPrice":0,"volume":11240,"lastSalePrice":959.47,"lastSaleSize":60,"lastSaleTime":1506953317571,"lastUpdated":1506953357633,"sector":"softwareservices","securityType":"commonstock"}
{"symbol":"MSFT","marketPercent":0.03244,"bidSize":100,"bidPrice":74.8,"askSize":100,"askPrice":74.83,"volume":70572,"lastSalePrice":74.825,"lastSaleSize":100,"lastSaleTime":1506953355123,"lastUpdated":1506953359118,"sector":"softwareservices","securityType":"commonstock"}
{"symbol":"FB","marketPercent":0.01211,"bidSize":100,"bidPrice":169.9,"askSize":100,"askPrice":170.67,"volume":39042,"lastSalePrice":170.67,"lastSaleSize":100,"lastSaleTime":1506953351912,"lastUpdated":1506953359641,"sector":"softwareservices","securityType":"commonstock"}
{"symbol":"GOOG","marketPercent":0.04795,"bidSize":100,"bidPrice":959.19,"askSize":0,"askPrice":0,"volume":11240,"lastSalePrice":959.47,"lastSaleSize":60,"lastSaleTime":1506953317571,"lastUpdated":1506953360949,"sector":"softwareservices","securityType":"commonstock"}
{"symbol":"FB","marketPercent":0.0121,"bidSize":100,"bidPrice":169.9,"askSize":100,"askPrice":170.7,"volume":39042,"lastSalePrice":170.67,"lastSaleSize":100,"lastSaleTime":1506953351912,"lastUpdated":1506953362205,"sector":"softwareservices","securityType":"commonstock"}
{"symbol":"GOOG","marketPercent":0.04795,"bidSize":114,"bidPrice":870,"askSize":0,"askPrice":0,"volume":11240,"lastSalePrice":959.47,"lastSaleSize":60,"lastSaleTime":1506953317571,"lastUpdated":1506953362629,"sector":"softwareservices","securityType":"commonstock"}

Note

Метка времени события — lastUpdated, в формате времени эпохи Unix.

Прогнозная модель для высокочастотной торговли

В этой демонстрации в примере используется линейная модель, описанная в работе Стратегия на основе дисбаланса ордеров в высокочастотной алгоритмической торговле.

Дисбаланс объёма заявок (VOI) — это функция текущих цены покупки/продажи и объёма, а также цены покупки/продажи и объёма с последнего тика. В документе определяется корреляция между VOI и будущим движением цен. Она строит линейную модель на основе предыдущих пяти значений VOI и изменения цены в течение следующих 10 тиков. Модель обучается на данных предыдущего дня с помощью линейной регрессии.

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

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

В следующих разделах показано, как выразить операции обучения и прогнозирования в задании Azure Stream Analytics. Полный запрос — это одна WITH инструкция, состоящая из общих табличных выражений (CTEs), которые образуют конвейер:

Этап CTE Purpose
typeconvertedquotes Преобразование необработанных полей ввода в правильные типы SQL
timefilteredquotes Отфильтровать котировки по времени торгов и удалить некорректные данные
shiftedquotes Используйте LAG, чтобы получить значения Bid/Ask предыдущего тика
currentPriceAndVOI Рассчитать дисбаланс объёма заявок (VOI) по текущему и предыдущему тику
shiftedPriceAndShiftedVOI Сформируйте последовательности из 10 последовательных средних цен и 2 последовательных значений VOI
modelInput Перепечатка данных в векторы функций (VOI как x, ценовая дельта как y)
modelagg / modelparambs / model Обучение модели линейной регрессии двух переменных с помощью агрегатов SUM и AVG
shiftedVOI / VOIAndModel / VOIANDModelJoined Объедините текущие значения VOI с обученной моделью, полученной на данных предыдущего дня
prediction Вычисление ожидаемого будущего изменения цен (efpc) из модели
tradeSignal Создание сигналов о покупке и продаже, когда efpc превышает пороговое значение ±0,02

Note

Для этого запроса требуется уровень совместимости Azure Stream Analytics 1.1 или выше, который сохраняет регистр имен полей для предсказуемого поведения с пользовательскими агрегатными функциями (UDA).

Очистить и преобразовать поля ввода кавычек

Первый CTE в запросе Azure Stream Analytics преобразует необработанные данные котировок из Event Hubs в столбцы SQL с корректными типами данных. DATEADD преобразует время эпохи (в миллисекундах Unix) в datetime. TRY_CAST принудительно вводит типы данных без сбоя запроса. Приведение полей входных данных к ожидаемым типам данных, чтобы избежать непредвиденного поведения в манипуляции или сравнении полей.

WITH
typeconvertedquotes AS (
    /* convert all input fields to proper types */
    SELECT
        System.Timestamp AS lastUpdated,
        symbol,
        DATEADD(millisecond, CAST(lastSaleTime as bigint), '1970-01-01T00:00:00Z') AS lastSaleTime,
        TRY_CAST(bidSize as bigint) AS bidSize,
        TRY_CAST(bidPrice as float) AS bidPrice,
        TRY_CAST(askSize as bigint) AS askSize,
        TRY_CAST(askPrice as float) AS askPrice,
        TRY_CAST(volume as bigint) AS volume,
        TRY_CAST(lastSaleSize as bigint) AS lastSaleSize,
        TRY_CAST(lastSalePrice as float) AS lastSalePrice
    FROM quotes TIMESTAMP BY DATEADD(millisecond, CAST(lastUpdated as bigint), '1970-01-01T00:00:00Z')
),
timefilteredquotes AS (
    /* filter between 7am and 1pm PST, 14:00 to 20:00 UTC */
    /* clean up invalid data points */
	SELECT * FROM typeconvertedquotes
	WHERE DATEPART(hour, lastUpdated) >= 14 AND DATEPART(hour, lastUpdated) < 20 AND bidSize > 0 AND askSize > 0 AND bidPrice > 0 AND askPrice > 0
),

Получение предыдущих значений галока с помощью LAG

Следующий CTE в запросе Azure Stream Analytics использует функцию LAG для получения цен покупки и продажи, а также объёмов из предыдущего тика для каждого биржевого символа. Одночасовое значение LIMIT DURATION произвольно выбирается. При такой частоте котировок предыдущий тик можно найти, если вернуться на час назад.

shiftedquotes AS (
    /* get previous bid/ask price and size in order to calculate VOI */
	SELECT
		symbol,
		(bidPrice + askPrice)/2 AS midPrice,
		bidPrice,
		bidSize,
		askPrice,
		askSize,
		LAG(bidPrice) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS bidPricePrev,
		LAG(bidSize) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS bidSizePrev,
		LAG(askPrice) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS askPricePrev,
		LAG(askSize) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS askSizePrev
	FROM timefilteredquotes
),

Рассчитать дисбаланс объёма ордеров (VOI)

Следующий CTE вычисляет значение VOI на основе данных bid/ask текущего и предыдущего тика. Запрос отфильтровывает значения NULL в случаях, когда предыдущий тик отсутствует.

currentPriceAndVOI AS (
    /* calculate VOI */
	SELECT
		symbol,
		midPrice,
		(CASE WHEN (bidPrice < bidPricePrev) THEN 0
            ELSE (CASE WHEN (bidPrice = bidPricePrev) THEN (bidSize - bidSizePrev) ELSE bidSize END)
         END) -
        (CASE WHEN (askPrice < askPricePrev) THEN askSize
            ELSE (CASE WHEN (askPrice = askPricePrev) THEN (askSize - askSizePrev) ELSE 0 END)
         END) AS VOI
	FROM shiftedquotes
	WHERE
		bidPrice IS NOT NULL AND
		bidSize IS NOT NULL AND
		askPrice IS NOT NULL AND
		askSize IS NOT NULL AND
		bidPricePrev IS NOT NULL AND
		bidSizePrev IS NOT NULL AND
		askPricePrev IS NOT NULL AND
		askSizePrev IS NOT NULL
),

Создание последовательностей признаков для обучения модели

Следующий CTE снова использует LAG для создания последовательности с 2 последовательными значениями VOI, а затем 10 последовательных средних ценовых значений. Эти последовательности формируют обучающие данные для модели линейной регрессии.

shiftedPriceAndShiftedVOI AS (
    /* get 10 future prices and 2 previous VOIs */
    SELECT
		symbol,
		midPrice AS midPrice10,
		LAG(midPrice, 1) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS midPrice9,
		LAG(midPrice, 2) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS midPrice8,
		LAG(midPrice, 3) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS midPrice7,
		LAG(midPrice, 4) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS midPrice6,
		LAG(midPrice, 5) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS midPrice5,
		LAG(midPrice, 6) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS midPrice4,
		LAG(midPrice, 7) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS midPrice3,
		LAG(midPrice, 8) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS midPrice2,
		LAG(midPrice, 9) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS midPrice1,
		LAG(midPrice, 10) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS midPrice,
		LAG(VOI, 10) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS VOI1,
		LAG(VOI, 11) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS VOI2
	FROM currentPriceAndVOI
),

Преобразовать данные в векторы признаков

Следующий CTE преобразует последовательности цен и значений VOI в векторы признаков для линейной модели с двумя переменными, где значения VOI являются независимыми переменными (x1, x2), а среднее будущее изменение цены — зависимой переменной (y). События с неполными данными отфильтровываются.

modelInput AS (
    /* create feature vector, x being VOI, y being delta price */
	SELECT
		symbol,
		(midPrice1 + midPrice2 + midPrice3 + midPrice4 + midPrice5 + midPrice6 + midPrice7 + midPrice8 + midPrice9 + midPrice10)/10.0 - midPrice AS y,
		VOI1 AS x1,
		VOI2 AS x2
	FROM shiftedPriceAndShiftedVOI
	WHERE
		midPrice1 IS NOT NULL AND
		midPrice2 IS NOT NULL AND
		midPrice3 IS NOT NULL AND
		midPrice4 IS NOT NULL AND
		midPrice5 IS NOT NULL AND
		midPrice6 IS NOT NULL AND
		midPrice7 IS NOT NULL AND
		midPrice8 IS NOT NULL AND
		midPrice9 IS NOT NULL AND
		midPrice10 IS NOT NULL AND
		midPrice IS NOT NULL AND
		VOI1 IS NOT NULL AND
		VOI2 IS NOT NULL
),

Обучите модель линейной регрессии с использованием SUM и AVG

Так как Azure Stream Analytics не имеет встроенной функции линейной регрессии, запрос использует SUM и AVG для вычисления коэффициентов (a, b1, b2) для модели двух переменной линейной регрессии. Модель ежедневно переобучается с использованием 24-часового кувыркающегося окна.

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

modelagg AS (
    /* get aggregates for linear regression calculation,
     http://faculty.cas.usf.edu/mbrannick/regression/Reg2IV.html */
	SELECT
		symbol,
		SUM(x1 * x1) AS x1x1,
		SUM(x2 * x2) AS x2x2,
		SUM(x1 * y) AS x1y,
		SUM(x2 * y) AS x2y,
		SUM(x1 * x2) AS x1x2,
		AVG(y) AS avgy,
		AVG(x1) AS avgx1,
		AVG(x2) AS avgx2
	FROM modelInput
	GROUP BY symbol, TumblingWindow(hour, 24, -4)
),
modelparambs AS (
    /* calculate b1 and b2 for the linear model */
	SELECT
		symbol,
		(x2x2 * x1y - x1x2 * x2y)/(x1x1 * x2x2 - x1x2 * x1x2) AS b1,
		(x1x1 * x2y - x1x2 * x1y)/(x1x1 * x2x2 - x1x2 * x1x2) AS b2,
		avgy,
		avgx1,
		avgx2
	FROM modelagg
),
model AS (
    /* calculate a for the linear model */
	SELECT
		symbol,
		avgy - b1 * avgx1 - b2 * avgx2 AS a,
		b1,
		b2
	FROM modelparambs
),

Оценка текущих цитат с помощью модели предыдущего дня

Чтобы использовать обученную модель линейной регрессии предыдущего дня для оценки текущего события, запрос присоединяет кавычки к коэффициентам модели. Вместо JOIN запрос использует UNION, чтобы объединить события модели и события котировок в единый поток. Затем она использует LAG для связывания событий с моделью предыдущего дня, поэтому вы получите ровно одно совпадение. Из-за выходных запрос охватывает последние три дня (72 часа). Если бы использовался простой JOIN, вы получили бы три модели для каждого события котировки.

shiftedVOI AS (
    /* get two consecutive VOIs */
	SELECT
		symbol,
		midPrice,
		VOI AS VOI1,		
		LAG(VOI, 1) OVER (PARTITION BY symbol LIMIT DURATION(hour, 1)) AS VOI2
	FROM currentPriceAndVOI
),
VOIAndModel AS (
    /* combine VOIs and models */
	SELECT
		'voi' AS type,
		symbol,
		midPrice,
		VOI1,
		VOI2,
        0.0 AS a,
        0.0 AS b1,
        0.0 AS b2
	FROM shiftedVOI
	UNION
	SELECT
		'model' AS type,
		symbol,
        0.0 AS midPrice,
        0 AS VOI1,
        0 AS VOI2,
		a,
		b1,
		b2
	FROM model
),
VOIANDModelJoined AS (
    /* match VOIs with the latest model within 3 days (72 hours, to take the weekend into account) */
	SELECT
		symbol,
		midPrice,
		VOI1 as x1,
		VOI2 as x2,
		LAG(a, 1) OVER (PARTITION BY symbol LIMIT DURATION(hour, 72) WHEN type = 'model') AS a,
		LAG(b1, 1) OVER (PARTITION BY symbol LIMIT DURATION(hour, 72) WHEN type = 'model') AS b1,
		LAG(b2, 1) OVER (PARTITION BY symbol LIMIT DURATION(hour, 72) WHEN type = 'model') AS b2
	FROM VOIAndModel
	WHERE type = 'voi'
),

Создание торговых сигналов из прогнозов

Окончательные CTE вычисляют ожидаемое будущее изменение цены (efpc), применяя формулу линейной регрессии (a + b1 * x1 + b2 * x2), а затем генерируют сигналы на покупку и продажу на основе порогового значения ±0,02. Торговое значение 10 означает покупку. Торговое значение -10 означает продажу.

prediction AS (
    /* make prediction if there is a model */
	SELECT
		symbol,
		midPrice,
		a + b1 * x1 + b2 * x2 AS efpc
	FROM VOIANDModelJoined
	WHERE
		a IS NOT NULL AND
		b1 IS NOT NULL AND
		b2 IS NOT NULL AND
        x1 IS NOT NULL AND
        x2 IS NOT NULL
),
tradeSignal AS (
    /* generate buy/sell signals */
	SELECT
        DateAdd(hour, -7, System.Timestamp) AS time,
		symbol,		
		midPrice,
        efpc,
		CASE WHEN (efpc > 0.02) THEN 10 ELSE (CASE WHEN (efpc < -0.02) THEN -10 ELSE 0 END) END AS trade,
		DATETIMEFROMPARTS(DATEPART(year, System.Timestamp), DATEPART(month, System.Timestamp), DATEPART(day, System.Timestamp), 0, 0, 0, 0) as date
	FROM prediction
),

Тестирование стратегии торговли с помощью имитации

После генерации торговых сигналов проверьте, насколько эффективна торговая стратегия, не торгуя реальными средствами.

В этом тесте используется UDA со скачкообразным окном, которое сдвигается каждую минуту. Группировка по дате и условие HAVING обеспечивают, что окно учитывает только события, относящиеся к одному и тому же дню. Для интервала в два дня дата GROUP BY разделяет группирование на предыдущий день и текущий день. Выражение HAVING отфильтровывает окна, которые заканчиваются текущим днём, но группируются по предыдущему дню.

simulation AS
(
    /* perform trade simulation for the past 7 hours to cover an entire trading day, and generate output every minute */
	SELECT
        DateAdd(hour, -7, System.Timestamp) AS time,
		symbol,
		date,
		uda.TradeSimulation(tradeSignal) AS s
	FROM tradeSignal
	GROUP BY HoppingWindow(minute, 420, 1), symbol, date
	Having DateDiff(day, date, time) < 1 AND DATEPART(hour, time) < 13
)

JavaScript UDA инициализирует все аккумуляторы в функции init, вычисляет изменение состояния при каждом добавлении события в окно и возвращает результаты моделирования в конце окна. Симуляция открывает длинную или короткую позицию по 10 акций в каждой сделке. Стоимость транзакции фиксированная: $8. В следующей таблице показаны четыре торговые действия, выполняемые UDA:

Condition Сигнал Действие Позиция после
Нет текущего удержания Купить (10) Купить, чтобы открыть Long
Нет текущего удержания Продать (-10) Продажа в открытие (короткая позиция) Short
Длинная позиция Продать (-10) Продажа для закрытия, затем продажа для открытия короткой позиции Short
Короткая позиция Купить (10) Покупка для закрытия, затем покупка для открытия Long
function main() {
	var TRADE_COST = 8.0;
	var SHARES = 10;
	this.init = function () {
		this.own = false;
		this.pos = 0;
		this.pnl = 0.0;
		this.tradeCosts = 0.0;
		this.buyPrice = 0.0;
		this.sellPrice = 0.0;
		this.buySize = 0;
		this.sellSize = 0;
		this.buyTotal = 0.0;
		this.sellTotal = 0.0;
	}
	this.accumulate = function (tradeSignal, timestamp) {
		if(!this.own && tradeSignal.trade == 10) {
		  // Buy to open
		  this.own = true;
		  this.pos = 1;
		  this.buyPrice = tradeSignal.midprice;
		  this.tradeCosts += TRADE_COST;
		  this.buySize += SHARES;
		  this.buyTotal += SHARES * tradeSignal.midprice;
		} else if(!this.own && tradeSignal.trade == -10) {
		  // Sell to open
		  this.own = true;
		  this.pos = -1
		  this.sellPrice = tradeSignal.midprice;
		  this.tradeCosts += TRADE_COST;
		  this.sellSize += SHARES;
		  this.sellTotal += SHARES * tradeSignal.midprice;
		} else if(this.own && this.pos == 1 && tradeSignal.trade == -10) {
		  // Sell to close
		  this.own = false;
		  this.pos = 0;
		  this.sellPrice = tradeSignal.midprice;
		  this.tradeCosts += TRADE_COST;
		  this.pnl += (this.sellPrice - this.buyPrice)*SHARES - 2*TRADE_COST;
		  this.sellSize += SHARES;
		  this.sellTotal += SHARES * tradeSignal.midprice;
		  // Sell to open
		  this.own = true;
		  this.pos = -1;
		  this.sellPrice = tradeSignal.midprice;
		  this.tradeCosts += TRADE_COST;
		  this.sellSize += SHARES;		  
		  this.sellTotal += SHARES * tradeSignal.midprice;
		} else if(this.own && this.pos == -1 && tradeSignal.trade == 10) {
		  // Buy to close
		  this.own = false;
		  this.pos = 0;
		  this.buyPrice = tradeSignal.midprice;
		  this.tradeCosts += TRADE_COST;
		  this.pnl += (this.sellPrice - this.buyPrice)*SHARES - 2*TRADE_COST;
		  this.buySize += SHARES;
		  this.buyTotal += SHARES * tradeSignal.midprice;
		  // Buy to open
		  this.own = true;
		  this.pos = 1;
		  this.buyPrice = tradeSignal.midprice;
		  this.tradeCosts += TRADE_COST;
		  this.buySize += SHARES;		  
		  this.buyTotal += SHARES * tradeSignal.midprice;
		}
	}
	this.computeResult = function () {
		var result = {
			"pnl": this.pnl,
			"buySize": this.buySize,
			"sellSize": this.sellSize,
			"buyTotal": this.buyTotal,
			"sellTotal": this.sellTotal,
			"tradeCost": this.tradeCost
			};
		return result;
	}
}

Note

Планируется вывести из эксплуатации соединитель вывода Power BI для Azure Stream Analytics. Рассмотрите возможность использования альтернативных назначений выходных данных, таких как Azure Data Explorer, Azure Synapse Analytics или хранилище данных, к которому Power BI может подключаться через DirectQuery или импорт. Дополнительные сведения см. в разделе Вывод данных из Azure Stream Analytics в Power BI.

Наконец, выведите данные на панель мониторинга Power BI для визуализации.

SELECT * INTO tradeSignalDashboard FROM tradeSignal /* output tradeSignal to PBI */
SELECT
    symbol,
    time,
    date,
    TRY_CAST(s.pnl as float) AS pnl,
    TRY_CAST(s.buySize as bigint) AS buySize,
    TRY_CAST(s.sellSize as bigint) AS sellSize,
    TRY_CAST(s.buyTotal as float) AS buyTotal,
    TRY_CAST(s.sellTotal as float) AS sellTotal
    INTO pnlDashboard
FROM simulation /* output trade simulation to PBI */

Chart, отображающий торговые сигналы, визуализированные на панели мониторинга Power BI для моделирования торговли.

Chart, отображающий результаты прибыли и потери, визуализированные на панели мониторинга Power BI для моделирования торговли.

Сводка

В этой статье показано, как реализовать реалистичную высокочастотную модель торговли с умеренно сложным запросом в Azure Stream Analytics. Модель использует две входные переменные вместо пяти, так как Azure Stream Analytics не включает встроенную функцию линейной регрессии. Однако вы также можете реализовать более сложные алгоритмы для пространств большей размерности в виде агрегатных функций, определяемых пользователем, на JavaScript.

Вы можете протестировать и отладить большую часть запроса, кроме UDA JavaScript, с помощью средств Azure Stream Analytics для Visual Studio Code для разработки запросов, тестирования и отладки.