Skip to content

Commit 3a5a2c5

Browse files
Lazy OnEndOfDay ScheduledEvent
- Only add OnEndOfDay ScheduledEvent if the algorithm implements the method. Adding unit tests - Avoid creating a new baseData instance at `SubscriptionDataSourceReader` - Adding static `FineFundamental` instance since creating new ones is expensive
1 parent 32b7a0c commit 3a5a2c5

17 files changed

Lines changed: 195 additions & 46 deletions

AlgorithmFactory/Python/Wrappers/AlgorithmPythonWrapper.cs

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -47,6 +47,11 @@ public class AlgorithmPythonWrapper : IAlgorithm
4747
private readonly dynamic _onMarginCall;
4848
private readonly IAlgorithm _baseAlgorithm;
4949

50+
/// <summary>
51+
/// True if the underlying python algorithm implements "OnEndOfDay"
52+
/// </summary>
53+
public bool IsOnEndOfDayImplemented { get; }
54+
5055
/// <summary>
5156
/// <see cref = "AlgorithmPythonWrapper"/> constructor.
5257
/// Creates and wraps the algorithm written in python.
@@ -90,6 +95,8 @@ public AlgorithmPythonWrapper(string moduleName)
9095
_onMarginCall = pyAlgorithm.GetPythonMethod("OnMarginCall");
9196

9297
_onOrderEvent = pyAlgorithm.GetAttr("OnOrderEvent");
98+
99+
IsOnEndOfDayImplemented = pyAlgorithm.GetPythonMethod("OnEndOfDay") != null;
93100
}
94101
attr.Dispose();
95102
}

