I guess that you allready know that FILE_SHARE_READ buffering—is an architectural anti-pattern. MetaTrader runs OnTick() synchronously. If the disk controller stalls, or if you are running 20 charts across 5 brokers on your HP Z1 workstation, your I/O thread will eventually bottleneck and miss ticks during high-volatility news events.
The true enterprise architecture decouples the MetaTrader client from the storage layer. We will turn MT4/MT5 into a pure, stateless telemetry sensor.
This version queues ticks in a memory buffer (an array of structs) and flushes them as a batched JSON payload via HTTP POST to a local ASP.NET Core microservice. Your C# backend can then handle the database ingestion (e.g., to MS SQL Server via Entity Framework Core) asynchronously.
How to ask support for historical spread data and what they send
Re: How to ask support for historical spread data and what they send
Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
Re: How to ask support for historical spread data and what they send
1. The Universal MT4/MT5 Telemetry Sensor
This EA builds a raw JSON array in memory to avoid external MQL JSON library dependencies, then uses MetaTrader's native WebRequest to fire the batch to your C# API.
Note: You must add http://localhost:5000 to the allowed WebRequest URLs in MetaTrader (Tools -> Options -> Expert Advisors).
This EA builds a raw JSON array in memory to avoid external MQL JSON library dependencies, then uses MetaTrader's native WebRequest to fire the batch to your C# API.
Note: You must add http://localhost:5000 to the allowed WebRequest URLs in MetaTrader (Tools -> Options -> Expert Advisors).
Code: Select all
//+------------------------------------------------------------------+
//| EnterpriseSpreadSensor.mq4/5 |
//| Stateless JSON Telemetry over HTTP POST |
//+------------------------------------------------------------------+
#property copyright "Copyright 2026"
#property strict
//--- Inputs
input string InpApiEndpoint = "http://localhost:5000/api/ticks"; // C# Microservice URL
input int InpBatchSize = 50; // Ticks per JSON payload
input string InpSessionStart = "07:00:00"; // UTC Start
input string InpSessionEnd = "10:00:00"; // UTC End
//--- Tick Buffer Struct
struct TickRecord {
ulong time_msc;
double bid;
double ask;
double spread_pips;
};
TickRecord g_tickQueue[];
int g_queueCount = 0;
double g_lastBid = 0.0;
double g_lastAsk = 0.0;
//+------------------------------------------------------------------+
int OnInit()
{
ArrayResize(g_tickQueue, InpBatchSize);
Print("[Telemetry Sensor] Initialized. Target: ", InpApiEndpoint);
return(INIT_SUCCEEDED);
}
//+------------------------------------------------------------------+
void OnTick()
{
MqlTick tick;
if(!SymbolInfoTick(_Symbol, tick)) return;
// Filter redundant quotes
if(tick.bid == g_lastBid && tick.ask == g_lastAsk) return;
datetime gmtNow = TimeGMT();
if(!IsWithinSession(gmtNow, InpSessionStart, InpSessionEnd)) return;
double pipSize = GetPipUnit();
double spreadPips = (pipSize > 0) ? (tick.ask - tick.bid) / pipSize : 0.0;
// Enqueue Tick
g_tickQueue[g_queueCount].time_msc = tick.time_msc;
g_tickQueue[g_queueCount].bid = tick.bid;
g_tickQueue[g_queueCount].ask = tick.ask;
g_tickQueue[g_queueCount].spread_pips = spreadPips;
g_queueCount++;
g_lastBid = tick.bid;
g_lastAsk = tick.ask;
// Flush Batch over HTTP
if(g_queueCount >= InpBatchSize)
{
FlushTelemetry();
g_queueCount = 0;
}
}
//+------------------------------------------------------------------+
//| HTTP POST & JSON Serialization |
//+------------------------------------------------------------------+
void FlushTelemetry()
{
// Manually build JSON array to avoid external DLLs/libraries in MT
string json = "[";
for(int i = 0; i < g_queueCount; i++)
{
string record = StringFormat("{\"symbol\":\"%s\",\"time_msc\":%I64u,\"bid\":%.5f,\"ask\":%.5f,\"spread\":%.2f}",
_Symbol, g_tickQueue[i].time_msc,
g_tickQueue[i].bid, g_tickQueue[i].ask, g_tickQueue[i].spread_pips);
json += record;
if(i < g_queueCount - 1) json += ",";
}
json += "]";
char postData[];
StringToCharArray(json, postData, 0, WHOLE_ARRAY, CP_UTF8);
// Strip trailing null terminator injected by StringToCharArray
int dataSize = ArraySize(postData) - 1;
char result[];
string resultHeaders;
string headers = "Content-Type: application/json\r\n";
int res = WebRequest("POST", InpApiEndpoint, headers, 1000, postData, dataSize, result, resultHeaders);
if(res != 200)
{
Print(StringFormat("[Telemetry Error] HTTP %d. Is the C# backend running?", res));
}
}
//+------------------------------------------------------------------+
//| Utilities |
//+------------------------------------------------------------------+
double GetPipUnit() { return (_Digits == 3 || _Digits == 5) ? _Point * 10.0 : _Point; }
bool IsWithinSession(datetime time, string startStr, string endStr)
{
datetime sod = StringToTime(TimeToString(time, TIME_DATE));
int currentSec = (int)(time - sod);
string sParts[], eParts[];
StringSplit(startStr, ':', sParts);
StringSplit(endStr, ':', eParts);
int startSec = (int)StringToInteger(sParts[0])*3600 + (int)StringToInteger(sParts[1])*60 + (int)StringToInteger(sParts[2]);
int endSec = (int)StringToInteger(eParts[0])*3600 + (int)StringToInteger(eParts[1])*60 + (int)StringToInteger(eParts[2]);
return (startSec <= endSec) ? (currentSec >= startSec && currentSec <= endSec)
: (currentSec >= startSec || currentSec <= endSec);
}Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
Re: How to ask support for historical spread data and what they send
2. The ASP.NET Core "Catcher" (High-Performance Ingestion)
To handle the incoming HTTP POSTs without blocking, we use an ASP.NET Core Minimal API coupled with System.Threading.Channels. This creates an in-memory queue on the C# side, allowing the API to instantly return HTTP 200 to MetaTrader while a background service bulk-inserts the data into MS SQL Server or a master CSV.
Create a new .NET 8 Web API project and replace Program.cs:
To handle the incoming HTTP POSTs without blocking, we use an ASP.NET Core Minimal API coupled with System.Threading.Channels. This creates an in-memory queue on the C# side, allowing the API to instantly return HTTP 200 to MetaTrader while a background service bulk-inserts the data into MS SQL Server or a master CSV.
Create a new .NET 8 Web API project and replace Program.cs:
Code: Select all
using System.Threading.Channels;
var builder = WebApplication.CreateBuilder(args);
// Create an unbounded channel for high-throughput, non-blocking ingestion
var tickChannel = Channel.CreateUnbounded<TickRecord>();
builder.Services.AddSingleton(tickChannel.Writer);
builder.Services.AddSingleton(tickChannel.Reader);
builder.Services.AddHostedService<TickIngestionWorker>();
var app = builder.Build();
// Minimal API Endpoint - O(1) execution time, instantly frees the MT4/MT5 thread
app.MapPost("/api/ticks", async (TickRecord[] payload, ChannelWriter<TickRecord> writer) =>
{
foreach (var tick in payload)
{
await writer.WriteAsync(tick);
}
return Results.Ok();
});
app.Run("http://localhost:5000");
// --- Models ---
public record TickRecord(string Symbol, ulong Time_msc, double Bid, double Ask, double Spread);
// --- Background Worker (Decoupled Database Writer) ---
public class TickIngestionWorker : BackgroundService
{
private readonly ChannelReader<TickRecord> _reader;
private readonly ILogger<TickIngestionWorker> _logger;
public TickIngestionWorker(ChannelReader<TickRecord> reader, ILogger<TickIngestionWorker> logger)
{
_reader = reader;
_logger = logger;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
_logger.LogInformation("Tick Ingestion Worker started.");
// Read dynamically as ticks arrive from the channel
await foreach (var tick in _reader.ReadAllAsync(stoppingToken))
{
// TODO: Map to Entity Framework Core -> MS SQL Server
// or batch write to a daily partitioned CSV file.
_logger.LogInformation("[{Symbol}] Spread: {Spread} | Bid: {Bid}", tick.Symbol, tick.Spread, tick.Bid);
}
}
}Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
Re: How to ask support for historical spread data and what they send
Here is the pure MQL4 implementation. Because MT4 Expert Advisors are strictly single-threaded, WebRequest blocks the entire EA execution until the HTTP response is received. This makes the decoupled architecture even more critical for MT4 than MT5—your C# backend must ingest the payload and return an HTTP status immediately.
1. The MT4 Sensor (SpreadTelemetry.mq4)
This script is compiled specifically for the MetaEditor 4 compiler. It builds the JSON payload natively to avoid relying on external MT4 DLLs that often break across terminal updates.
1. The MT4 Sensor (SpreadTelemetry.mq4)
This script is compiled specifically for the MetaEditor 4 compiler. It builds the JSON payload natively to avoid relying on external MT4 DLLs that often break across terminal updates.
Code: Select all
//+------------------------------------------------------------------+
//| SpreadTelemetry.mq4 |
//| Stateless JSON Telemetry over HTTP POST |
//+------------------------------------------------------------------+
#property copyright "Enterprise Telemetry"
#property strict
//--- Inputs
input string InpApiEndpoint = "http://localhost:5000/api/ticks";
input int InpBatchSize = 50;
input string InpSessionStart = "07:00:00";
input string InpSessionEnd = "10:00:00";
//--- Structs & Globals
struct TickRecord {
long time_msc;
double bid;
double ask;
double spread_pips;
};
TickRecord g_tickQueue[];
int g_queueCount = 0;
double g_lastBid = 0.0;
double g_lastAsk = 0.0;
//+------------------------------------------------------------------+
int OnInit()
{
ArrayResize(g_tickQueue, InpBatchSize);
Print("[MT4 Telemetry] Initialized. Target: ", InpApiEndpoint);
return(INIT_SUCCEEDED);
}
//+------------------------------------------------------------------+
void OnTick()
{
MqlTick tick;
if(!SymbolInfoTick(_Symbol, tick)) return;
if(tick.bid == g_lastBid && tick.ask == g_lastAsk) return;
datetime gmtNow = TimeGMT();
if(!IsWithinSession(gmtNow, InpSessionStart, InpSessionEnd)) return;
double pipSize = (_Digits == 3 || _Digits == 5) ? _Point * 10.0 : _Point;
double spreadPips = (pipSize > 0) ? (tick.ask - tick.bid) / pipSize : 0.0;
// Enqueue
g_tickQueue[g_queueCount].time_msc = tick.time_msc;
g_tickQueue[g_queueCount].bid = tick.bid;
g_tickQueue[g_queueCount].ask = tick.ask;
g_tickQueue[g_queueCount].spread_pips = spreadPips;
g_queueCount++;
g_lastBid = tick.bid;
g_lastAsk = tick.ask;
if(g_queueCount >= InpBatchSize)
{
FlushTelemetry();
g_queueCount = 0;
}
}
//+------------------------------------------------------------------+
void FlushTelemetry()
{
string json = "[";
for(int i = 0; i < g_queueCount; i++)
{
string record = StringFormat("{\"symbol\":\"%s\",\"time_msc\":%I64d,\"bid\":%.5f,\"ask\":%.5f,\"spread\":%.2f}",
_Symbol, g_tickQueue[i].time_msc,
g_tickQueue[i].bid, g_tickQueue[i].ask, g_tickQueue[i].spread_pips);
json += record;
if(i < g_queueCount - 1) json += ",";
}
json += "]";
char postData[];
StringToCharArray(json, postData, 0, WHOLE_ARRAY, CP_UTF8);
// MT4 StringToCharArray appends a null terminator, strip it for clean JSON
int dataSize = ArraySize(postData) - 1;
char result[];
string resultHeaders;
string headers = "Content-Type: application/json\r\n";
// Timeout set to 500ms to prevent MT4 thread locking
int res = WebRequest("POST", InpApiEndpoint, headers, 500, postData, dataSize, result, resultHeaders);
if(res != 200)
{
if(res == 4014)
Print("[MT4 Telemetry] ERR 4014: Add ", InpApiEndpoint, " to allowed WebRequest URLs in Options -> Expert Advisors.");
else
Print("[MT4 Telemetry] HTTP Error: ", res);
}
}
//+------------------------------------------------------------------+
bool IsWithinSession(datetime time, string startStr, string endStr)
{
datetime sod = StringToTime(TimeToString(time, TIME_DATE));
int currentSec = (int)(time - sod);
string sParts[], eParts[];
StringSplit(startStr, ':', sParts);
StringSplit(endStr, ':', eParts);
int startSec = (int)StringToInteger(sParts[0])*3600 + (int)StringToInteger(sParts[1])*60 + (int)StringToInteger(sParts[2]);
int endSec = (int)StringToInteger(eParts[0])*3600 + (int)StringToInteger(eParts[1])*60 + (int)StringToInteger(eParts[2]);
return (startSec <= endSec) ? (currentSec >= startSec && currentSec <= endSec)
: (currentSec >= startSec || currentSec <= endSec);
}Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
Re: How to ask support for historical spread data and what they send
2. The ASP.NET Core Ingestion Service (EF Core & MS SQL)
Because MT4 will freeze if the HTTP response is delayed, the C# API acts solely as a high-speed router. It accepts the JSON, drops it into an unbounded memory channel, and instantly returns 200 OK. A background IHostedService asynchronously flushes the channel into an MS SQL Server database via Entity Framework Core.
Because MT4 will freeze if the HTTP response is delayed, the C# API acts solely as a high-speed router. It accepts the JSON, drops it into an unbounded memory channel, and instantly returns 200 OK. A background IHostedService asynchronously flushes the channel into an MS SQL Server database via Entity Framework Core.
Code: Select all
using System.Threading.Channels;
using Microsoft.EntityFrameworkCore;
var builder = WebApplication.CreateBuilder(args);
// Configure Entity Framework Core for MS SQL Server
builder.Services.AddDbContext<TelemetryDbContext>(options =>
options.UseSqlServer("Server=localhost;Database=MarketTelemetry;Trusted_Connection=True;TrustServerCertificate=True;"));
// Create high-throughput channel
var tickChannel = Channel.CreateUnbounded<TickRecord>();
builder.Services.AddSingleton(tickChannel.Writer);
builder.Services.AddSingleton(tickChannel.Reader);
// Register Background Worker
builder.Services.AddHostedService<SqlIngestionWorker>();
var app = builder.Build();
// O(1) Minimal API - Unblocks the MT4 thread instantly
app.MapPost("/api/ticks", async (TickRecord[] payload, ChannelWriter<TickRecord> writer) =>
{
foreach (var tick in payload)
{
await writer.WriteAsync(tick);
}
return Results.Ok();
});
app.Run("http://localhost:5000");
// --- Models ---
public record TickRecord(string Symbol, long Time_msc, double Bid, double Ask, double Spread);
public class TelemetryEntity
{
public int Id { get; set; }
public string Symbol { get; set; } = string.Empty;
public DateTime TimestampUtc { get; set; }
public double Bid { get; set; }
public double Ask { get; set; }
public double SpreadPips { get; set; }
}
public class TelemetryDbContext : DbContext
{
public TelemetryDbContext(DbContextOptions<TelemetryDbContext> options) : base(options) { }
public DbSet<TelemetryEntity> Ticks { get; set; }
}
// --- Background Worker ---
public class SqlIngestionWorker : BackgroundService
{
private readonly ChannelReader<TickRecord> _reader;
private readonly IServiceProvider _serviceProvider;
private readonly ILogger<SqlIngestionWorker> _logger;
public SqlIngestionWorker(ChannelReader<TickRecord> reader, IServiceProvider serviceProvider, ILogger<SqlIngestionWorker> logger)
{
_reader = reader;
_serviceProvider = serviceProvider;
_logger = logger;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
_logger.LogInformation("SQL Ingestion Worker Active.");
while (!stoppingToken.IsCancellationRequested)
{
var batch = new List<TelemetryEntity>();
// Read up to 500 records at a time from the channel
while (batch.Count < 500 && _reader.TryRead(out var tick))
{
batch.Add(new TelemetryEntity
{
Symbol = tick.Symbol,
TimestampUtc = DateTimeOffset.FromUnixTimeMilliseconds(tick.Time_msc).UtcDateTime,
Bid = tick.Bid,
Ask = tick.Ask,
SpreadPips = tick.Spread
});
}
if (batch.Any())
{
using var scope = _serviceProvider.CreateScope();
var dbContext = scope.ServiceProvider.GetRequiredService<TelemetryDbContext>();
await dbContext.Ticks.AddRangeAsync(batch, stoppingToken);
await dbContext.SaveChangesAsync(stoppingToken);
_logger.LogInformation("Inserted {Count} ticks into MS SQL Database.", batch.Count);
}
else
{
// Wait for data if channel is empty
await _reader.WaitToReadAsync(stoppingToken);
}
}
}
}Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
Re: How to ask support for historical spread data and what they send
Here is the pure MT5 implementation. MetaTrader 5 offers a much more robust execution environment than MT4, but EAs still run synchronously on a single thread per chart. To prevent the EA from missing ticks during NFP or high-liquidity sweeps, the WebRequest timeout is kept ultra-short, and the batching architecture remains strictly decoupled.
I have updated the headers and optimized the MQL5 syntax specifically for MT5’s memory management.
I have updated the headers and optimized the MQL5 syntax specifically for MT5’s memory management.
Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
Re: How to ask support for historical spread data and what they send
1. The MT5 Sensor (SpreadTelemetry.mq5)
Because this avoids external DLLs or third-party JSON libraries, you can compile this directly in MetaEditor 5 and drop it onto any chart.
Because this avoids external DLLs or third-party JSON libraries, you can compile this directly in MetaEditor 5 and drop it onto any chart.
Code: Select all
//+------------------------------------------------------------------+
//| SpreadTelemetry.mq5 |
//| Copyright 2026, AI Profi Solutions |
//| https://aiprofisolutions.com |
//+------------------------------------------------------------------+
#property strict
//--- Inputs
input string InpApiEndpoint = "http://localhost:5000/api/ticks";
input int InpBatchSize = 50;
input string InpSessionStart = "07:00:00";
input string InpSessionEnd = "10:00:00";
//--- Structs & Globals
struct TickRecord {
ulong time_msc;
double bid;
double ask;
double spread_pips;
};
TickRecord g_tickQueue[];
int g_queueCount = 0;
double g_lastBid = 0.0;
double g_lastAsk = 0.0;
//+------------------------------------------------------------------+
int OnInit()
{
ArrayResize(g_tickQueue, InpBatchSize);
Print("[MT5 Telemetry] Initialized. Target: ", InpApiEndpoint);
return(INIT_SUCCEEDED);
}
//+------------------------------------------------------------------+
void OnTick()
{
MqlTick tick;
if(!SymbolInfoTick(_Symbol, tick)) return;
// Ignore ticks where bid/ask haven't changed (e.g., volume-only updates)
if(tick.bid == g_lastBid && tick.ask == g_lastAsk) return;
datetime gmtNow = TimeTradeServer();
if(!IsWithinSession(gmtNow, InpSessionStart, InpSessionEnd)) return;
double pipSize = (_Digits == 3 || _Digits == 5) ? _Point * 10.0 : _Point;
double spreadPips = (pipSize > 0) ? (tick.ask - tick.bid) / pipSize : 0.0;
// Enqueue
g_tickQueue[g_queueCount].time_msc = tick.time_msc;
g_tickQueue[g_queueCount].bid = tick.bid;
g_tickQueue[g_queueCount].ask = tick.ask;
g_tickQueue[g_queueCount].spread_pips = spreadPips;
g_queueCount++;
g_lastBid = tick.bid;
g_lastAsk = tick.ask;
if(g_queueCount >= InpBatchSize)
{
FlushTelemetry();
g_queueCount = 0;
}
}
//+------------------------------------------------------------------+
void FlushTelemetry()
{
// Manual JSON array construction for zero-dependency compilation
string json = "[";
for(int i = 0; i < g_queueCount; i++)
{
string record = StringFormat("{\"symbol\":\"%s\",\"time_msc\":%I64u,\"bid\":%.5f,\"ask\":%.5f,\"spread\":%.2f}",
_Symbol, g_tickQueue[i].time_msc,
g_tickQueue[i].bid, g_tickQueue[i].ask, g_tickQueue[i].spread_pips);
json += record;
if(i < g_queueCount - 1) json += ",";
}
json += "]";
char postData[];
StringToCharArray(json, postData, 0, WHOLE_ARRAY, CP_UTF8);
// StringToCharArray appends a null terminator, strip it for clean JSON payload
int dataSize = ArraySize(postData) - 1;
char result[];
string resultHeaders;
string headers = "Content-Type: application/json\r\n";
// Timeout set to 500ms. MT5 will not block longer than this.
int res = WebRequest("POST", InpApiEndpoint, headers, 500, postData, dataSize, result, resultHeaders);
if(res != 200)
{
if(res == 4014)
Print("[MT5 Telemetry] ERR 4014: Add ", InpApiEndpoint, " to allowed WebRequest URLs in Tools -> Options -> Expert Advisors.");
else
Print("[MT5 Telemetry] HTTP Error: ", res);
}
}
//+------------------------------------------------------------------+
bool IsWithinSession(datetime time, string startStr, string endStr)
{
datetime sod = time - (time % 86400); // Fast MT5 Start of Day calculation
int currentSec = (int)(time - sod);
string sParts[], eParts[];
StringSplit(startStr, ':', sParts);
StringSplit(endStr, ':', eParts);
int startSec = (int)StringToInteger(sParts[0])*3600 + (int)StringToInteger(sParts[1])*60 + (int)StringToInteger(sParts[2]);
int endSec = (int)StringToInteger(eParts[0])*3600 + (int)StringToInteger(eParts[1])*60 + (int)StringToInteger(eParts[2]);
return (startSec <= endSec) ? (currentSec >= startSec && currentSec <= endSec)
: (currentSec >= startSec || currentSec <= endSec);
}Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.
Re: How to ask support for historical spread data and what they send
2. The ASP.NET Core Ingestion Service
For a multi-core workstation like the HP Z1, this C# backend fully utilizes hardware concurrency. The Minimal API instantly accepts the JSON payload and pushes it into an unbounded memory channel (returning 200 OK to MT5 in less than 2 milliseconds).
A background worker service then pulls from that channel and performs bulk inserts into your local MS SQL Server instance via Entity Framework Core.
For a multi-core workstation like the HP Z1, this C# backend fully utilizes hardware concurrency. The Minimal API instantly accepts the JSON payload and pushes it into an unbounded memory channel (returning 200 OK to MT5 in less than 2 milliseconds).
A background worker service then pulls from that channel and performs bulk inserts into your local MS SQL Server instance via Entity Framework Core.
Code: Select all
using System.Threading.Channels;
using Microsoft.EntityFrameworkCore;
var builder = WebApplication.CreateBuilder(args);
// Configure Entity Framework Core for MS SQL Server
builder.Services.AddDbContext<TelemetryDbContext>(options =>
options.UseSqlServer("Server=localhost;Database=MarketTelemetry;Trusted_Connection=True;TrustServerCertificate=True;"));
// Create high-throughput channel for thread decoupling
var tickChannel = Channel.CreateUnbounded<TickRecord>();
builder.Services.AddSingleton(tickChannel.Writer);
builder.Services.AddSingleton(tickChannel.Reader);
// Register Background Worker
builder.Services.AddHostedService<SqlIngestionWorker>();
var app = builder.Build();
// O(1) Minimal API - Unblocks the MT5 thread instantly
app.MapPost("/api/ticks", async (TickRecord[] payload, ChannelWriter<TickRecord> writer) =>
{
foreach (var tick in payload)
{
await writer.WriteAsync(tick);
}
return Results.Ok();
});
app.Run("http://localhost:5000");
// --- Models ---
public record TickRecord(string Symbol, ulong Time_msc, double Bid, double Ask, double Spread);
public class TelemetryEntity
{
public int Id { get; set; }
public string Symbol { get; set; } = string.Empty;
public DateTime TimestampUtc { get; set; }
public double Bid { get; set; }
public double Ask { get; set; }
public double SpreadPips { get; set; }
}
public class TelemetryDbContext : DbContext
{
public TelemetryDbContext(DbContextOptions<TelemetryDbContext> options) : base(options) { }
public DbSet<TelemetryEntity> Ticks { get; set; }
}
// --- Background Worker ---
public class SqlIngestionWorker : BackgroundService
{
private readonly ChannelReader<TickRecord> _reader;
private readonly IServiceProvider _serviceProvider;
private readonly ILogger<SqlIngestionWorker> _logger;
public SqlIngestionWorker(ChannelReader<TickRecord> reader, IServiceProvider serviceProvider, ILogger<SqlIngestionWorker> logger)
{
_reader = reader;
_serviceProvider = serviceProvider;
_logger = logger;
}
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
_logger.LogInformation("SQL Ingestion Worker Active.");
while (!stoppingToken.IsCancellationRequested)
{
var batch = new List<TelemetryEntity>();
// Read up to 1000 records at a time from the channel
while (batch.Count < 1000 && _reader.TryRead(out var tick))
{
batch.Add(new TelemetryEntity
{
Symbol = tick.Symbol,
TimestampUtc = DateTimeOffset.FromUnixTimeMilliseconds((long)tick.Time_msc).UtcDateTime,
Bid = tick.Bid,
Ask = tick.Ask,
SpreadPips = tick.Spread
});
}
if (batch.Any())
{
using var scope = _serviceProvider.CreateScope();
var dbContext = scope.ServiceProvider.GetRequiredService<TelemetryDbContext>();
// Bulk insert into MS SQL Server
await dbContext.Ticks.AddRangeAsync(batch, stoppingToken);
await dbContext.SaveChangesAsync(stoppingToken);
_logger.LogInformation("Inserted {Count} ticks into MS SQL Database.", batch.Count);
}
else
{
// Yield thread if channel is empty
await _reader.WaitToReadAsync(stoppingToken);
}
}
}
}Preserve your own money. Scale with the market's money. Exponential growth is the ultimate key.