Simulation de trading à haute fréquence avec Stream Analytics

Azure Stream Analytics prend en charge l’analytique avancée par le biais de la combinaison de langage SQL, de fonctions définies par l’utilisateur (UDF) JavaScript et d’agrégats définis par l’utilisateur (UUDA). Les analyses avancées peuvent inclure l’apprentissage automatique en ligne et la notation, ainsi que la simulation des processus avec état. Cet article explique comment effectuer une régression linéaire dans un travail Azure Stream Analytics qui effectue une formation continue et un scoring dans un scénario de trading à haute fréquence.

Prerequisites

Flux de travail de trading à haute fréquence

Le flux logique du trading à haute fréquence est le suivant :

  1. Obtention des cours en temps réel depuis une bourse de valeurs.
  2. Création d’un modèle prédictif autour des devis pour anticiper le mouvement des prix.
  3. Placer des commandes d’achat ou de vente pour gagner de l’argent à partir de la prédiction réussie des mouvements de prix.

Ce scénario nécessite :

  • Un flux de cotations en temps réel.
  • Un modèle prédictif capable de fonctionner à partir de cotations en temps réel.
  • Simulation de trading qui illustre le bénéfice ou la perte de l’algorithme de trading.

Flux de devis en temps réel

Important

L’APIiextrading.com WebSocket de trading IEX référencée dans cette section a été supprimée. IEX Cloud fournit désormais des données de marché via IEX Cloud avec différentes authentifications et points de terminaison. Mettez à jour l’URL et l’authentification dans votre implémentation en conséquence.

Important

Les packages NuGet SocketIoClientDotNet et WindowsAzure.ServiceBus utilisés dans cet exemple sont obsolètes. Pour les nouveaux projets, utilisez une bibliothèque de client Socket.IO actuelle et le package Azure.Messaging.EventHubs avec EventHubProducerClient au lieu du EventHubClient hérité.

Investors Exchange (IEX) proposait auparavant gratuitement des cotations acheteur et vendeur en temps réel via socket.io. Vous pouvez écrire un simple programme console pour recevoir des cotations en temps réel et les transmettre à Azure Event Hubs comme source de données. Le code suivant est un squelette du programme. Le code omet la gestion des erreurs pour la concision. Vous devez également inclure les packages NuGet SocketIoClientDotNet et WindowsAzure.ServiceBus dans votre projet.

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

Cet exemple de code est destiné uniquement à l’illustration. Le point de terminaison de l’API WebSocket IEX et les packages NuGet utilisés ici ne sont plus disponibles. N’utilisez pas ce code en production. Consultez les notes IMPORTANTES plus haut dans cette section pour connaître les alternatives actuelles.

Voici quelques exemples d’événements générés :

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

L’horodatage de l’événement est lastUpdated, en heure d’époque.

Modèle prédictif pour le trading à haute fréquence

Pour cette démonstration, l’exemple utilise un modèle linéaire décrit dans Stratégie basée sur le déséquilibre de commande dans la négociation algorithmique haute fréquence.

Le déséquilibre du volume de commande (VOI) est une fonction d’offres/demandes qui concerne les prix et le volume, elle s’applique aux offres/demandes en cours ou depuis le dernier cycle. Le document identifie la corrélation entre VOI et le mouvement futur des prix. Il crée un modèle linéaire entre les cinq dernières valeurs VOI et le changement de prix dans les 10 cycles suivants. Le modèle s’entraîne sur les données du jour précédent avec une régression linéaire.

Le modèle entraîné génère ensuite en temps réel des prévisions de variation des prix à partir des cotations de la séance de bourse en cours. Lorsque le modèle prédit un changement de prix suffisamment important, il exécute un commerce. Selon le paramètre de seuil, un seul stock peut générer des milliers de transactions pendant une journée de négociation.

Diagramme montrant la formule de définition de déséquilibre des commandes en volume utilisée dans le trading à haute fréquence.

