Concurrent DataStream-Chunks

This commit is contained in:
Michael Schimmel
2025-07-16 16:03:10 +02:00
parent 342eb07c42
commit 120c62083e
2 changed files with 30 additions and 13 deletions
+13 -5
View File
@@ -145,6 +145,14 @@ end;
procedure TForm1.StopButtonClick(Sender: TObject); procedure TForm1.StopButtonClick(Sender: TObject);
begin begin
FTerminate.Notify; FTerminate.Notify;
TaskManager.WaitFor(FProcessDone);
var Layout := CurrLayout<TVertScrollBox>;
if Layout <> nil then
begin
Layout.Content.DeleteChildren;
Layout.Repaint;
end;
end; end;
procedure TForm1.TreeViewDblClick(Sender: TObject); procedure TForm1.TreeViewDblClick(Sender: TObject);
@@ -362,12 +370,12 @@ begin
var Closes := Ohlc.Field<Double>('Close'); var Closes := Ohlc.Field<Double>('Close');
var Hull := Closes.MakeParallel.Chain<Double>(TIndicators.CreateHMA(150)); var Hull := Closes.MakeParallel.Chain<Double>(TIndicators.CreateHMA(150));
var Sma := Closes.Chain<Double>(TIndicators.CreateSMA(50)); var Sma := Closes.MakeParallel.Chain<Double>(TIndicators.CreateSMA(50));
var Ema := Closes.Chain<Double>(TIndicators.CreateEMA(21)); var Ema := Closes.MakeParallel.Chain<Double>(TIndicators.CreateEMA(21));
var Boli := Closes.MakeParallel.Chain<TBollingerBandsResult>(TIndicators.CreateBollingerBands(20, 2.0)); var Boli := Closes.MakeParallel.Chain<TBollingerBandsResult>(TIndicators.CreateBollingerBands(20, 2.0));
var Rsi := Closes.MakeParallel.Chain<Double>(TIndicators.CreateRSI(14)); var Rsi := Closes.MakeParallel.Chain<Double>(TIndicators.CreateRSI(14));
var Macd := Closes.MakeParallel.Chain<TMacdResult>(TIndicators.CreateMACD(12, 26, 9)); var Macd := Closes.MakeParallel.Chain<TMacdResult>(TIndicators.CreateMACD(12, 26, 9));
var Stoch := Ohlc.Chain<TStochasticResult>(TIndicators.CreateStochastic(14, 3)); var Stoch := Ohlc.MakeParallel.Chain<TStochasticResult>(TIndicators.CreateStochastic(14, 3));
chart.SetXAxisSeries(timeframe, Timestamps.Sender); chart.SetXAxisSeries(timeframe, Timestamps.Sender);
@@ -430,7 +438,7 @@ type
pnl: Double; pnl: Double;
end; end;
begin begin
var timeframe := TTimeframe.D; var timeframe := TTimeframe.M15;
var ticker := TConverter.CreateTicker<TDataPoint<TAskBidItem>>; var ticker := TConverter.CreateTicker<TDataPoint<TAskBidItem>>;
@@ -453,7 +461,7 @@ begin
var Lowest: Double := Double.MaxValue; var Lowest: Double := Double.MaxValue;
var Highest: Double := Double.MinValue; var Highest: Double := Double.MinValue;
var ATR := Ohlc[0].Chain<Double>(TIndicators.CreateATR(15)); var ATR := Ohlc[0].Chain<Double>(TIndicators.CreateATR(50)).MakeParallel;
var ATRSeries := TConverter.CreateEndpoint<Double>(ATR.Sender, 5); var ATRSeries := TConverter.CreateEndpoint<Double>(ATR.Sender, 5);
var HullSeries := TConverter.CreateEndpoint<Double>(Hull.Sender, 5); var HullSeries := TConverter.CreateEndpoint<Double>(Hull.Sender, 5);
+17 -8
View File
@@ -62,7 +62,7 @@ type
function ProcessChunks( function ProcessChunks(
const DataChunks: TArray<TArray<TDataPoint<T>>>; const DataChunks: TArray<TArray<TDataPoint<T>>>;
const Terminated: TState; Terminated: TState;
Processor: IMycProcessor<TArray<TDataPoint<T>>> Processor: IMycProcessor<TArray<TDataPoint<T>>>
): TState; ): TState;
@@ -362,19 +362,28 @@ end;
function TAuraDataServer<T>.ProcessChunks( function TAuraDataServer<T>.ProcessChunks(
const DataChunks: TArray<TArray<TDataPoint<T>>>; const DataChunks: TArray<TArray<TDataPoint<T>>>;
const Terminated: TState; Terminated: TState;
Processor: IMycProcessor<TArray<TDataPoint<T>>> Processor: IMycProcessor<TArray<TDataPoint<T>>>
): TState; ): TState;
begin begin
var done := TLatch.CreateLatch(Length(DataChunks)); var callProcess :=
function(Idx: Integer): TFunc<TState>
begin
var cChunk := DataChunks[Idx];
Result :=
function: TState
begin
if not Terminated.IsSet then
Result := Processor.ProcessData(cChunk);
end;
end;
for var i := 0 to High(DataChunks) do for var i := 0 to High(DataChunks) do
if not Terminated.IsSet then if not Terminated.IsSet then
Processor.ProcessData(DataChunks[i]).Signal.Subscribe(done) begin
else var q := TaskManager.RunTask(Result, callProcess(i));
done.Notify; Result := q;
end
Result := done.State;
end; end;
function TAuraDataServer<T>.ProcessFile( function TAuraDataServer<T>.ProcessFile(