Remarque
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de vous connecter ou de modifier des répertoires.
L’accès à cette page nécessite une autorisation. Vous pouvez essayer de modifier des répertoires.
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
- Un abonnement Azure. Si vous n’en avez pas, créez un compte gratuit.
- Un travail Azure Stream Analytics.
- Un espace de noms Azure Event Hubs et un hub d’événements.
- Connaissance du langage de requête Stream Analytics.
- (Facultatif) Un compte Power BI si vous souhaitez visualiser la sortie.
Flux de travail de trading à haute fréquence
Le flux logique du trading à haute fréquence est le suivant :
- Obtention des cours en temps réel depuis une bourse de valeurs.
- Création d’un modèle prédictif autour des devis pour anticiper le mouvement des prix.
- 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.
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.
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 */
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.