Les sections suivantes montrent comment exprimer les opérations d’entraînement et de prédiction dans un travail Azure Stream Analytics. La requête complète est une instruction unique WITH composée d’expressions de table communes (CTEs) qui forment un pipeline :

Étape CTE Purpose
typeconvertedquotes Convertir les champs d’entrée brutes en types SQL appropriés
timefilteredquotes Filtrer les cotations en fonction des heures de négociation et supprimer les données invalides
shiftedquotes Utiliser LAG pour récupérer les valeurs d’offre/demande du cycle précédent
currentPriceAndVOI Calculer le déséquilibre des commandes volumineuses (VOI) à partir du cycle actuel et du précédent
shiftedPriceAndShiftedVOI Construisez des séquences de 10 prix médians consécutifs et de 2 valeurs VOI consécutives
modelInput Remodeler les données en vecteurs de caractéristiques (VOI as x, price delta as y)
modelagg / modelparambs / model Entraîner un modèle de régression linéaire à deux variables à l’aide des agrégats SUM et AVG
shiftedVOI / VOIAndModel / VOIANDModelJoined Associer les valeurs VOI actuelles au modèle entraîné la veille
prediction Calculer le changement de prix futur attendu (efpc) à partir du modèle
tradeSignal Générer des signaux d’achat/vente lorsque efp dépasse le seuil ±0,02

Note

dCette requête nécessite le niveau de compatibilité 1.1 ou ultérieur d’Azure Stream Analytics, qui préserve la casse des noms de champs pour un comportement prévisible avec les agrégats définis par l’utilisateur.

Nettoyer et convertir des champs d’entrée de devis

La première expression de table commune dans la requête Azure Stream Analytics convertit les données de cotation brutes d’Event Hubs en colonnes SQL correctement typées. DATEADD convertit l’heure d’époque (millisecondes Unix) en datetime. TRY_CAST convertit les types de données sans faire échouer la requête. Cassez les champs d’entrée vers les types de données attendus pour éviter un comportement inattendu dans la manipulation ou la comparaison des champs.

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

Récupérer les valeurs du cycle précédent avec LAG

L’expression de table commune dans la requête Azure Stream Analytics utilise la fonction LAG pour obtenir le prix de l’offre/demande et la taille du cycle précédent pour chaque symbole de stock. Une heure de la valeur LIMIT DURATION est arbitrairement choisie. Compte tenu de la fréquence des devis, vous pouvez trouver le cycle précédent en remontant d’une heure.

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

Calculer le déséquilibre de la commande volumineuse (VOI)

L’expression de table commune suivante calcule la valeur VOI à partir des données de l’offre/demande du cycle actuel et du précédent. La requête filtre les valeurs nulles lorsqu’aucun tick précédent n’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
),

Générer des séquences de fonctionnalités pour l’entraînement du modèle

Le CTE suivant utilise à nouveau LAG pour générer une séquence comportant 2 valeurs VOI consécutives, suivies de 10 valeurs de prix médian consécutives. Ces séquences forment les données d’apprentissage du modèle de régression linéaire.

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

Remodeler les données en vecteurs de caractéristiques

La CTE suivante transforme les séquences de prix et de VOI en vecteurs de caractéristiques pour un modèle linéaire à deux variables, où les valeurs de VOI sont les variables indépendantes (x1, x2) et la variation moyenne future du prix est la variable dépendante (y). Les événements avec des données incomplètes sont filtrés.

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

Entraîner le modèle de régression linéaire avec SUM et AVG

Étant donné que Azure Stream Analytics n'a pas de fonction de régression linéaire intégrée, la requête utilise des agrégats SUM et AVG pour calculer les coefficients (a, b1, b2) pour le modèle de régression linéaire à deux variables. Le modèle réentraîne quotidiennement à l’aide d’une fenêtre bascule de 24 heures.

Diagramme montrant la formule mathématique de régression linéaire pour les coefficients de modèle de calcul.

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

