Högfrekvent handelssimulering med Stream Analytics

Azure Stream Analytics stöder avancerad analys genom kombinationen av SQL-språk, Användardefinierade JavaScript-funktioner (UDF: er) och användardefinierade aggregeringar (UDA). Avancerad analys omfattar utbildning och bedömning av maskininlärning online samt tillståndskänslig processsimulering. Den här artikeln beskriver hur du utför linjär regression i ett Azure Stream Analytics jobb som utför kontinuerlig träning och bedömning i ett scenario med högfrekvent handel.

Förutsättningar

Arbetsflöde för högfrekvent handel

Det logiska flödet för handel med hög frekvens är:

  1. Hämta realtidskurser från en värdepappersbörs.
  2. Bygga en prediktiv modell kring noteringarna för att förutse prisrörelsen.
  3. Placera köp- eller säljorder för att tjäna pengar på den framgångsrika förutsägelsen av prisrörelserna.

Det här scenariot kräver:

  • Ett flöde för realtidscitat.
  • En förutsägelsemodell som kan användas med citattecken i realtid.
  • En handelssimulering som visar vinst eller förlust för handelsalgoritmen.

Kursflöde i realtid

Important

IEX Trading WebSocket API (iextrading.com) som refereras i det här avsnittet har dragits tillbaka. IEX Cloud tillhandahåller nu marknadsdata via IEX Cloud med olika autentiserings- och slutpunkter. Uppdatera URL:en och autentiseringen i implementeringen i enlighet med detta.

Important

NuGet-paketen SocketIoClientDotNet och WindowsAzure.ServiceBus som används i det här exemplet är inaktuella. För nya projekt använder du ett aktuellt Socket.IO-klientbibliotek och paketet Azure.Messaging.EventHubs med EventHubProducerClient i stället för den äldre EventHubClient.

Investors Exchange (IEX) erbjöd tidigare gratis realtidskurser för köp och sälj genom att använda socket.io. Du kan skriva ett enkelt konsolprogram för att ta emot citattecken i realtid och skicka dem till Azure Event Hubs som datakälla. Följande kod är ett skelett av programmet. Koden utelämnar felhantering för korthet. Du måste också inkludera NuGet-paketen SocketIoClientDotNet och WindowsAzure.ServiceBus i projektet.

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

Det här kodexemplet är endast till för bild. IEX WebSocket API-slutpunkten och NuGet-paketen som används här är inte längre tillgängliga. Använd inte den här koden i produktion. Se viktig information tidigare i det här avsnittet för aktuella alternativ.

Här följer några genererade exempelhändelser:

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

Tidsstämpeln för händelsen är lastUpdated, i epoktid.

Förutsägelsemodell för handel med hög frekvens

I den här demonstrationen använder exemplet en linjär modell som beskrivs i Order Imbalance Based Strategy in High Frequency Algorithmic Trading.

Obalans i volymordning (VOI) är en funktion av aktuellt bud/fråga pris och volym samt bud/fråga pris och volym från den senaste ticken. Dokumentet identifierar korrelationen mellan VOI och framtida prisrörelse. Den bygger en linjär modell mellan de senaste fem VOI-värdena och prisändringen under de kommande 10 ticken. Modellen tränar på föregående dags data med linjär regression.

Den tränade modellen gör sedan prisändringsförutsägelser på offerter under den aktuella handelsdagen i realtid. När modellen förutsäger en tillräckligt stor prisändring körs en handel. Beroende på tröskelvärdet kan en enskild aktie generera tusentals affärer under en handelsdag.

Diagram som visar den definitionsformel för obalans i volymordningen som används i högfrekvent handel.

Följande avsnitt visar hur du uttrycker tränings- och förutsägelseåtgärder i ett Azure Stream Analytics jobb. Den fullständiga frågan är en enda WITH instruktion som består av vanliga tabelluttryck (CTE) som utgör en pipeline:

CTE-steg Purpose
typeconvertedquotes Konvertera råa indatafält till rätt SQL-typer
timefilteredquotes Filtrera kurser till handelstid och ta bort ogiltiga data
shiftedquotes Använd LAG för att hämta föregående ticks bud-/fråga-värden
currentPriceAndVOI Beräkna obalans i volymordning (VOI) från aktuell och föregående tick
shiftedPriceAndShiftedVOI Byggsekvenser med 10 på varandra följande mellanpriser och 2 på varandra följande VOI-värden
modelInput Omforma data till funktionsvektorer (VOI som x, prisdelta som y)
modelagg / modelparambs / model Träna en linjär regressionsmodell med två variabler med SUM- och AVG-aggregat
shiftedVOI / VOIAndModel / VOIANDModelJoined Koppla aktuella VOI-värden till föregående dags tränade modell
prediction Beräkna förväntad framtida prisändring (efpc) från modellen
tradeSignal Generera köp-/säljsignaler när efpc överskrider tröskelvärdet ±0,02

Note

Den här frågan kräver Azure Stream Analytics kompatibilitetsnivå 1.1 eller senare, vilket bevarar fältnamnshöljet för förutsägbart beteende med UDA:er.

Rensa och konvertera indatafält för offert

Den första CTE i Azure Stream Analytics-frågan konverterar rådata från Event Hubs till korrekt skrivna SQL-kolumner. DATEADD konverterar epoktid (Unix millisekunder) till datetime. TRY_CAST framtvingar datatyper utan att misslyckas med frågan. Omvandla indatafält till de förväntade datatyperna för att undvika oväntat beteende vid manipulering eller jämförelse av fälten.

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

