/*
* 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
{
///
/// Data writer for saving an IEnumerable of BaseData into the LEAN data directory.
///
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;
///
/// Create a new lean data writer to this base data directory.
///
/// Symbol string
/// Base data directory
/// Resolution of the desired output data
/// The tick type
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);
}
}
///
/// Create a new lean data writer to this base data directory.
///
/// Base data directory
/// Resolution of the desired output data
/// The security type
/// The tick type
public LeanDataWriter(string dataDirectory, Resolution resolution, SecurityType securityType, TickType tickType)
{
_dataDirectory = dataDirectory;
_resolution = resolution;
_securityType = securityType;
_tickType = tickType;
_appendToZips = securityType == SecurityType.Future;
}
///
/// Given the constructor parameters, write out the data in LEAN format.
///
/// IEnumerable source of the data: sorted from oldest to newest.
public void Write(IEnumerable source)
{
switch (_resolution)
{
case Resolution.Daily:
case Resolution.Hour:
WriteDailyOrHour(source);
break;
case Resolution.Minute:
case Resolution.Second:
case Resolution.Tick:
WriteMinuteOrSecondOrTick(source);
break;
}
}
///
/// Downloads historical data from the brokerage and saves it in LEAN format.
///
/// The brokerage from where to fetch the data
/// The list of symbols
/// The starting date/time (UTC)
/// The ending date/time (UTC)
public void DownloadAndSave(IBrokerage brokerage, List 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>>();
var historyBySymbolDailyOrHour = new Dictionary>();
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 symbols,
Symbol canonicalSymbol,
IReadOnlyDictionary> 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(historyBySymbol[symbol]
.ToDictionary(x => x.Time, x => LeanData.GenerateLine(x, _securityType, _resolution)));
var rows = new SortedDictionary();
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 symbols,
DateTime startTimeUtc,
DateTime endTimeUtc,
Symbol canonicalSymbol,
IReadOnlyDictionary>> 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);
}
}
///
/// Write out the data in LEAN format (minute, second or tick resolutions)
///
/// IEnumerable source of the data: sorted from oldest to newest.
/// This function overwrites existing data files
private void WriteMinuteOrSecondOrTick(IEnumerable 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);
}
}
///
/// Write out the data in LEAN format (daily or hour resolutions)
///
/// IEnumerable source of the data: sorted from oldest to newest.
/// This function performs a merge (insert/append/overwrite) with the existing Lean zip file
private void WriteDailyOrHour(IEnumerable 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(source.ToDictionary(x => x.Time, x => LeanData.GenerateLine(x, _securityType, _resolution)));
SortedDictionary 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);
}
}
///
/// Loads an existing hourly or daily Lean zip file into a SortedDictionary
///
private static SortedDictionary LoadHourlyOrDailyFile(string fileName)
{
var rows = new SortedDictionary();
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;
}
///
/// Write this file to disk.
///
/// The full path to the new file
/// The data to write as a string
/// The date the data represents
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);
}
}
///
/// Get the output zip file
///
/// Base output directory for the zip file
/// Date/time for the data we're writing
/// The full path to the output zip file
private string GetZipOutputFileName(string baseDirectory, DateTime time)
{
return LeanData.GenerateZipFilePath(baseDirectory, _symbol, time, _resolution, _tickType);
}
}
}