Simulace vysokofrekvenčního obchodování se Stream Analytics

Azure Stream Analytics podporuje pokročilou analýzu prostřednictvím kombinace jazyka SQL, uživatelem definovaných funkcí JavaScriptu (UDF) a uživatelem definovaných agregací (UDA). Pokročilá analýza zahrnuje online trénování a vyhodnocování strojového učení a simulaci stavových procesů. Tento článek popisuje, jak provádět lineární regresi v Azure Stream Analytics úloze, která provádí průběžné trénování a vyhodnocování ve scénáři vysokofrekvenčního obchodování.

Předpoklady

Pracovní postup vysokofrekvenčního obchodování

Logický tok vysokofrekvenčního obchodování je:

  1. Získávání kurzů v reálném čase z burzy cenných papírů
  2. Vytvoření prediktivního modelu kolem uvozovek za účelem předvídání pohybu cen
  3. Zadávání nákupních a prodejních pokynů s cílem vydělat na základě správné předpovědi cenových pohybů.

Tento scénář vyžaduje:

  • Informační kanál nabídek v reálném čase.
  • Prediktivní model, který může pracovat s uvozovkami v reálném čase.
  • Simulace obchodování, která ukazuje zisk nebo ztrátu algoritmu obchodování.

Informační kanál nabídek v reálném čase

Important

Rozhraní IEX trading WebSocket API (iextrading.com) odkazované v této části bylo vyřazeno. IEX Cloud teď poskytuje data trhu prostřednictvím IEX Cloudu s různými ověřováními a koncovými body. Odpovídajícím způsobem aktualizujte adresu URL a ověřování v implementaci.

Important

Balíčky NuGet SocketIoClientDotNet a WindowsAzure.ServiceBus použité v této ukázce jsou zastaralé. Pro nové projekty použijte aktuální klientskou knihovnu Socket.IO a balíček Azure.Messaging.EventHubs s EventHubProducerClient místo starší verze EventHubClient.

Investors Exchange (IEX) dříve nabízela bezplatné kotace nákupních a prodejních cen v reálném čase prostřednictvím socket.io. Můžete napsat jednoduchý konzolový program, který bude přijímat nabídky v reálném čase a nabízet je do Azure Event Hubs jako zdroj dat. Následující kód je kostra programu. Kód vynechá zpracování chyb kvůli stručnosti. Do projektu musíte také zahrnout balíčky SocketIoClientDotNetNuGetWindowsAzure.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);
});

Caution

Tento vzorový kód je určen pouze pro ilustraci. Koncový bod rozhraní API IEX WebSocket a zde použité balíčky NuGet už nejsou k dispozici. Nepoužívejte tento kód v produkčním prostředí. Aktuální alternativy najdete v poznámkách DŮLEŽITÉ výše v této části.

Tady jsou některé vygenerované ukázkové události:

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

Časové razítko události je ve formátu epoch time: lastUpdated.

Prediktivní model pro vysokofrekvenční obchodování

Pro tuto ukázku se v ukázkovém příkladu používá lineární model popsaný v Strategii založené na nerovnováze objednávek při vysokofrekvenčním algoritmickém obchodování.

Objemová nerovnováha objednávek (VOI) je funkcí aktuální nákupní a prodejní ceny a jejich objemu a také nákupní a prodejní ceny a jejich objemu z posledního tiku. Tento dokument identifikuje korelaci mezi VOI a budoucím pohybem cen. Vytváří lineární model mezi posledními pěti hodnotami VOI a změnou ceny v následujících 10 ticích. Model se trénuje na datech z předchozího dne pomocí lineární regrese.

Vytrénovaný model pak v reálném čase predikuje změny cen kotací během aktuálního obchodního dne. Když model předpovídá dostatečně velkou změnu ceny, provede obchod. V závislosti na nastavení prahové hodnoty může jedna akcie během obchodního dne generovat tisíce obchodů.

