Notitie
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen u aan te melden of de directory te wijzigen.
Voor toegang tot deze pagina is autorisatie vereist. U kunt proberen de mappen te wijzigen.
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
- Een Azure-abonnement. Als u nog geen account hebt, kunt u een gratis account maken.
- Een Azure Stream Analytics-taak.
- Een Azure Event Hubs naamruimte en event hub.
- Bekendheid met Stream Analytics Query Language.
- (Optioneel) Een Power BI-account als u de uitvoer wilt visualiseren.
Werkstroom voor high-frequency trading
Het logische werkingsverloop van high-frequency trading is:
- Realtimekoersen ophalen van een effectenbeurs.
- Het bouwen van een voorspellend model rond de koersen om te anticiperen op de prijsverplaatsing.
- 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.
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.
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 */
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.