Simulação de negociação de alta frequência com o Stream Analytics

Azure Stream Analytics dá suporte à análise avançada por meio da combinação de linguagem SQL, UDFs (funções definidas pelo usuário) e UDAs (agregações definidas pelo usuário). As análises avançadas podem incluir treinamento de aprendizado de máquina online e pontuação, além da simulação de processo com estado. Este artigo descreve como executar a regressão linear em um trabalho Azure Stream Analytics que faz treinamento contínuo e pontuação em um cenário de negociação de alta frequência.

Pré-requisitos

Fluxo de trabalho de negociação de alta frequência

O fluxo lógico da negociação de alta frequência é:

  1. Obtenção de cotações em tempo real de uma bolsa de valores.
  2. Construindo um modelo preditivo com base nas cotações para antecipar a movimentação dos preços.
  3. Fazendo pedidos de compra ou venda para ganhar dinheiro com a previsão bem-sucedida dos movimentos de preços.

Este cenário requer:

  • De um feed de cotação em tempo real.
  • Um modelo preditivo que pode operar com cotações em tempo real.
  • Uma simulação de negociação que demonstra o lucro ou a perda do algoritmo de negociação.

Feed de cotações em tempo real

Importante

A API websocket de negociação IEX (iextrading.com) referenciada nesta seção foi desativada. A IEX Cloud agora disponibiliza dados de mercado por meio da IEX Cloud, com autenticação e endpoints diferentes. Atualize a URL e a autenticação em sua implementação adequadamente.

Importante

Os pacotes NuGet SocketIoClientDotNet e WindowsAzure.ServiceBus usados neste exemplo estão obsoletos. Para novos projetos, use uma biblioteca de clientes Socket.IO atual e o pacote Azure.Messaging.EventHubs com EventHubProducerClient em vez do EventHubClient herdado.

A Investors Exchange (IEX) oferecia anteriormente cotações de compra e venda em tempo real gratuitamente usando o socket.io. Você pode escrever um programa simples de console para receber cotações em tempo real e enviá-las para o Hubs de Eventos do Azure como fonte de dados. O código a seguir é um esqueleto do programa. O código omite o tratamento de erros para fins de brevidade. Você também precisa incluir os pacotes NuGet SocketIoClientDotNeteWindowsAzure.ServiceBus em seu projeto.

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);
});

Caution

Este exemplo de código destina-se apenas à ilustração. O ponto de extremidade da API WebSocket IEX e os pacotes NuGet usados aqui não estão mais disponíveis. Não use esse código em produção. Consulte as notas IMPORTANTEs anteriormente nesta seção para obter alternativas atuais.

Aqui estão alguns eventos de exemplo gerados:

{"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

O carimbo de data/hora do evento é lastUpdated, em época.

Modelo preditivo para negociação de alta frequência

Para essa demonstração, o exemplo usa um modelo linear descrito na estratégia baseada em desequilíbrio de ordem na negociação algorítmica de alta frequência.

Desequilíbrio de ordem de volume (VOI) é uma função do preço de compra e venda atual e do volume e preço de compra e venda do último tique. O artigo identifica a correlação entre VOI e movimentação futura de preços. Ele cria um modelo linear entre os últimos cinco valores de VOI, e a mudança de preço nos próximos 10 tiques. O modelo é treinado com base nos dados do dia anterior por meio de regressão linear.

Em seguida, o modelo treinado faz previsões de alteração de preço nas cotações do dia de negociação atual em tempo real. Quando o modelo prevê uma alteração de preço grande o suficiente, ele executa uma negociação. Dependendo da configuração do limite, uma única ação pode gerar milhares de negociações durante um dia de negociação.

Diagrama que mostra a fórmula de definição do desequilíbrio no volume das ordens usada em negociações de alta frequência.

As seções a seguir mostram como expressar as operações de treinamento e previsão em um trabalho Azure Stream Analytics. A consulta completa é uma única instrução WITH composta por CTEs (expressões de tabela comum) que formam um pipeline:

Estágio CTE Purpose
typeconvertedquotes Converter campos de entrada brutos em tipos SQL adequados
timefilteredquotes Filtrar cotações para o horário de negociação e remover dados inválidos
shiftedquotes Usar LAG para recuperar os valores de oferta/lance da marca anterior
currentPriceAndVOI Calcular o VOI (desequilíbrio de volume das ordens) a partir da marca atual e da anterior
shiftedPriceAndShiftedVOI Crie sequências de 10 preços médios consecutivos e 2 valores consecutivos de VOI
modelInput Remodelar dados em vetores de recursos (VOI como x, delta de preço como y)
modelagg / modelparambs / model Treinar um modelo de regressão linear de duas variáveis usando agregações SUM e AVG
shiftedVOI / VOIAndModel / VOIANDModelJoined Combinar os valores VOI atuais com o modelo treinado do dia anterior
prediction Calcular a variação futura esperada do preço (efpc) com base no modelo
tradeSignal Gerar sinais de compra/venda quando o efpc exceder o limite de ±0,02

Note

Esta consulta requer o nível de compatibilidade 1.1 ou posterior do Azure Stream Analytics, que preserva a capitalização dos nomes dos campos para um comportamento previsível com UDAs.

Limpar e converter campos de entrada de cota

A primeira CTE na consulta do Azure Stream Analytics converte os dados brutos de cotas do Event Hubs em colunas SQL tipadas da forma correta. DATEADD converte a hora da época (Unix milissegundos) em datetime. TRY_CAST impõe os tipos de dados sem falhar na consulta. Converta os campos de entrada para os tipos de dados esperados para evitar comportamentos inesperados ao manipular ou comparar os campos.

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
),

