eb1181f5f7
* Adds preliminary universe selection for Future Options
* Fixes scaling issues with Future Options
* Fixes scaling multiplying by 10000x instead of using _scaleFactor
* Fixes scaling for Tick
* Revert changes to Tick since it divides the scaling factor
* Changes stale method name to new method name after rebase
* Fixes selection bugs, adds new methods, and adds unit tests
* Fixes bug where Equity Symbol was created for an underlying
non-equity Symbol, resulting in equity data trying to be loaded
* Adds unit tests covering changes to Tick, QuoteBar, TradeBar and
LeanData
* Adds regression test for AddUniverseOption filter contract selection
for Future Options
* Addresses review - modifies the AddFutureOption signature
* Adds new AddUniverseOptions method overload
* Removes and adds a new unit test
* Misc. modifications to account for new changes
* Fixes bug where futures were loaded using default SID Date
* Refactors and removes unnecessary work
* Fixes regression algorithm, which previously made no trades
* Adds future option data
* Adds the corresponding underlying data, in this case, futures data
to enable usage of future options data
* Replaces data with new data (ES18Z20)
* Improves Future chain filtering and updates regression stats
* Add AddFutureOptionContract API
* Expands regression and unit tests to test in finer detail
* Adds Python regression algorithms for AddFutureOption[Contract] methods
* Adds new unit test for BacktestingOptionChainProvider
* Fixes bug with BacktesingOptionChainProvider where we
attempted to load the Trades option chain first, resulting
in breakage of backwards compatibility and limitation of the
option chain.
* Adds new regression algorithms (Py) to Algorithm.Python project
* Adds FutureOptionMarginBuyingPowerModel
* Modifies code paths used to select margin model
* Adds related unit tests for margin model
* Fixes issue with unit test and MHDB/SPDB lookup for Future Options
* Preliminary regression algorithm testing ITM call/put option buying
* Fixes bug where fee model used did not find non-US market
options fee model. We now use the futures fee model for future
options because IB charges the same commissions per contract
between futures and futures options
* Adds proper regression algorithm for ITM future options expiration
* Pushing broken algorithm for review
* Currently, algorithm does not fill forward, causing
a single future option to not get exercised when it is delisted.
* Adds FutureOptionPutITMExpiryRegressionAlgorithm
* Improves existing regression algorithm for call side
* Fixes bug in existing regression algorithm
* Adds AAPL daily data to advance enumerator for ^^^ fix
* Adds additional future option regression algorithms
* Adds Buy OTM expiration regression algorithms
* Adds Sell ITM/OTM expiration regression algorithms
* Adds missing Python regression algorithms
* Adds remaining Python regression algorithms and fixes issues
* Fixes naming issues and statistics
* Adds short option OTM regression algorithms (Py)
* Add license header and class comments to python algorithms
* Cleans up comments and docstrings
* Create Buy/Sell call intraday regression algo
* Redirects future options symbol properties to futures symbol properties
* Asserts exercise/assignment price and updates stats in regression algos
* Adds new unit test covering changes to SecurityService
* Adds comments and fixes failing test
* Partially fixes future option mis-calculated profit/loss
* Adjusts portfolio model to calculate FOP as a no upfront pay asset class
* Updates regression algorithm statistics
* Begin IB FOP support
* Initial support for FOP IB data streaming, live í¾
* Adds additional functionality to LiveOptionChainProvider
- Allows querying CME API to retrieve option chains for CME products
- Ultimately, it's also the groundwork for the CME
LiveFutureChainProvider
* Edits IDataQueueUniverseProvider interface to provide greater
control to implementors of it
* Misc. bug fixes required to get FOP data streaming through IB
* Adds comments, adds missing rategate call, and cleans up code
* Force exchange for FOP and Futures when no exchange is provided
* Fixes bug with Portfolio modeling across all asset classes
* Adds LiveOptionChainProvider tests for Future Options
* IB brokerage option symbol bug fixes and improvements
* Fixes contract multiplier lookup bug
* Fixes issue where we attempted to subscribe to IB data feed with canonical security
* Adds ES MHDB entry
* Reverts portfolio modeling changes for Futures Options
* Since IB eats into our account's cash balance when
a new FOP contract is purchased, we must model by applying funds
to our cash whenever a new purchase/sell occurs.
If we choose to model FOPs exactly as we do with futures, we
will end up with an invalid TotalPortfolioValue on algorithm
restart. By all means and purposes, FOPs are modeled exactly
the same as equity options with respect to the portfolio.
* Adds comments clarifying portfolio modeling and clarifies
existing portfolio modeling comments with additional context.
* Fixes IB symbol lookup for future options
* Fixes LiveOptionChainProvider looping 5 times per option chain
request, even on success
* Sets OptionChainedUniverseSelectionModel to produce a canonical
future/future option/option Symbol to avoid creating two Symbols
* Adds GLOBEX future option symbol mapping from future -> fop
* Fixes LiveOptionChainProvider loading wrong contract option chains
* Fixes loading of futures options ZIP files when backtesting
* Adds a string -> decimal JSON converter
* Additional fixes/refactoring to the LiveOptionChainProvider
* Adds tests for changes to Symbol and LeanData
* Reverts changes to IB-symbol-map
* Fixes Value for mapped future options tickers
* Fixes Symbol test
* Changes path of future options to future's expiry date
* Extra changes made to remove scaling from writing CSV
* Added method to map from FOP Globex -> FUT Globex
* Fixes MOO and MOC orders for future options
* Note: this order type might not be supported by IB or CME.
* Bug fixes and updates unit tests
* Update regression tests and data format
* Rebase changes
* 1. Multiple bug fixes for LiveOptionChainProvider, reverts IQFeed changes
2. Address review (partial): Code reuse and cleanup
1.
* Modifies check in
`AddFutureOptionShort(Call|Put)ITMExpiryRegressionAlgorithm`
to ensure no buys have negative quantity
* Code reuse changes in IB brokerage
* Bug fix in IB brokerage where we assigned the FOP expiry
as the futures expiry (requires verification)
* Doc changes and adds missing summaries/license banners
* Disposes of HTTP client resources in LiveOptionChainProvider
* Renames classes and adds FutureOption folder in Common/Securities
2.
* We revert back to the quotes API for the option chain,
since the settlement API sometimes had missing strikes.
* Fixes future option expiry being set as future's expiry
in LiveOptionChainProvider
* Fixes bug where wrong option chain was selected because of bad
expiry lookup in the futures expiries returned from CME
* Fixes multiple looping bug in LiveOptionChainProvider
* Adds strike price scaling for LiveOptionChainProvider
* Reverts IQFeed changes and simplifies interface upgrade changes
Some additional challenges we'll have to solve as part of FOPs:
- The `OptionSymbol.IsStandard` method makes the assumption that
weeklies contracts follow the pattern equities follows, which
does not apply to Futures Options
- The Subscription created in:
`OptionChainUniverseSubscriptionEnumeratorFactory`
...adds a Trade config. For illiquid contracts, this
will delay universe selection for the option symbol
until we get a trade. However, if we add a quote config,
the data would instead be loaded based on the first quote
we received from the brokerage.
But since we're currently using a trade config, illiquid
contracts won't start streaming data until it receives a trade.
NOTE: this commit is a WIP to addressing the reviews received in the PR,
but has been committed early for efficiency in the review process
* Fixes regression algorithms and misc. bugs
* Fixes map file lookup for non-equity options
* Adds extra assertion at end of algorithm to ensure no holdings are
left when the algorithm ends.
* Adds FutureOptionSymbol, allowing all contracts through as standard
* Changes SPDB to allow defaulting to underlying future symbol
properties if no entry is found for the given FOP
* Fixes calls to SPDB in SecurityService, IBBrokerage
* Reverts AAPL daily ZIP file to fix majority of regression algorithms
* Adds FOPs symbol properties
* Fixes existing symbol properties for a few futures
* Adds tests for changes to Symbol Properties Database
* Removes string SPDB lookup method
* Updates tests and misc callees of previous method
* Updates all regression tests to use data of already expired contracts
* Adds Futures Options Expiry Functions tests
* Adds required futures data for 2020-01-05
* Address review (partial): Expands test coverage and fixes tests
* Set option chain tests parallelism to fixture only
* Fixes broken test for contract month delta for FuturesOptionsExpiryFunctions
* Changes delisting date logic for Futures Options
* Address review: removes duplicate code, misc code fixes
* Bug fix in MarketHoursDatabase.GetDatabaseSymbolKey() where
we would use the underlying's Symbol for lookup in the MHDB
* Adds missing license banner
* Removes Futures Options entries from MHDB
* Adds new tests
* Adds SecurityType.FutureOption
* Converts any underlying comparisons and uses SecurityType directly
instead for FOP specific behavior
* Extra code modifications to acommodate new SecurityType
* Addresses review: fixes order fee bug on exercise
* Additional bug fixes and adding of SecurityType.FutureOption
* Updates regression algorithms OrderListHash
* Fixes various bugs in IB live implementation
* Fixes bug setting the right contract expiration date for FOP
generated by LiveOptionChainProvider
* Adds new function to FuturesOptionsExpiryFunctions
* Clarifies parameter names better in some functions/methods
* Fixes bugs in IB brokerage for FOPs
* Address review - code cleanup and refactor
* Remove MappingEventProvider, SplitEventProvider, and
DividendEventProvider for Futures Options in
CorporateEventEnumeratorFactory
* Address review: Use MHDB key resolver in SPDB
* Makes regression tests pass and adds comment for expiry issue
* Fixes MHDB lookup on string symbol method
* Adds Futures Options greeks regression algorithm (C# only)
* Adds explanitory comment on MHDB FOP lookup
* Remove python from FutureOptionCallITMGreeksExpiryRegressionAlgorithm
514 lines
20 KiB
C#
514 lines
20 KiB
C#
/*
|
|
* QUANTCONNECT.COM - Democratizing Finance, Empowering Individuals.
|
|
* Lean Algorithmic Trading Engine v2.0. Copyright 2014 QuantConnect Corporation.
|
|
*
|
|
* Licensed under the Apache License, Version 2.0 (the "License");
|
|
* you may not use this file except in compliance with the License.
|
|
* You may obtain a copy of the License at http://www.apache.org/licenses/LICENSE-2.0
|
|
*
|
|
* Unless required by applicable law or agreed to in writing, software
|
|
* distributed under the License is distributed on an "AS IS" BASIS,
|
|
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
* See the License for the specific language governing permissions and
|
|
* limitations under the License.
|
|
*/
|
|
|
|
using System;
|
|
using System.Collections.Generic;
|
|
using System.IO;
|
|
using System.Linq;
|
|
using System.Text;
|
|
using Ionic.Zip;
|
|
using QuantConnect.Data;
|
|
using QuantConnect.Interfaces;
|
|
using QuantConnect.Logging;
|
|
using QuantConnect.Securities;
|
|
using QuantConnect.Util;
|
|
|
|
namespace QuantConnect.ToolBox
|
|
{
|
|
/// <summary>
|
|
/// Data writer for saving an IEnumerable of BaseData into the LEAN data directory.
|
|
/// </summary>
|
|
public class LeanDataWriter
|
|
{
|
|
private readonly Symbol _symbol;
|
|
private readonly string _dataDirectory;
|
|
private readonly TickType _tickType;
|
|
private readonly bool _appendToZips;
|
|
private readonly Resolution _resolution;
|
|
private readonly SecurityType _securityType;
|
|
|
|
/// <summary>
|
|
/// Create a new lean data writer to this base data directory.
|
|
/// </summary>
|
|
/// <param name="symbol">Symbol string</param>
|
|
/// <param name="dataDirectory">Base data directory</param>
|
|
/// <param name="resolution">Resolution of the desired output data</param>
|
|
/// <param name="tickType">The tick type</param>
|
|
public LeanDataWriter(Resolution resolution, Symbol symbol, string dataDirectory, TickType tickType = TickType.Trade) : this(
|
|
dataDirectory,
|
|
resolution,
|
|
symbol.ID.SecurityType,
|
|
tickType
|
|
)
|
|
{
|
|
_symbol = symbol;
|
|
// All fx data is quote data.
|
|
if (_securityType == SecurityType.Forex || _securityType == SecurityType.Cfd)
|
|
{
|
|
_tickType = TickType.Quote;
|
|
}
|
|
|
|
if (_securityType != SecurityType.Equity && _securityType != SecurityType.Forex && _securityType != SecurityType.Cfd && _securityType != SecurityType.Crypto && _securityType != SecurityType.Future && _securityType != SecurityType.Option && _securityType != SecurityType.FutureOption)
|
|
{
|
|
throw new Exception("Sorry this security type is not yet supported by the LEAN data writer: " + _securityType);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Create a new lean data writer to this base data directory.
|
|
/// </summary>
|
|
/// <param name="dataDirectory">Base data directory</param>
|
|
/// <param name="resolution">Resolution of the desired output data</param>
|
|
/// <param name="securityType">The security type</param>
|
|
/// <param name="tickType">The tick type</param>
|
|
public LeanDataWriter(string dataDirectory, Resolution resolution, SecurityType securityType, TickType tickType)
|
|
{
|
|
_dataDirectory = dataDirectory;
|
|
_resolution = resolution;
|
|
_securityType = securityType;
|
|
_tickType = tickType;
|
|
_appendToZips = securityType == SecurityType.Future;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Given the constructor parameters, write out the data in LEAN format.
|
|
/// </summary>
|
|
/// <param name="source">IEnumerable source of the data: sorted from oldest to newest.</param>
|
|
public void Write(IEnumerable<BaseData> source)
|
|
{
|
|
switch (_resolution)
|
|
{
|
|
case Resolution.Daily:
|
|
case Resolution.Hour:
|
|
WriteDailyOrHour(source);
|
|
break;
|
|
|
|
case Resolution.Minute:
|
|
case Resolution.Second:
|
|
case Resolution.Tick:
|
|
WriteMinuteOrSecondOrTick(source);
|
|
break;
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Downloads historical data from the brokerage and saves it in LEAN format.
|
|
/// </summary>
|
|
/// <param name="brokerage">The brokerage from where to fetch the data</param>
|
|
/// <param name="symbols">The list of symbols</param>
|
|
/// <param name="startTimeUtc">The starting date/time (UTC)</param>
|
|
/// <param name="endTimeUtc">The ending date/time (UTC)</param>
|
|
public void DownloadAndSave(IBrokerage brokerage, List<Symbol> symbols, DateTime startTimeUtc, DateTime endTimeUtc)
|
|
{
|
|
if (symbols.Count == 0)
|
|
{
|
|
throw new ArgumentException("DownloadAndSave(): The symbol list cannot be empty.");
|
|
}
|
|
|
|
if (_tickType != TickType.Trade && _tickType != TickType.Quote)
|
|
{
|
|
throw new ArgumentException("DownloadAndSave(): The tick type must be Trade or Quote.");
|
|
}
|
|
|
|
if (_securityType != SecurityType.Future && _securityType != SecurityType.Option && _securityType != SecurityType.FutureOption)
|
|
{
|
|
throw new ArgumentException($"DownloadAndSave(): The security type must be {SecurityType.Future} or {SecurityType.Option}.");
|
|
}
|
|
|
|
if (symbols.Any(x => x.SecurityType != _securityType))
|
|
{
|
|
throw new ArgumentException($"DownloadAndSave(): All symbols must have {_securityType} security type.");
|
|
}
|
|
|
|
if (symbols.DistinctBy(x => x.ID.Symbol).Count() > 1)
|
|
{
|
|
throw new ArgumentException("DownloadAndSave(): All symbols must have the same root ticker.");
|
|
}
|
|
|
|
var dataType = LeanData.GetDataType(_resolution, _tickType);
|
|
|
|
var marketHoursDatabase = MarketHoursDatabase.FromDataFolder();
|
|
|
|
var ticker = symbols.First().ID.Symbol;
|
|
var market = symbols.First().ID.Market;
|
|
|
|
var canonicalSymbol = Symbol.Create(ticker, _securityType, market);
|
|
|
|
var exchangeHours = marketHoursDatabase.GetExchangeHours(canonicalSymbol.ID.Market, canonicalSymbol, _securityType);
|
|
var dataTimeZone = marketHoursDatabase.GetDataTimeZone(canonicalSymbol.ID.Market, canonicalSymbol, _securityType);
|
|
|
|
var historyBySymbol = new Dictionary<Symbol, List<IGrouping<DateTime, BaseData>>>();
|
|
var historyBySymbolDailyOrHour = new Dictionary<Symbol, List<BaseData>>();
|
|
|
|
foreach (var symbol in symbols)
|
|
{
|
|
var historyRequest = new HistoryRequest(
|
|
startTimeUtc,
|
|
endTimeUtc,
|
|
dataType,
|
|
symbol,
|
|
_resolution,
|
|
exchangeHours,
|
|
dataTimeZone,
|
|
_resolution,
|
|
true,
|
|
false,
|
|
DataNormalizationMode.Raw,
|
|
_tickType
|
|
);
|
|
|
|
var history = brokerage.GetHistory(historyRequest)
|
|
.Select(
|
|
x =>
|
|
{
|
|
x.Time = x.Time.ConvertTo(exchangeHours.TimeZone, dataTimeZone);
|
|
return x;
|
|
})
|
|
.ToList();
|
|
|
|
if (_resolution == Resolution.Daily || _resolution == Resolution.Hour)
|
|
{
|
|
historyBySymbolDailyOrHour.Add(symbol, history);
|
|
}
|
|
else
|
|
{
|
|
// group by date in DataTimeZone
|
|
var historyByDate = history.GroupBy(x => x.Time.Date).ToList();
|
|
historyBySymbol.Add(symbol, historyByDate);
|
|
}
|
|
}
|
|
|
|
if (_resolution == Resolution.Daily || _resolution == Resolution.Hour)
|
|
{
|
|
SaveDailyOrHour(symbols, canonicalSymbol, historyBySymbolDailyOrHour);
|
|
}
|
|
else
|
|
{
|
|
SaveMinuteOrSecondOrTick(symbols, startTimeUtc, endTimeUtc, canonicalSymbol, historyBySymbol);
|
|
}
|
|
}
|
|
|
|
private void SaveDailyOrHour(
|
|
List<Symbol> symbols,
|
|
Symbol canonicalSymbol,
|
|
IReadOnlyDictionary<Symbol, List<BaseData>> historyBySymbol)
|
|
{
|
|
var zipFileName = Path.Combine(
|
|
_dataDirectory,
|
|
LeanData.GenerateRelativeZipFilePath(canonicalSymbol, DateTime.MinValue, _resolution, _tickType));
|
|
|
|
var folder = Path.GetDirectoryName(zipFileName);
|
|
if (!Directory.Exists(folder))
|
|
{
|
|
Directory.CreateDirectory(folder);
|
|
}
|
|
|
|
using (var zip = new ZipFile(zipFileName))
|
|
{
|
|
foreach (var symbol in symbols)
|
|
{
|
|
// Load new data rows into a SortedDictionary for easy merge/update
|
|
var newRows = new SortedDictionary<DateTime, string>(historyBySymbol[symbol]
|
|
.ToDictionary(x => x.Time, x => LeanData.GenerateLine(x, _securityType, _resolution)));
|
|
|
|
var rows = new SortedDictionary<DateTime, string>();
|
|
|
|
var zipEntryName = LeanData.GenerateZipEntryName(symbol, DateTime.MinValue, _resolution, _tickType);
|
|
|
|
if (zip.ContainsEntry(zipEntryName))
|
|
{
|
|
// If file exists, we load existing data and perform merge
|
|
using (var stream = new MemoryStream())
|
|
{
|
|
zip[zipEntryName].Extract(stream);
|
|
stream.Seek(0, SeekOrigin.Begin);
|
|
|
|
using (var reader = new StreamReader(stream))
|
|
{
|
|
string line;
|
|
while ((line = reader.ReadLine()) != null)
|
|
{
|
|
var time = Parse.DateTimeExact(line.Substring(0, DateFormat.TwelveCharacter.Length), DateFormat.TwelveCharacter);
|
|
rows[time] = line;
|
|
}
|
|
}
|
|
}
|
|
|
|
foreach (var kvp in newRows)
|
|
{
|
|
rows[kvp.Key] = kvp.Value;
|
|
}
|
|
}
|
|
else
|
|
{
|
|
// No existing file, just use the new data
|
|
rows = newRows;
|
|
}
|
|
|
|
// Loop through the SortedDictionary and write to zip entry
|
|
var sb = new StringBuilder();
|
|
foreach (var kvp in rows)
|
|
{
|
|
// Build the line and append it to the file
|
|
sb.AppendLine(kvp.Value);
|
|
}
|
|
|
|
// Write the zip entry
|
|
if (sb.Length > 0)
|
|
{
|
|
if (zip.ContainsEntry(zipEntryName))
|
|
{
|
|
zip.RemoveEntry(zipEntryName);
|
|
}
|
|
|
|
zip.AddEntry(zipEntryName, sb.ToString());
|
|
}
|
|
}
|
|
|
|
if (zip.Count > 0)
|
|
{
|
|
zip.Save();
|
|
}
|
|
}
|
|
}
|
|
|
|
private void SaveMinuteOrSecondOrTick(
|
|
List<Symbol> symbols,
|
|
DateTime startTimeUtc,
|
|
DateTime endTimeUtc,
|
|
Symbol canonicalSymbol,
|
|
IReadOnlyDictionary<Symbol, List<IGrouping<DateTime, BaseData>>> historyBySymbol)
|
|
{
|
|
var date = startTimeUtc;
|
|
while (date <= endTimeUtc)
|
|
{
|
|
var zipFileName = Path.Combine(
|
|
_dataDirectory,
|
|
LeanData.GenerateRelativeZipFilePath(canonicalSymbol, date, _resolution, _tickType));
|
|
|
|
var folder = Path.GetDirectoryName(zipFileName);
|
|
if (!Directory.Exists(folder))
|
|
{
|
|
Directory.CreateDirectory(folder);
|
|
}
|
|
|
|
if (File.Exists(zipFileName) && !_appendToZips)
|
|
{
|
|
File.Delete(zipFileName);
|
|
}
|
|
|
|
using (var zip = new ZipFile(zipFileName))
|
|
{
|
|
foreach (var symbol in symbols)
|
|
{
|
|
var zipEntryName = LeanData.GenerateZipEntryName(symbol, date, _resolution, _tickType);
|
|
|
|
foreach (var group in historyBySymbol[symbol])
|
|
{
|
|
if (group.Key == date.Date)
|
|
{
|
|
var sb = new StringBuilder();
|
|
foreach (var row in group)
|
|
{
|
|
var line = LeanData.GenerateLine(row, _securityType, _resolution);
|
|
sb.AppendLine(line);
|
|
}
|
|
|
|
if (_appendToZips && zip.ContainsEntry(zipEntryName))
|
|
{
|
|
zip.RemoveEntry(zipEntryName);
|
|
}
|
|
|
|
zip.AddEntry(zipEntryName, sb.ToString());
|
|
break;
|
|
}
|
|
}
|
|
}
|
|
|
|
if (zip.Count > 0)
|
|
{
|
|
zip.Save();
|
|
}
|
|
}
|
|
|
|
date = date.AddDays(1);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Write out the data in LEAN format (minute, second or tick resolutions)
|
|
/// </summary>
|
|
/// <param name="source">IEnumerable source of the data: sorted from oldest to newest.</param>
|
|
/// <remarks>This function overwrites existing data files</remarks>
|
|
private void WriteMinuteOrSecondOrTick(IEnumerable<BaseData> source)
|
|
{
|
|
var sb = new StringBuilder();
|
|
var lastTime = new DateTime();
|
|
|
|
// Loop through all the data and write to file as we go
|
|
foreach (var data in source)
|
|
{
|
|
// Ensure the data is sorted
|
|
if (data.Time < lastTime) throw new Exception("The data must be pre-sorted from oldest to newest");
|
|
|
|
// Based on the security type and resolution, write the data to the zip file
|
|
if (lastTime != DateTime.MinValue && data.Time.Date > lastTime.Date)
|
|
{
|
|
// Write and clear the file contents
|
|
var outputFile = GetZipOutputFileName(_dataDirectory, lastTime);
|
|
WriteFile(outputFile, sb.ToString(), lastTime);
|
|
sb.Clear();
|
|
}
|
|
|
|
lastTime = data.Time;
|
|
|
|
// Build the line and append it to the file
|
|
sb.Append(LeanData.GenerateLine(data, _securityType, _resolution) + Environment.NewLine);
|
|
}
|
|
|
|
// Write the last file
|
|
if (sb.Length > 0)
|
|
{
|
|
var outputFile = GetZipOutputFileName(_dataDirectory, lastTime);
|
|
WriteFile(outputFile, sb.ToString(), lastTime);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Write out the data in LEAN format (daily or hour resolutions)
|
|
/// </summary>
|
|
/// <param name="source">IEnumerable source of the data: sorted from oldest to newest.</param>
|
|
/// <remarks>This function performs a merge (insert/append/overwrite) with the existing Lean zip file</remarks>
|
|
private void WriteDailyOrHour(IEnumerable<BaseData> source)
|
|
{
|
|
var sb = new StringBuilder();
|
|
var lastTime = new DateTime();
|
|
|
|
// Determine file path
|
|
var outputFile = GetZipOutputFileName(_dataDirectory, lastTime);
|
|
|
|
// Load new data rows into a SortedDictionary for easy merge/update
|
|
var newRows = new SortedDictionary<DateTime, string>(source.ToDictionary(x => x.Time, x => LeanData.GenerateLine(x, _securityType, _resolution)));
|
|
SortedDictionary<DateTime, string> rows;
|
|
|
|
if (File.Exists(outputFile))
|
|
{
|
|
// If file exists, we load existing data and perform merge
|
|
rows = LoadHourlyOrDailyFile(outputFile);
|
|
foreach (var kvp in newRows)
|
|
{
|
|
rows[kvp.Key] = kvp.Value;
|
|
}
|
|
}
|
|
else
|
|
{
|
|
// No existing file, just use the new data
|
|
rows = newRows;
|
|
}
|
|
|
|
// Loop through the SortedDictionary and write to file contents
|
|
foreach (var kvp in rows)
|
|
{
|
|
// Build the line and append it to the file
|
|
sb.Append(kvp.Value + Environment.NewLine);
|
|
}
|
|
|
|
// Write the file contents
|
|
if (sb.Length > 0)
|
|
{
|
|
WriteFile(outputFile, sb.ToString(), lastTime);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Loads an existing hourly or daily Lean zip file into a SortedDictionary
|
|
/// </summary>
|
|
private static SortedDictionary<DateTime, string> LoadHourlyOrDailyFile(string fileName)
|
|
{
|
|
var rows = new SortedDictionary<DateTime, string>();
|
|
|
|
using (var zip = ZipFile.Read(fileName))
|
|
{
|
|
using (var stream = new MemoryStream())
|
|
{
|
|
zip[0].Extract(stream);
|
|
stream.Seek(0, SeekOrigin.Begin);
|
|
|
|
using (var reader = new StreamReader(stream))
|
|
{
|
|
string line;
|
|
while ((line = reader.ReadLine()) != null)
|
|
{
|
|
var time = Parse.DateTimeExact(line.Substring(0, DateFormat.TwelveCharacter.Length), DateFormat.TwelveCharacter);
|
|
rows[time] = line;
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
return rows;
|
|
}
|
|
|
|
/// <summary>
|
|
/// Write this file to disk.
|
|
/// </summary>
|
|
/// <param name="filePath">The full path to the new file</param>
|
|
/// <param name="data">The data to write as a string</param>
|
|
/// <param name="date">The date the data represents</param>
|
|
private void WriteFile(string filePath, string data, DateTime date)
|
|
{
|
|
var tempFilePath = filePath + ".tmp";
|
|
|
|
data = data.TrimEnd();
|
|
if (File.Exists(filePath) && !_appendToZips)
|
|
{
|
|
File.Delete(filePath);
|
|
Log.Trace("LeanDataWriter.Write(): Existing deleted: " + filePath);
|
|
}
|
|
|
|
// Create the directory if it doesnt exist
|
|
Directory.CreateDirectory(Path.GetDirectoryName(filePath));
|
|
|
|
if (_appendToZips)
|
|
{
|
|
var entryName = LeanData.GenerateZipEntryName(_symbol, date, _resolution, _tickType);
|
|
Compression.ZipCreateAppendData(filePath, entryName, data, true);
|
|
Log.Trace("LeanDataWriter.Write(): Appended: " + filePath);
|
|
}
|
|
else
|
|
{
|
|
// Write out this data string to a zip file
|
|
Compression.Zip(data, tempFilePath, LeanData.GenerateZipEntryName(_symbol, date, _resolution, _tickType));
|
|
|
|
// Move temp file to the final destination with the appropriate name
|
|
File.Move(tempFilePath, filePath);
|
|
Log.Trace("LeanDataWriter.Write(): Created: " + filePath);
|
|
}
|
|
}
|
|
|
|
/// <summary>
|
|
/// Get the output zip file
|
|
/// </summary>
|
|
/// <param name="baseDirectory">Base output directory for the zip file</param>
|
|
/// <param name="time">Date/time for the data we're writing</param>
|
|
/// <returns>The full path to the output zip file</returns>
|
|
private string GetZipOutputFileName(string baseDirectory, DateTime time)
|
|
{
|
|
return LeanData.GenerateZipFilePath(baseDirectory, _symbol, time, _resolution, _tickType);
|
|
}
|
|
|
|
}
|
|
}
|