Évaluer les cotations actuelles à l’aide du modèle de la veille

Pour utiliser le modèle de régression linéaire entraînée du jour précédent pour noter l’événement actuel, la requête joint les guillemets aux coefficients du modèle. Au lieu d’utiliser JOIN, la requête utilise UNION pour combiner des événements de modèle et des événements de citation dans un seul flux. Ensuite, il utilise le lag pour associer les événements au modèle du jour précédent, de sorte que vous obtenez exactement une correspondance. En raison du week-end, la requête porte sur les trois jours précédents (72 heures). En utilisant simplement JOIN, nous obtiendrions trois modèles pour chaque événement de devis.

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

Générer des signaux commerciaux à partir de prédictions

Les CTE finales calculent le changement de prix futur attendu (efpc) en appliquant la formule de régression linéaire (a + b1 * x1 + b2 * x2) puis en générant des signaux d’achat/vente en fonction d’un seuil ±0,02. Une valeur de transaction de 10 correspond à un achat. Une valeur de transaction de -10 correspond à une vente.

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

Tester la stratégie de trading avec une simulation

Après avoir généré les signaux commerciaux, testez l’efficacité de la stratégie commerciale sans trading pour de vrai.

Ce test utilise un agrégat défini par l’utilisateur avec une fenêtre récurrente qui saute toutes les minutes. Le regroupement à la date et la clause HAVING garantissent que la fenêtre ne tient compte que des événements appartenant au même jour. Pour une fenêtre glissante sur deux jours, la date dans GROUP BY répartit le regroupement entre le jour précédent et le jour en cours. La clause HAVING filtre les fenêtres se terminant durant le jour en cours, mais regroupe en fonction de la journée précédente.

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
)

L’UDA JavaScript initialise tous les accumulateurs dans la fonction init, calcule la transition d’état à chaque événement ajouté à la fenêtre et renvoie les résultats de la simulation à la fin de la fenêtre. La simulation détient ou vend à découvert 10 actions d’un titre par transaction. Le coût de la transaction est forfaitaire : $8. Le tableau suivant présente les quatre actions commerciales effectuées par l’UDA :

Condition Signal Action Position après
Aucune position actuelle Acheter (10) Acheter pour ouvrir Long
Aucune position actuelle Vendre (-10) Vendre pour ouvrir (court) Short
Position longue Vendre (-10) Vendre pour fermer, puis vendre pour ouvrir (court) Short
Position courte Acheter (10) Acheter pour clôturer, puis acheter pour ouvrir 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

Le connecteur de sortie Power BI pour Azure Stream Analytics est planifié pour la mise hors service. Envisagez d’utiliser d’autres destinations de sortie telles que Azure Data Explorer, Azure Synapse Analytics ou un magasin de données auquel Power BI pouvez se connecter via DirectQuery ou l’importation. Pour plus d’informations, consultez Sortie Azure Stream Analytics vers Power BI.

Enfin, exportez les données vers le tableau de bord Power BI pour la visualisation.

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 */

Graphique illustrant les signaux de trading visualisés dans un tableau de bord Power BI pour la simulation de trading.

Chart qui montre les résultats des bénéfices et des pertes visualisées dans un tableau de bord Power BI pour la simulation de trading.

Résumé

Cet article montre comment implémenter un modèle de trading à fréquence élevée réaliste avec une requête modérément complexe dans Azure Stream Analytics. Le modèle utilise deux variables d'entrée au lieu de cinq, car Azure Stream Analytics n'inclut pas de fonction de régression linéaire intégrée. Toutefois, vous pouvez également implémenter des algorithmes plus sophistiqués avec des dimensions plus élevées en tant qu’UUDAs JavaScript.

Vous pouvez tester et déboguer la plupart de la requête, autre que javaScript UDA, à l’aide d’outils Azure Stream Analytics pour Visual Studio Code pour le développement, le test et le débogage des requêtes.