Recuperar valores de marca anteriores com LAG

A próxima CTE na consulta do Azure Stream Analytics usa a função LAG para obter o preço e o tamanho de compra/venda da marca anterior para cada símbolo da ação. Uma hora de valor LIMIT DURATION é escolhida arbitrariamente. Considerando a frequência de cotas, você pode encontrar a marca anterior retrocedendo uma hora.

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
),

Calcular o desequilíbrio da ordem de volume (VOI)

O próximo CTE calcula o valor de VOI a partir dos dados de compra/venda da marca atual e da anterior. A consulta filtra valores nulos para casos em que nenhuma marca anterior existe.

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
),

Criar sequências de atributos para treinar o modelo

O próximo CTE usa LAG novamente para criar uma sequência com 2 valores consecutivos de VOI, seguidos por 10 valores consecutivos de preço médio. Essas sequências formam os dados de treinamento para o modelo de regressão linear.

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
),

Reestruturar dados em vetores de características

O próximo CTE reorganiza as sequências de preço e de VOI em vetores de características para um modelo linear de duas variáveis, em que os valores de VOI são as variáveis independentes (x1, x2) e a variação média futura do preço é a variável dependente (y). Eventos com dados incompletos são filtrados.

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
),

Treinar o modelo de regressão linear com SUM e AVG

Como Azure Stream Analytics não tem uma função de regressão linear interna, a consulta usa agregações SUM e AVG para calcular os coeficientes (a, b1, b2) para o modelo de regressão linear de duas variáveis. O modelo é treinado diariamente usando uma janela em cascata de 24 horas.

Diagrama que mostra a fórmula matemática de regressão linear para coeficientes de modelo de computação.

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
),

Atribuir pontuação às cotações atuais usando o modelo do dia anterior

Para usar o modelo de regressão linear treinado do dia anterior para pontuar o evento atual, a consulta une as aspas com os coeficientes do modelo. Em vez de usar JOIN, a consulta usa UNION para combinar eventos de modelo e eventos de aspas em um único fluxo. Em seguida, ele usa LAG para emparelhar os eventos com o modelo do dia anterior, para que você obtenha exatamente uma correspondência. Devido ao fim de semana, a consulta retroage três dias (72 horas). Se fosse usado um JOIN simples, você obteria três modelos para cada evento de cotação.

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'
),

Gerar sinais comerciais a partir de previsões

Os CTEs finais calculam a alteração de preço futura esperada (efpc) aplicando a fórmula de regressão linear (a + b1 * x1 + b2 * x2) e, em seguida, geram sinais de compra/venda com base em um limite de ±0,02. Um valor de transação de 10 representa compra. Um valor de negociação de -10 significa venda.

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
),

Testar a estratégia de negociação com uma simulação

Depois de gerar os sinais de negociação, teste o quão eficaz é a estratégia de negociação sem negociação real.

Este teste usa um UDA com uma janela que salta a cada um minuto. O agrupamento na data e a cláusula HAVING garantem que a janela apenas contabilize eventos que pertencem ao mesmo dia. Para uma janela de salto que inclua dois dias, GROUP BY por data, separa o agrupamento em dia anterior e atual. A cláusula HAVING filtra as janelas que terminam no dia atual, mas o agrupamento no dia anterior.

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
)

A UDA do JavaScript inicializa todos os acumuladores na init função, calcula a transição de estado com cada evento adicionado à janela e retorna os resultados da simulação no final da janela. A simulação mantém uma posição comprada ou vendida de 10 ações de um papel em cada operação. O custo da transação é fixo $8. A tabela a seguir mostra as quatro ações de negociação executadas pela UDA:

Condition Sinal Ação Posição após
Nenhuma retenção atual Comprar (10) Comprar para abrir Long
Nenhuma retenção atual Vender (-10) Vender para abrir (curto) Short
Posição longa Vender (-10) Vender para fechar e, em seguida, vender para abrir (curto) Short
Posição curta Comprar (10) Comprar para fechar e, em seguida, comprar para abrir 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

O conector de saída do Power BI para o Azure Stream Analytics está programado para ser descontinuado. Considere o uso de destinos de saída alternativos, como Azure Data Explorer, Azure Synapse Analytics ou um repositório de dados ao qual Power BI pode se conectar por meio do DirectQuery ou importação. Para obter mais informações, consulte saída do Azure Stream Analytics para o Power BI.

Por fim, envie os dados para o painel do Power BI para visualização.

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 que mostra sinais comerciais visualizados em um painel de Power BI para a simulação de negociação.

Chart que mostra os resultados de lucro e perda visualizados em um painel de Power BI para a simulação de negociação.

Resumo

Este artigo mostra como implementar um modelo de negociação realista de alta frequência com uma consulta moderadamente complexa em Azure Stream Analytics. O modelo usa duas variáveis de entrada em vez de cinco porque Azure Stream Analytics não inclui uma função de regressão linear interna. No entanto, você também pode implementar algoritmos mais sofisticados com dimensões mais altas como UDAs JavaScript.

Você pode testar e depurar a maior parte da consulta, exceto o JavaScript UDA, usando as ferramentas do Azure Stream Analytics para Visual Studio Code para desenvolver, testar e depurar consultas.