Diagram znázorňující vzorec nerovnováhy pořadí objemů používaný při vysokofrekvenčním obchodování

Následující části ukazují, jak definovat operace trénování a predikce v úloze Azure Stream Analytics. Úplný dotaz je jediný WITH příkaz složený z běžných výrazů tabulek (CTE), které tvoří kanál:

Fáze CTE Purpose
typeconvertedquotes Převod nezpracovaných vstupních polí na správné typy SQL
timefilteredquotes Filtrování uvozovek na obchodní hodiny a odebrání neplatných dat
shiftedquotes Použijte LAG k načtení hodnot bid/ask z předchozího ticku
currentPriceAndVOI Vypočítejte objemovou nerovnováhu objednávek (VOI) z aktuálního a předchozího ticku
shiftedPriceAndShiftedVOI Sestavení sekvencí 10 po sobě jdoucích středních cen a 2 po sobě jdoucích hodnot VOI
modelInput Změna tvaru dat na vektory funkcí (VOI jako x, price delta jako y)
modelagg / modelparambs / model Trénování modelu lineární regrese se dvěma proměnnými pomocí agregací SUMa a AVG
shiftedVOI / VOIAndModel / VOIANDModelJoined Spojit aktuální hodnoty VOI s modelem natrénovaným předchozí den
prediction Výpočet očekávané budoucí změny ceny (efpc) z modelu
tradeSignal Generování signálů nákupu/prodeje, když efpc překročí prahovou hodnotu ±0,02

Note

Tento dotaz vyžaduje úroveň kompatibility 1.1 nebo novější služby Azure Stream Analytics, která zachovává použití velkých a malých písmen v názvech polí kvůli předvídatelnému chování funkcí UDA.

Vyčištění a převod vstupních polí nabídky

První CTE v dotazu Azure Stream Analytics převádí nezpracovaná data nabídek ze služby Event Hubs na správně typované sloupce SQL. FUNKCE DATEADD převede epochový čas (Unix milisekundy) na datetime. TRY_CAST převede datové typy bez selhání dotazu. Přetypujte vstupní pole na očekávané datové typy, abyste se vyhnuli neočekávanému chování při manipulaci nebo porovnání polí.

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

Načtení předchozích hodnot ticku pomocí funkce LAG

Další CTE v dotazu Azure Stream Analytics využívá funkci LAG k získání nabídkové a poptávkové ceny a velikosti z předchozího tikového záznamu pro každý akciový symbol. Jedna hodina hodnoty LIMIT DURATION je zvolena libovolně. Vzhledem k četnosti nabídek můžete předchozí zaškrtnutí najít tak, že se podíváte zpět na hodinu.

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

Výpočet nerovnováhy pořadí svazků (VOI)

Další CTE vypočítá hodnotu VOI z dat o nabídce a poptávce z aktuálního a předchozího ticku. Dotaz vyfiltruje hodnoty null pro případy, kdy neexistuje žádná předchozí značka.

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

Vytváření sekvencí funkcí pro trénování modelu

Další CTE znovu použije lag k vytvoření sekvence se 2 po sobě jdoucími hodnotami VOI, následovanými 10 po sobě jdoucími středními cenami. Tyto sekvence tvoří trénovací data pro model lineární regrese.

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

Změna tvaru dat na vektory funkcí

Další CTE přetváří cenové a VOI sekvence na vektory příznaků pro dvouproměnný lineární model, v němž hodnoty VOI představují nezávislé proměnné (x1, x2) a průměrná budoucí změna cen je závislá proměnná (y). Události s neúplnými daty se odfiltrují.

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

Trénování lineárního regresního modelu pomocí funkce SUM a AVG

Protože Azure Stream Analytics nemá integrovanou lineární regresní funkci, dotaz používá SUM a AVG agreguje koeficienty (a, b1, b2) pro model lineární regrese se dvěma proměnnými. Model se denně přetrénuje pomocí 24hodinového přeskakujícího okna.