Hämta tidigare tickvärden med LAG

Nästa CTE i frågan i Azure Stream Analytics använder funktionen LAG för att hämta köp-/säljpriset och volymen från föregående tick för varje aktiesymbol. En timmes gränsvaraktighetsvärde väljs godtyckligt. Med tanke på offertfrekvensen kan du hitta föregående kryss genom att titta tillbaka en timme.

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

Beräkna obalans i volymordning (VOI)

Nästa CTE beräknar VOI-värdet utifrån bid-/ask-data från den aktuella och föregående ticken. Frågan filtrerar bort null-värden för fall där det inte finns någon tidigare bockmarkering.

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

Skapa funktionssekvenser för modellträning

Nästa CTE använder LAG igen för att skapa en sekvens med 2 på varandra följande VOI-värden, följt av 10 på varandra följande mellanprisvärden. Dessa sekvenser utgör träningsdata för den linjära regressionsmodellen.

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

Omforma data till funktionsvektorer

Nästa CTE omformar pris- och VOI-sekvenserna till funktionsvektorer för en linjär modell med två variabler, där VOI-värden är de oberoende variablerna (x1, x2) och den genomsnittliga framtida prisändringen är den beroende variabeln (y). Händelser med ofullständiga data filtreras bort.

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äna den linjära regressionsmodellen med SUM och AVG

Eftersom Azure Stream Analytics inte har någon inbyggd linjär regressionsfunktion använder frågan SUM och AVG aggregeringar för att beräkna koefficienterna (a, b1, b2) för den linjära regressionsmodellen med två variabler. Modellen tränas om dagligen med ett rullande 24-timmarsfönster.

Diagram som visar den linjära regressionsmatematiska formeln för beräkning av modellkoefficienter.

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

Poängsätta aktuella citattecken med föregående dags modell

Om du vill använda föregående dags tränade linjära regressionsmodell för att bedöma den aktuella händelsen ansluter frågan citattecknarna med modellkoefficienterna. I stället för att använda JOIN använder frågan UNION för att kombinera modellhändelser och offerthändelser till en enda ström. Sedan används LAG för att para ihop händelserna med föregående dags modell, så att du får exakt en matchning. På grund av helgen går sökningen tillbaka tre dagar (72 timmar). Om en enkel JOIN användes skulle du få tre modeller för varje offerthändelse.

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

Generera handelssignaler från förutsägelser

De slutliga CTE:erna beräknar den förväntade framtida prisändringen (efpc) genom att tillämpa den linjära regressionsformeln (a + b1 * x1 + b2 * x2) och sedan generera köp-/säljsignaler baserat på ett tröskelvärde på ±0,02. Ett handelsvärde på 10 är köp. Ett handelsvärde på -10 innebär sälj.

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

Testa handelsstrategin med en simulering

Efter att ha genererat handelssignalerna, testa hur effektiv handelsstrategin är utan handel på riktigt.

Det här testet använder en UDA med ett hoppfönster som hoppar var minut. Grupperingen på datum och HAVING-satsen säkerställer att fönstret endast omfattar händelser som tillhör samma dag. För ett hoppfönster över två dagar delar GROUP BY-datumet upp grupperingen i föregående dag och innevarande dag. HAVING-villkoret filtrerar bort de fönster som slutar den innevarande dagen men grupperas på föregående dag.

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 initierar alla ackumulatorer i init funktionen, beräknar tillståndsövergången med varje händelse som läggs till i fönstret och returnerar simuleringsresultatet i slutet av fönstret. I simuleringen äger eller blankar man 10 aktier i samma aktie i varje affär. Transaktionskostnaden är en fast $8. I följande tabell visas de fyra handelsåtgärder som UDA utför:

Tillstånd Signal Action Position efter
Inget aktuellt innehav Köp (10) Köp för att öppna Long
Inget aktuellt innehav Sälj (-10) Sälj för att öppna (kort) Short
Långt läge Sälj (-10) Sälj för att stänga och sedan sälja för att öppna (kort) Short
Kort position Köp (10) Köp för att stänga och köp sedan för att öppna 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-utdataanslutningen för Azure Stream Analytics kommer att dras tillbaka. Överväg att använda alternativa utdatamål som Azure Data Explorer, Azure Synapse Analytics eller ett datalager som Power BI kan ansluta till via DirectQuery eller importera. Mer information finns i Azure Stream Analytics utdata till Power BI.

Slutligen utdata till Power BI instrumentpanel för visualisering.

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 som visar handelssignaler som visualiserats i en Power BI instrumentpanel för handelssimuleringen.

Chart som visar resultat för vinst och förlust visualiserade i en Power BI instrumentpanel för handelssimuleringen.

Sammanfattning

Den här artikeln visar hur du implementerar en realistisk handelsmodell med hög frekvens med en måttligt komplex fråga i Azure Stream Analytics. Modellen använder två indatavariabler i stället för fem eftersom Azure Stream Analytics inte innehåller någon inbyggd linjär regressionsfunktion. Men du kan också implementera mer avancerade algoritmer med högre dimensioner som JavaScript UDAs.

Du kan testa och felsöka större delen av frågan, förutom JavaScript UDA, med hjälp av Azure Stream Analytics verktyg för Visual Studio Code för frågeutveckling, testning och felsökning.