Simulatie van high-frequency trading met Stream Analytics

Azure Stream Analytics ondersteunt geavanceerde analyses via de combinatie van SQL-taal, door de gebruiker gedefinieerde JavaScript-functies (UDF's) en door de gebruiker gedefinieerde aggregaties (UDF's). Geavanceerde analyses omvatten online training en scoring van machine-learningmodellen en statusbewuste processimulatie. In dit artikel wordt beschreven hoe u lineaire regressie kunt uitvoeren in een Azure Stream Analytics-job die continue training en scoring uitvoert in een scenario voor hoogfrequente handel.

Prerequisites

Werkstroom voor high-frequency trading

Het logische werkingsverloop van high-frequency trading is:

  1. Realtimekoersen ophalen van een effectenbeurs.
  2. Het bouwen van een voorspellend model rond de koersen om te anticiperen op de prijsverplaatsing.
  3. Het plaatsen van koop- of verkooporders om geld te verdienen aan de succesvolle voorspelling van de prijsbewegingen.

Voor dit scenario is het volgende vereist:

  • Een realtime offertefeed.
  • Een predictief model dat kan werken op basis van realtime koersnoteringen.
  • Een handelssimulatie die de winst of het verlies van het handelsalgoritmen laat zien.

Realtime offertefeed

Important

De IEX trading WebSocket-API (iextrading.com) waarnaar in deze sectie wordt verwezen, is buiten gebruik gesteld. IEX Cloud biedt nu marktgegevens via IEX Cloud met verschillende verificatie en eindpunten. Werk de URL en verificatie in uw implementatie dienovereenkomstig bij.

Important

De NuGet-pakketten SocketIoClientDotNet en WindowsAzure.ServiceBus die in dit voorbeeld worden gebruikt, zijn verouderd. Gebruik voor nieuwe projecten een huidige Socket.IO-clientbibliotheek en het Azure.Messaging.EventHubs-pakket met EventHubProducerClient in plaats van de verouderde EventHubClient.

Investors Exchange (IEX) bood voorheen gratis real-time bied- en laatkoersen aan via socket.io. U kunt een eenvoudig consoleprogramma schrijven om realtime koersnoteringen te ontvangen en deze naar Azure Event Hubs te sturen als gegevensbron. De volgende code is een skelet van het programma. De code laat de foutafhandeling weg voor beknoptheid. U moet ook de SocketIoClientDotNetWindowsAzure.ServiceBus NuGet-pakketten in uw project opnemen.

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

Dit codevoorbeeld is alleen ter illustratie. Het IEX WebSocket-API-eindpunt en de NuGet-pakketten die hier worden gebruikt, zijn niet meer beschikbaar. Gebruik deze code niet in productie. Zie de belangrijke opmerkingen eerder in deze sectie voor de huidige alternatieven.

Hier volgen enkele gegenereerde voorbeeldgebeurtenissen:

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

Het tijdstempel van de gebeurtenis is lastUpdated, in epochetijd.

Voorspellend model voor high-frequency trading

Voor deze demonstratie gebruikt het voorbeeld een lineair model, beschreven in Strategie op basis van orderonevenwicht in hoogfrequente algoritmische handel.

Onevenwicht in het ordervolume (VOI) is een functie van de huidige bied-/laatprijs en het volume, en van de bied-/laatprijs en het volume van de vorige tick. Het document identificeert de correlatie tussen VOI en toekomstige prijsbewegingen. Het stelt een lineair model op op basis van de afgelopen vijf VOI-waarden en de prijsverandering in de volgende 10 ticks. Het model traint op de gegevens van de vorige dag met lineaire regressie.

Het getrainde model maakt vervolgens prijswijzigingsvoorspellingen op koersen in de huidige handelsdag in realtime. Wanneer het model een grote prijswijziging voorspelt, wordt er een transactie uitgevoerd. Afhankelijk van de drempelwaarde kan één aandeel duizenden transacties genereren tijdens een handelsdag.

Diagram dat de definitieformule voor onevenwicht in ordervolume toont die wordt gebruikt in hoogfrequente handel.

