DataPoint
This commit is contained in:
@@ -12,37 +12,37 @@ type
|
||||
TDataStreamProvider<T> = class(TInterfacedObject, TMutable<TDataSeries<T>>.IMutable)
|
||||
private
|
||||
FStream: IDataStream<T>;
|
||||
FDataSeries: TDataSeries<T>;
|
||||
FDataArray: IDataArray<T>;
|
||||
FNewDataAvailable: TFlag;
|
||||
FHasDataSubscription: TSignal.TSubscription;
|
||||
FChunk: TArray<TDataPoint<T>>;
|
||||
FLookback: Int64;
|
||||
function GetChanged: TSignal;
|
||||
function GetValue: TDataSeries<T>;
|
||||
public
|
||||
constructor Create(ALookback: Int64; const AStream: IDataStream<T>);
|
||||
constructor Create(ALookback, AMaxChunkSize: Int64; const AStream: IDataStream<T>);
|
||||
destructor Destroy; override;
|
||||
end;
|
||||
|
||||
TDataStreamProvider = record
|
||||
class function Create<T>(Lookback: Int64; const Stream: IDataStream<T>): TMutable<TDataSeries<T>>; static;
|
||||
class function Create<T>(Lookback, MaxChunkSize: Int64; const Stream: IDataStream<T>): TMutable<TDataSeries<T>>; static;
|
||||
end;
|
||||
|
||||
implementation
|
||||
|
||||
{ TDataStreamProvider<T> }
|
||||
|
||||
constructor TDataStreamProvider<T>.Create(ALookback: Int64; const AStream: IDataStream<T>);
|
||||
constructor TDataStreamProvider<T>.Create(ALookback, AMaxChunkSize: Int64; const AStream: IDataStream<T>);
|
||||
begin
|
||||
inherited Create;
|
||||
FLookback := ALookback;
|
||||
SetLength(FChunk, AMaxChunkSize);
|
||||
FStream := AStream;
|
||||
FHasDataSubscription := FStream.HasData.Subscribe(FNewDataAvailable);
|
||||
FDataSeries := TDataSeries<T>.Create( ALookback );
|
||||
FNewDataAvailable := TFlag.CreateObserver(FStream.HasData);
|
||||
FDataArray := TDataSeries<T>.CreateArray(ALookback);
|
||||
end;
|
||||
|
||||
destructor TDataStreamProvider<T>.Destroy;
|
||||
begin
|
||||
FHasDataSubscription.Unsubscribe;
|
||||
inherited;
|
||||
end;
|
||||
|
||||
@@ -55,22 +55,18 @@ function TDataStreamProvider<T>.GetValue: TDataSeries<T>;
|
||||
begin
|
||||
if FNewDataAvailable.Reset then
|
||||
begin
|
||||
var Data: array[0..1024 - 1] of TDataPoint<T>;
|
||||
var cnt := FStream.GetChunk(Data);
|
||||
while cnt > 0 do
|
||||
begin
|
||||
FDataSeries.Add(Data, cnt);
|
||||
cnt := FStream.GetChunk(Data);
|
||||
end;
|
||||
var cnt := FStream.GetChunk(FChunk);
|
||||
if cnt > 0 then
|
||||
FDataArray.Add(FChunk, cnt);
|
||||
end;
|
||||
Result := FDataSeries;
|
||||
Result := FDataArray;
|
||||
end;
|
||||
|
||||
{ TDataStreamProvider }
|
||||
|
||||
class function TDataStreamProvider.Create<T>(Lookback: Int64; const Stream: IDataStream<T>): TMutable<TDataSeries<T>>;
|
||||
class function TDataStreamProvider.Create<T>(Lookback, MaxChunkSize: Int64; const Stream: IDataStream<T>): TMutable<TDataSeries<T>>;
|
||||
begin
|
||||
Result := TDataStreamProvider<T>.Create(Lookback, Stream);
|
||||
Result := TDataStreamProvider<T>.Create(Lookback, MaxChunkSize, Stream);
|
||||
end;
|
||||
|
||||
end.
|
||||
|
||||
Reference in New Issue
Block a user