Diagram znázorňující matematický vzorec lineární regrese pro výpočetní koeficienty modelu

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

Určení aktuálních uvozovek pomocí modelu předchozího dne

Pokud chcete použít vytrénovaný lineární regresní model předchozího dne pro bodování aktuální události, dotaz spojí uvozovky s koeficienty modelu. Místo použití funkce JOIN dotaz používá funkci UNION ke kombinování událostí modelu a událostí uvozovek do jednoho datového proudu. Pak pomocí lag spáruje události s modelem předchozího dne, takže získáte přesně jednu shodu. Kvůli víkendu sahá dotaz tři dny zpětně (72 hodin). Pokud by se použilo jednoduché SPOJENÍ , získáte tři modely pro každou událost nabídky.

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

Generování obchodních signálů z předpovědí

Konečné hodnoty CTE vypočítají očekávanou budoucí změnu ceny (efpc) použitím vzorce lineární regrese (a + b1 * x1 + b2 * x2) a pak vygenerují signály nákupu a prodeje na základě prahové hodnoty ±0,02. Obchodní hodnota 10 znamená nákup. Obchodní hodnota -10 znamená prodej.

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

Otestování strategie obchodování pomocí simulace

Po vygenerování obchodních signálů otestujte, jak účinná je obchodní strategie, aniž byste obchodovali se skutečnými penězi.

Tento test používá UDA s posuvným oknem, které se posouvá každou minutu. Seskupení podle data a klauzule HAVING zajišťují, že okno zohledňuje pouze události, které náleží ke stejnému dni. V případě přeskakujícího okna mezi dvěma dny odděluje datum GROUP BY seskupení do předchozího dne a aktuálního dne. Klauzule HAVING odfiltruje okna, která končí v aktuálním dni, ale jsou seskupena podle předchozího dne.

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
)

UDA JavaScriptu inicializuje všechny akumulátory ve init funkci, vypočítá přechod stavu s každou událostí přidanou do okna a vrátí výsledky simulace na konci okna. Simulace v každém obchodu drží nebo prodává nakrátko 10 akcií jedné společnosti. Náklady na transakci jsou fixní $8. Následující tabulka uvádí čtyři obchodní akce, které UDA provádí:

Condition Signál Action Pozice za
Žádná aktuální pozice Koupit (10) Koupit k otevření Long
Žádná aktuální pozice Prodej (-10) Prodej k otevření (krátké) Short
Dlouhá pozice Prodej (-10) Prodej za účelem uzavření a prodej za účelem otevření (krátké) Short
Krátká pozice Koupit (10) Koupit k uzavření a pak koupit k otevření 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

Výstupní konektor Power BI pro Azure Stream Analytics je naplánován na vyřazení z provozu. Zvažte použití alternativních výstupních cílů, jako jsou Azure Data Explorer, Azure Synapse Analytics nebo úložiště dat, ke kterému se Power BI může připojit přes DirectQuery nebo import. Další informace najdete ve výstupu Azure Stream Analytics do Power BI.

Nakonec data odešlete na řídicí panel Power BI k vizualizaci.

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, který zobrazuje obchodní signály vizualizované na řídicím panelu Power BI pro simulaci obchodování.

Chart zobrazící výsledky zisku a ztráty vizualizované na řídicím panelu Power BI pro simulaci obchodování.

Shrnutí

Tento článek ukazuje, jak implementovat realistický model vysokofrekvenčního obchodování se středně složitým dotazem v Azure Stream Analytics. Model používá místo pěti vstupních proměnných dvě vstupní proměnné, protože Azure Stream Analytics neobsahuje integrovanou lineární regresní funkci. Můžete ale také implementovat sofistikovanější algoritmy s vyššími rozměry jako UDA JavaScriptu.

Většinu dotazu kromě UDA JavaScriptu můžete testovat a ladit pomocí nástrojů Azure Stream Analytics pro Visual Studio Code pro vývoj dotazů, testování a ladění.