In de volgende secties ziet u hoe u de trainings- en voorspellingsbewerkingen in een Azure Stream Analytics taak kunt uitdrukken. De volledige query is één WITH instructie die bestaat uit algemene tabelexpressies (CTE's) die een pijplijn vormen:

CTE-fase Purpose
typeconvertedquotes Onbewerkte invoervelden converteren naar de juiste SQL-typen
timefilteredquotes Koersnoteringen filteren op handelstijden en ongeldige data verwijderen
shiftedquotes Gebruik LAG om de bied-/laatwaarden van de vorige tick op te halen
currentPriceAndVOI Bereken het volume-orderonevenwicht (VOI) op basis van de huidige en vorige tick
shiftedPriceAndShiftedVOI Maak reeksen van 10 opeenvolgende middenprijzen en 2 opeenvolgende VOI-waarden
modelInput Gegevens omvormen tot kenmerkvectoren (VOI als x, prijsdelta als y)
modelagg / modelparambs / model Een lineair regressiemodel met twee variabelen trainen met SUM- en AVG-aggregaties
shiftedVOI / VOIAndModel / VOIANDModelJoined Huidige VOI-waarden samenvoegen met het getrainde model van de vorige dag
prediction Verwachte toekomstige prijswijziging (efpc) berekenen van het model
tradeSignal Koop-/verkoopsignalen genereren wanneer efpc de drempelwaarde voor ±0,02 overschrijdt

Note

Voor deze query is Azure Stream Analytics-compatibiliteitsniveau 1.1 of hoger vereist, waarbij het gebruik van hoofdletters en kleine letters in veldnamen behouden blijft voor voorspelbaar gedrag met UDA's.

Invoervelden voor offertes opschonen en converteren

De eerste CTE in de Azure Stream Analytics-query converteert de onbewerkte aanhalingsgegevens van Event Hubs naar correct getypte SQL-kolommen. DATEADD converteert tijdsduur (Unix milliseconden) naar datum/tijd. TRY_CAST gegevenstypen coërt zonder dat de query mislukt. Converteer invoervelden naar de verwachte gegevenstypes om onverwacht gedrag bij bewerking of vergelijking van de velden te voorkomen.

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

Vorige tickwaarden ophalen met LAG

De volgende CTE in de Azure Stream Analytics-query gebruikt de functie LAG om de bied-/vraagprijs en omvang van de vorige tick voor elk aandeelsymbool op te halen. Voor LIMIT DURATION is willekeurig gekozen voor een waarde van één uur. Gezien de frequentie van de noteringen kunt u de vorige tick vinden door één uur terug te kijken.

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

Volume order onevenwichtigheid berekenen (VOI)

De volgende CTE berekent de VOI-waarde op basis van de bied-/laatgegevens van de huidige en vorige tick. Deze query filtert null-waarden uit in gevallen waarin geen vorige tick aanwezig is.

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

Functiereeksen bouwen voor modeltraining

De volgende CTE gebruikt LAG opnieuw om een reeks te maken met 2 opeenvolgende VOI-waarden, gevolgd door 10 opeenvolgende gemiddelde prijswaarden. Deze reeksen vormen de trainingsgegevens voor het lineaire regressiemodel.

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

Gegevens opnieuw vormgeven in functievectoren

De volgende CTE hervormt de prijs- en VOI-reeksen in functievectoren voor een lineair model met twee variabelen, waarbij VOI-waarden de onafhankelijke variabelen zijn (x1, x2) en de gemiddelde toekomstige prijswijziging de afhankelijke variabele (y). Gebeurtenissen met onvolledige gegevens worden uitgefilterd.

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

Het lineaire regressiemodel trainen met SUM en AVG

Omdat Azure Stream Analytics geen ingebouwde lineaire regressiefunctie heeft, gebruikt de query SUM en AVG aggregaties voor het berekenen van de coëfficiënten (a, b1, b2) voor het lineaire regressiemodel met twee variabelen. Het model wordt dagelijks opnieuw getraind met een 24-uurs tumblingvenster.

Diagram met de lineaire regressieberekeningsformule voor het berekenen van modelcoëfficiënten.

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

Beoordeel de huidige koersen met het model van de vorige dag

Als u het getrainde lineaire regressiemodel van de vorige dag wilt gebruiken om de huidige gebeurtenis te scoren, voegt de query de aanhalingstekens samen met de modelcoëfficiënten. In plaats van JOIN te gebruiken, gebruikt de query UNION om modelgebeurtenissen en aanhalingstekens in één stream te combineren. Vervolgens wordt LAG gebruikt om de gebeurtenissen te koppelen aan het model van de vorige dag, zodat u precies één overeenkomst krijgt. Vanwege het weekend kijkt de query drie dagen (72 uur) terug. Als er een eenvoudige JOIN werd gebruikt, krijgt u drie modellen voor elke offertegebeurtenis.

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

Handelssignalen genereren van voorspellingen

De laatste CTE's berekenen de verwachte toekomstige prijswijziging (efpc) door de lineaire regressieformule (a + b1 * x1 + b2 * x2) toe te passen en vervolgens koop-/verkoopsignalen te genereren op basis van een drempelwaarde van ±0,02. Een handelswaarde van 10 staat voor koop. Een handelswaarde van -10 betekent verkopen.

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

De handelsstrategie testen met een simulatie

Nadat u de handelssignalen hebt gegenereerd, test u hoe effectief de handelsstrategie is zonder dat u echt hoeft te handelen.

Deze test gebruikt een UDA met een verspringend venster dat elke minuut verspringt. De groepering op datum en de HAVING-clausule zorgen ervoor dat het venster alleen rekening houdt met gebeurtenissen die op dezelfde dag vallen. Voor een hoppingvenster over twee dagen scheidt de GROUP BY-datum de groepering in de vorige dag en de huidige dag. De HAVING-clausule filtert de vensters eruit die eindigen op de huidige dag, maar op de vorige dag zijn gegroepeerd.

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
)