Engine/DataFeeds/Enumerators/Factories/BaseDataCollectionSubscriptionEnumeratorFactory.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -63,7 +63,7 @@ public IEnumerator<BaseData> CreateEnumerator(SubscriptionRequest request, IData
6363
foreach (var date in tradableDays)
6464
{
6565
var source = sourceFactory.GetSource(configuration, date, false);
66-
var factory = SubscriptionDataSourceReader.ForSource(source, dataCacheProvider, configuration, date, false);
66+
var factory = SubscriptionDataSourceReader.ForSource(source, dataCacheProvider, configuration, date, false, sourceFactory);
6767
var coarseFundamentalForDate = factory.Read(source);
6868
// shift all date of emitting the file forward one day to model emitting coarse midnight the next day.
6969
yield return new BaseDataCollection(date.AddDays(1), configuration.Symbol, coarseFundamentalForDate);

Engine/DataFeeds/Enumerators/Factories/BaseDataSubscriptionEnumeratorFactory.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -79,7 +79,7 @@ public IEnumerator<BaseData> CreateEnumerator(SubscriptionRequest request, IData
7979
request.Configuration.MappedSymbol = GetMappedSymbol(request.Configuration, date);
8080
}
8181
var source = sourceFactory.GetSource(request.Configuration, date, _isLiveMode);
82-
var factory = SubscriptionDataSourceReader.ForSource(source, dataCacheProvider, request.Configuration, date, _isLiveMode);
82+
var factory = SubscriptionDataSourceReader.ForSource(source, dataCacheProvider, request.Configuration, date, _isLiveMode, sourceFactory);
8383
var entriesForDate = factory.Read(source);
8484
foreach (var entry in entriesForDate)
8585
{

Engine/DataFeeds/Enumerators/Factories/FineFundamentalSubscriptionEnumeratorFactory.cs

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,8 @@ public class FineFundamentalSubscriptionEnumeratorFactory : ISubscriptionEnumera
3636
{
3737
private static readonly ConcurrentDictionary<int, List<DateTime>> FineFilesCache
3838
= new ConcurrentDictionary<int, List<DateTime>>();
39+
// creating a fine fundamental instance is expensive (its massive) so we keep our factory instance
40+
private static readonly FineFundamental FineFundamental = new FineFundamental();
3941

4042
private readonly bool _isLiveMode;
4143
private readonly Func<SubscriptionRequest, IEnumerable<DateTime>> _tradableDaysProvider;
@@ -64,13 +66,12 @@ public IEnumerator<BaseData> CreateEnumerator(SubscriptionRequest request, IData
6466
{
6567
var tradableDays = _tradableDaysProvider(request);
6668

67-
var fineFundamental = new FineFundamental();
6869
var fineFundamentalConfiguration = new SubscriptionDataConfig(request.Configuration, typeof(FineFundamental), request.Security.Symbol);
6970

7071
foreach (var date in tradableDays)
7172
{
72-
var fineFundamentalSource = GetSource(fineFundamental, fineFundamentalConfiguration, date);
73-
var fineFundamentalFactory = SubscriptionDataSourceReader.ForSource(fineFundamentalSource, dataCacheProvider, fineFundamentalConfiguration, date, _isLiveMode);
73+
var fineFundamentalSource = GetSource(FineFundamental, fineFundamentalConfiguration, date);
74+
var fineFundamentalFactory = SubscriptionDataSourceReader.ForSource(fineFundamentalSource, dataCacheProvider, fineFundamentalConfiguration, date, _isLiveMode, FineFundamental);
7475
var fineFundamentalForDate = (FineFundamental)fineFundamentalFactory.Read(fineFundamentalSource).FirstOrDefault();
7576

7677
yield return new FineFundamental

Engine/DataFeeds/Enumerators/Factories/LiveCustomDataSubscriptionEnumeratorFactory.cs

Lines changed: 6 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,7 @@ public IEnumerator<BaseData> CreateEnumerator(SubscriptionRequest request, IData
7272
var source = sourceFactory.GetSource(config, localDate, true);
7373

7474
// fetch the new source and enumerate the data source reader
75-
var enumerator = EnumerateDataSourceReader(config, dataProvider, frontier, source, localDate);
75+
var enumerator = EnumerateDataSourceReader(config, dataProvider, frontier, source, localDate, sourceFactory);
7676

7777
if (SourceRequiresFastForward(source))
7878
{
@@ -127,12 +127,12 @@ public IEnumerator<BaseData> CreateEnumerator(SubscriptionRequest request, IData
127127
return refresher;
128128
}
129129

130-
private IEnumerator<BaseData> EnumerateDataSourceReader(SubscriptionDataConfig config, IDataProvider dataProvider, Ref<DateTime> localFrontier, SubscriptionDataSource source, DateTime localDate)
130+
private IEnumerator<BaseData> EnumerateDataSourceReader(SubscriptionDataConfig config, IDataProvider dataProvider, Ref<DateTime> localFrontier, SubscriptionDataSource source, DateTime localDate, BaseData baseDataInstance)
131131
{
132132
using (var dataCacheProvider = new SingleEntryDataCacheProvider(dataProvider))
133133
{
134134
var newLocalFrontier = localFrontier.Value;
135-
var dataSourceReader = GetSubscriptionDataSourceReader(source, dataCacheProvider, config, localDate);
135+
var dataSourceReader = GetSubscriptionDataSourceReader(source, dataCacheProvider, config, localDate, baseDataInstance);
136136
foreach (var datum in dataSourceReader.Read(source))
137137
{
138138
// always skip past all times emitted on the previous invocation of this enumerator
@@ -177,10 +177,11 @@ private IEnumerator<BaseData> EnumerateDataSourceReader(SubscriptionDataConfig c
177177
protected virtual ISubscriptionDataSourceReader GetSubscriptionDataSourceReader(SubscriptionDataSource source,
178178
IDataCacheProvider dataCacheProvider,
179179
SubscriptionDataConfig config,
180-
DateTime date
180+
DateTime date,
181+
BaseData baseDataInstance
181182
)
182183
{
183-
return SubscriptionDataSourceReader.ForSource(source, dataCacheProvider, config, date, true);
184+
return SubscriptionDataSourceReader.ForSource(source, dataCacheProvider, config, date, true, baseDataInstance);
184185
}
185186

186187
private bool SourceRequiresFastForward(SubscriptionDataSource source)

Engine/DataFeeds/IndexSubscriptionDataSourceReader.cs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -104,7 +104,8 @@ public override IEnumerable<BaseData> Read(SubscriptionDataSource source)
104104
DataCacheProvider,
105105
_config,
106106
_date,
107-
IsLiveMode);
107+
IsLiveMode,
108+
_factory);
108109

109110
var enumerator = dataReader.Read(dataSource).GetEnumerator();
110111
while (enumerator.MoveNext())

Engine/DataFeeds/SubscriptionDataReader.cs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -469,7 +469,7 @@ private bool UpdateDataEnumerator(bool endOfEnumerator)
469469

470470
// save off for comparison next time
471471
_source = newSource;
472-
var subscriptionFactory = CreateSubscriptionFactory(newSource);
472+
var subscriptionFactory = CreateSubscriptionFactory(newSource, _dataFactory);
473473
_subscriptionFactoryEnumerator = subscriptionFactory.Read(newSource).GetEnumerator();
474474
return true;
475475
}
@@ -488,9 +488,9 @@ private bool UpdateDataEnumerator(bool endOfEnumerator)
488488
while (true);
489489
}
490490

491-
private ISubscriptionDataSourceReader CreateSubscriptionFactory(SubscriptionDataSource source)
491+
private ISubscriptionDataSourceReader CreateSubscriptionFactory(SubscriptionDataSource source, BaseData baseDataInstance)
492492
{
493-
var factory = SubscriptionDataSourceReader.ForSource(source, _dataCacheProvider, _config, _tradeableDates.Current, _isLiveMode);
493+
var factory = SubscriptionDataSourceReader.ForSource(source, _dataCacheProvider, _config, _tradeableDates.Current, _isLiveMode, baseDataInstance);
494494
AttachEventHandlers(factory, source);
495495
return factory;
496496
}

Engine/DataFeeds/SubscriptionDataSourceReader.cs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -35,8 +35,9 @@ public static class SubscriptionDataSourceReader
3535
/// <param name="config">The configuration of the subscription</param>
3636
/// <param name="date">The date to be processed</param>
3737
/// <param name="isLiveMode">True for live mode, false otherwise</param>
38+
/// <param name="factory">The base data instance factory</param>
3839
/// <returns>A new <see cref="ISubscriptionDataSourceReader"/> that can read the specified <paramref name="source"/></returns>
39-
public static ISubscriptionDataSourceReader ForSource(SubscriptionDataSource source, IDataCacheProvider dataCacheProvider, SubscriptionDataConfig config, DateTime date, bool isLiveMode)
40+
public static ISubscriptionDataSourceReader ForSource(SubscriptionDataSource source, IDataCacheProvider dataCacheProvider, SubscriptionDataConfig config, DateTime date, bool isLiveMode, BaseData factory)
4041
{
4142
ISubscriptionDataSourceReader reader;
4243
TextSubscriptionDataSourceReader textReader = null;
@@ -64,7 +65,6 @@ public static ISubscriptionDataSourceReader ForSource(SubscriptionDataSource sou
6465
// wire up event handlers for logging missing files
6566
if (source.TransportMedium == SubscriptionTransportMedium.LocalFile)
6667
{
67-
var factory = config.GetBaseDataInstance();
6868
if (!factory.IsSparseData())
6969
{
7070
reader.InvalidSource += (sender, args) => Log.Error($"SubscriptionDataSourceReader.InvalidSource(): File not found: {args.Source.Source}");

Engine/DataFeeds/TextSubscriptionDataSourceReader.cs

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -34,9 +34,9 @@ namespace QuantConnect.Lean.Engine.DataFeeds
3434
public class TextSubscriptionDataSourceReader : BaseSubscriptionDataSourceReader
3535
{
3636
private readonly bool _implementsStreamReader;
37-
private readonly BaseData _factory;
3837
private readonly DateTime _date;
3938
private readonly SubscriptionDataConfig _config;
39+
private BaseData _factory;
4040
private bool _shouldCacheDataPoints;
4141
private static readonly MemoryCache BaseDataSourceCache = new MemoryCache("BaseDataSourceCache",
4242
// Cache can use up to 70% of the installed physical memory
@@ -77,7 +77,6 @@ public TextSubscriptionDataSourceReader(IDataCacheProvider dataCacheProvider, Su
7777
{
7878
_date = date;
7979
_config = config;
80-
_factory = config.GetBaseDataInstance();
8180
_shouldCacheDataPoints = !_config.IsCustomData && _config.Resolution >= Resolution.Hour
8281
&& _config.Type != typeof(FineFundamental) && _config.Type != typeof(CoarseFundamental)
8382
&& !DataCacheProvider.IsDataEphemeral;
@@ -114,6 +113,12 @@ public override IEnumerable<BaseData> Read(SubscriptionDataSource source)
114113
OnCreateStreamReaderError(_date, source);
115114
yield break;
116115
}
116+
117+
if (_factory == null)
118+
{
119+
// only create a factory if the stream isn't null
120+
_factory = _config.GetBaseDataInstance();
121+
}
117122
// while the reader has data
118123
while (!reader.EndOfStream)
119124
{

Engine/RealTime/BacktestingRealTimeHandler.cs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -52,7 +52,7 @@ public void Setup(IAlgorithm algorithm, AlgorithmNodePacket job, IResultHandler
5252

5353
// create events for algorithm's end of tradeable dates
5454
// set up the events for each security to fire every tradeable date before market close
55-
base.Setup(Algorithm.StartDate, Algorithm.EndDate);
55+
base.Setup(Algorithm.StartDate, Algorithm.EndDate, job.Language);
5656

5757
foreach (var scheduledEvent in GetScheduledEventsSortedByTime())
5858
{

0 commit comments

Comments
 (0)