De JavaScript-UDA initialiseert alle accumulators in de init functie, berekent de statusovergang met elke gebeurtenis die aan het venster wordt toegevoegd en retourneert de simulatieresultaten aan het einde van het venster. De simulatie houdt per transactie 10 aandelen van één aandeel aan of gaat er short in. De transactiekosten zijn plat $8. In de volgende tabel ziet u de vier handelsacties die de UDA uitvoert:

Condition Signaal Action Positie na afloop
Geen huidige holding Kopen (10) Kopen om een positie te openen Long
Geen huidige holding Verkoop (-10) Verkopen om te openen (kort) Short
Lange positie Verkoop (-10) Verkopen om een positie te sluiten en vervolgens verkopen om een shortpositie te openen Short
Korte positie Kopen (10) Kopen om een positie te sluiten en vervolgens kopen om een positie te openen 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

De Power BI uitvoerconnector voor Azure Stream Analytics is gepland voor buitengebruikstelling. Overweeg alternatieve uitvoerbestemmingen te gebruiken, zoals Azure Data Explorer, Azure Synapse Analytics of een gegevensarchief waarmee Power BI verbinding kan maken via DirectQuery of importeren. Zie Azure Stream Analytics-uitvoer voor Power BI voor meer informatie.

Ten slotte voert u uitvoer uit naar het Power BI dashboard voor visualisatie.

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 die handelssignalen weergeeft die zijn gevisualiseerd in een Power BI dashboard voor de handelssimulatie.

Chart waarin winst- en verliesresultaten worden weergegeven die zijn gevisualiseerd in een Power BI dashboard voor de tradingsimulatie.

Overzicht

In dit artikel wordt beschreven hoe u een realistisch high-frequency trading model implementeert met een redelijk complexe query in Azure Stream Analytics. Het model gebruikt twee invoervariabelen in plaats van vijf omdat Azure Stream Analytics geen ingebouwde lineaire regressiefunctie bevat. U kunt echter ook geavanceerdere algoritmen met hogere dimensies implementeren als JavaScript-UDF's.

U kunt de meeste query's, behalve de JavaScript-UDA, testen en fouten opsporen met behulp van Azure Stream Analytics-hulpprogramma's voor Visual Studio Code voor het ontwikkelen, testen en foutopsporing van query's.