Compare commits

..
32 Commits
Author SHA1 Message Date
Victor Grau Serrat 8b141a0c28 Fix floats for volume in data.history 2017-10-03 09:11:59 -06:00
VictorandGitHub 7f602d7fcc Update requirements.txt 2017-09-21 11:27:35 -06:00
Victor Grau Serrat 1f56325895 fix price resolution in 1-minute data bundle: 8 decimal places 2017-09-20 23:37:55 -06:00
Victor Grau Serrat 7359cdc48f fix data.history error with tz-aware dataframe 2017-09-20 16:20:57 -06:00
Victor Grau Serrat 42566ca92c Merge branch 'master' of github.com:enigmampc/catalyst 2017-09-20 15:37:44 -06:00
Victor Grau Serrat 09bf875d6c Merge branch 'develop': adds 1-min OHLCV data resolution, fractional coins and 9 decimals of price resolution 2017-09-20 15:20:30 -06:00
Victor Grau Serrat b354837b83 Merge branch 'poloniex-1min-curate' into develop 2017-09-20 15:17:06 -06:00
Victor Grau Serrat a1bc174740 Wrapping up 1min data for Poloniex in backtesting 2017-09-20 15:11:18 -06:00
VictorandGitHub ea27346876 Merge pull request #34 from abnera/patch-1
Added environment.yml for simpler conda installation
2017-09-20 12:50:56 -06:00
Victor Grau Serrat 81bd2d84f0 >=0.2.dev2 for catalyst, since we're in active dev and will change periodically 2017-09-20 12:47:31 -06:00
Abner Ayala-AcevedoandGitHub 06f48cf158 Update to include python-dev 2017-09-20 10:41:42 -07:00
Abner Ayala-AcevedoandGitHub 05a69cfc92 Added environment.yml for simpler conda installation
Simpler conda installation by using environment.yml requirements.
`conda env create -f python2.7-environment.yml`
`Linux or Mac: source activate catalyst`
`Windows: activate catalyst`
2017-09-20 10:16:29 -07:00
Victor Grau Serrat 36c2564bb0 splitting plot styles: dark for live, default for backtesting 2017-09-20 11:05:09 -06:00
Victor Grau Serrat 91e71c5e38 WIP: bundling 1min data 2017-09-20 09:15:39 -06:00
Victor Grau Serrat e761433d06 Merge branch 'poloniex-1min-curate' of github.com:enigmampc/catalyst into poloniex-1min-curate 2017-09-18 09:45:12 -06:00
Victor Grau Serrat 6fddb92563 WIP: trades to disk - no append/no ingestion 2017-09-18 09:44:19 -06:00
Victor Grau Serrat 4a4277d9d1 WIP: curating 1min Poloniex data - no append 2017-09-18 09:44:19 -06:00
Victor Grau Serrat 3361b09ac2 Merge branch 'fractional-coins' into develop 2017-09-15 16:06:30 -06:00
Victor Grau Serrat 01eefd67e0 ingestion switch to create_writers when ingesting locally 2017-09-15 10:51:04 -06:00
Victor Grau Serrat d124125258 WIP: trades to disk - no append/no ingestion 2017-09-15 09:32:01 -06:00
Victor Grau Serrat 72e07e242f WIP: curating 1min Poloniex data - no append 2017-09-15 09:32:01 -06:00
Victor Grau Serrat c3897cfa5a ENH: wrapping up asset min_trade_size for fractional coinsup to 1/100000000th of a coin 2017-09-14 15:22:51 -06:00
Andrew CampbellandVictor Grau Serrat e5a137f205 ENH: Bound trade amount with asset specific min trade size 2017-09-14 10:19:05 -06:00
Victor Grau Serrat 48143d3212 WIP: Fixes 1/1000 price issue in history, and works with full coins. Requires matching-version 'catalyst ingest' 2017-09-13 17:22:11 -06:00
fredfortier 8850657f26 Fixed path issue with obsolete branch 2017-09-07 14:26:55 -04:00
fredfortier 6e6c62533b Fixes to the graph timeline axis 2017-09-05 11:15:34 -04:00
fredfortier 6a98a937dd Minor fix in the graph logic 2017-09-05 01:15:40 -04:00
fredfortier c4900af088 Minor fix in the graph logic 2017-09-05 00:58:11 -04:00
fredfortier 6e3017010f Added the initial version of a live graph 2017-09-05 00:51:38 -04:00
fredfortier 7247b761d5 Fixed an issue with failed orders 2017-09-04 11:33:21 -04:00
fredfortier a58e3522a9 Improved handling of insufficient funds on bittrex 2017-09-04 11:30:32 -04:00
fredfortier ad95369028 Improved handling of insufficient funds on bittrex 2017-09-04 11:23:40 -04:00
27 changed files with 762 additions and 188 deletions
+10 -2
View File
@@ -210,6 +210,12 @@ def ipython_only(option):
help='The base currency used to calculate statistics ' help='The base currency used to calculate statistics '
'(e.g. usd, btc, eth).', '(e.g. usd, btc, eth).',
) )
@click.option(
'--live-graph/--no-live-graph',
is_flag=True,
default=False,
help='Display live graph.',
)
@click.pass_context @click.pass_context
def run(ctx, def run(ctx,
algofile, algofile,
@@ -227,7 +233,8 @@ def run(ctx,
live, live,
exchange_name, exchange_name,
algo_namespace, algo_namespace,
base_currency): base_currency,
live_graph):
"""Run a backtest for the given algorithm. """Run a backtest for the given algorithm.
""" """
@@ -283,7 +290,8 @@ def run(ctx,
live=live, live=live,
exchange=exchange_name, exchange=exchange_name,
algo_namespace=algo_namespace, algo_namespace=algo_namespace,
base_currency=base_currency base_currency=base_currency,
live_graph=live_graph
) )
if output == '-': if output == '-':
+6 -9
View File
@@ -125,6 +125,7 @@ from catalyst.utils.factory import create_simulation_parameters
from catalyst.utils.math_utils import ( from catalyst.utils.math_utils import (
tolerant_equals, tolerant_equals,
round_if_near_integer, round_if_near_integer,
round_nearest
) )
from catalyst.utils.pandas_utils import clear_dataframe_indexer_caches from catalyst.utils.pandas_utils import clear_dataframe_indexer_caches
from catalyst.utils.preprocess import preprocess from catalyst.utils.preprocess import preprocess
@@ -1488,7 +1489,7 @@ class TradingAlgorithm(object):
def _calculate_order(self, asset, amount, def _calculate_order(self, asset, amount,
limit_price=None, stop_price=None, style=None): limit_price=None, stop_price=None, style=None):
amount = self.round_order(amount) amount = self.round_order(amount, asset)
# Raises a ZiplineError if invalid parameters are detected. # Raises a ZiplineError if invalid parameters are detected.
self.validate_order_params(asset, self.validate_order_params(asset,
@@ -1505,16 +1506,13 @@ class TradingAlgorithm(object):
return amount, style return amount, style
@staticmethod @staticmethod
def round_order(amount): def round_order(amount, asset):
""" """
Convert number of shares to an integer. Converts the number of shares to the smallest tradable lot size for
the asset being ordered.
By default, truncates to the integer share count that's either within
.0001 of amount or closer to zero.
E.g. 3.9999 -> 4.0; 5.5 -> 5.0; -5.5 -> -5.0
""" """
return int(round_if_near_integer(amount)) return round_nearest(amount, asset.min_trade_size)
def validate_order_params(self, def validate_order_params(self,
asset, asset,
@@ -1550,7 +1548,6 @@ class TradingAlgorithm(object):
self.updated_portfolio(), self.updated_portfolio(),
self.get_datetime(), self.get_datetime(),
self.trading_client.current_data) self.trading_client.current_data)
@staticmethod @staticmethod
def __convert_order_params_for_blotter(limit_price, stop_price, style): def __convert_order_params_for_blotter(limit_price, stop_price, style):
""" """
+11 -4
View File
@@ -59,6 +59,7 @@ cdef class Asset:
cdef readonly object exchange cdef readonly object exchange
cdef readonly object exchange_full cdef readonly object exchange_full
cdef readonly object min_trade_size
_kwargnames = frozenset({ _kwargnames = frozenset({
'sid', 'sid',
@@ -70,6 +71,7 @@ cdef class Asset:
'auto_close_date', 'auto_close_date',
'exchange', 'exchange',
'exchange_full', 'exchange_full',
'min_trade_size',
}) })
def __init__(self, def __init__(self,
@@ -81,7 +83,8 @@ cdef class Asset:
object end_date=None, object end_date=None,
object first_traded=None, object first_traded=None,
object auto_close_date=None, object auto_close_date=None,
object exchange_full=None): object exchange_full=None,
object min_trade_size=None):
self.sid = sid self.sid = sid
self.sid_hash = hash(sid) self.sid_hash = hash(sid)
@@ -94,6 +97,7 @@ cdef class Asset:
self.end_date = end_date self.end_date = end_date
self.first_traded = first_traded self.first_traded = first_traded
self.auto_close_date = auto_close_date self.auto_close_date = auto_close_date
self.min_trade_size = min_trade_size
def __int__(self): def __int__(self):
return self.sid return self.sid
@@ -148,7 +152,8 @@ cdef class Asset:
def __repr__(self): def __repr__(self):
attrs = ('symbol', 'asset_name', 'exchange', attrs = ('symbol', 'asset_name', 'exchange',
'start_date', 'end_date', 'first_traded', 'auto_close_date') 'start_date', 'end_date', 'first_traded', 'auto_close_date',
'min_trade_size')
tuples = ((attr, repr(getattr(self, attr, None))) tuples = ((attr, repr(getattr(self, attr, None)))
for attr in attrs) for attr in attrs)
strings = ('%s=%s' % (t[0], t[1]) for t in tuples) strings = ('%s=%s' % (t[0], t[1]) for t in tuples)
@@ -170,7 +175,8 @@ cdef class Asset:
self.end_date, self.end_date,
self.first_traded, self.first_traded,
self.auto_close_date, self.auto_close_date,
self.exchange_full)) self.exchange_full,
self.min_trade_size))
cpdef to_dict(self): cpdef to_dict(self):
""" """
@@ -186,6 +192,7 @@ cdef class Asset:
'auto_close_date': self.auto_close_date, 'auto_close_date': self.auto_close_date,
'exchange': self.exchange, 'exchange': self.exchange,
'exchange_full': self.exchange_full, 'exchange_full': self.exchange_full,
'min_trade_size': self.min_trade_size
} }
@classmethod @classmethod
@@ -234,7 +241,7 @@ cdef class Equity(Asset):
def __repr__(self): def __repr__(self):
attrs = ('symbol', 'asset_name', 'exchange', attrs = ('symbol', 'asset_name', 'exchange',
'start_date', 'end_date', 'first_traded', 'auto_close_date', 'start_date', 'end_date', 'first_traded', 'auto_close_date',
'exchange_full') 'exchange_full', 'min_trade_size')
tuples = ((attr, repr(getattr(self, attr, None))) tuples = ((attr, repr(getattr(self, attr, None)))
for attr in attrs) for attr in attrs)
strings = ('%s=%s' % (t[0], t[1]) for t in tuples) strings = ('%s=%s' % (t[0], t[1]) for t in tuples)
+2 -1
View File
@@ -39,7 +39,8 @@ equities = sa.Table(
sa.Column('first_traded', sa.Integer), sa.Column('first_traded', sa.Integer),
sa.Column('auto_close_date', sa.Integer), sa.Column('auto_close_date', sa.Integer),
sa.Column('exchange', sa.Text), sa.Column('exchange', sa.Text),
sa.Column('exchange_full', sa.Text) sa.Column('exchange_full', sa.Text),
sa.Column('min_trade_size', sa.Float)
) )
equity_symbol_mappings = sa.Table( equity_symbol_mappings = sa.Table(
+3
View File
@@ -73,6 +73,7 @@ _equities_defaults = {
'exchange': None, 'exchange': None,
# optional, something like "New York Stock Exchange" # optional, something like "New York Stock Exchange"
'exchange_full': None, 'exchange_full': None,
'min_trade_size': 1
} }
# Default values for the futures DataFrame # Default values for the futures DataFrame
@@ -390,6 +391,8 @@ class AssetDBWriter(object):
The date on which to close any positions in this asset. The date on which to close any positions in this asset.
exchange : str exchange : str
The exchange where this asset is traded. The exchange where this asset is traded.
min_trade_size: float, optional
The minimum denomination this asset can be traded.
The index of this dataframe should contain the sids. The index of this dataframe should contain the sids.
futures : pd.DataFrame, optional futures : pd.DataFrame, optional
+184 -71
View File
@@ -1,12 +1,10 @@
import json, time, csv import json, time, csv
from datetime import datetime from datetime import datetime
import pandas as pd import pandas as pd
import os import os, time, shutil, requests, logbook
import time
import requests
import logbook
DT_START = time.mktime(datetime(2010, 1, 1, 0, 0).timetuple()) DT_START = int(time.mktime(datetime(2010, 1, 1, 0, 0).timetuple()))
DT_END = int(time.time())
CSV_OUT_FOLDER = '/var/tmp/catalyst/data/poloniex/' CSV_OUT_FOLDER = '/var/tmp/catalyst/data/poloniex/'
CONN_RETRIES = 2 CONN_RETRIES = 2
@@ -14,9 +12,9 @@ logbook.StderrHandler().push_application()
log = logbook.Logger(__name__) log = logbook.Logger(__name__)
class PoloniexCurator(object): class PoloniexCurator(object):
""" '''
OHLCV data feed generator for crypto data. Based on Poloniex market data OHLCV data feed generator for crypto data. Based on Poloniex market data
""" '''
_api_path = 'https://poloniex.com/public?' _api_path = 'https://poloniex.com/public?'
currency_pairs = [] currency_pairs = []
@@ -29,6 +27,9 @@ class PoloniexCurator(object):
log.error('Failed to create data folder: %s' % CSV_OUT_FOLDER) log.error('Failed to create data folder: %s' % CSV_OUT_FOLDER)
log.exception(e) log.exception(e)
'''
Retrieves and returns all currency pairs from the exchange
'''
def get_currency_pairs(self): def get_currency_pairs(self):
url = self._api_path + 'command=returnTicker' url = self._api_path + 'command=returnTicker'
@@ -47,98 +48,210 @@ class PoloniexCurator(object):
log.debug('Currency pairs retrieved successfully: %d' % (len(self.currency_pairs))) log.debug('Currency pairs retrieved successfully: %d' % (len(self.currency_pairs)))
def _get_start_date(self, csv_fn):
''' Function returns latest appended date, if the file has been previously written '''
the last line is an empty one, so we have to read the second to last line Helper function that reads tradeID and date fields from CSV readline
'''
def _retrieve_tradeID_date(self, row):
tId = int(row.split(',')[0])
d = pd.to_datetime( row.split(',')[1], infer_datetime_format=True).value // 10 ** 9
return tId, d
'''
Retrieves TradeHistory from exchange for a given currencyPair between start and end dates.
If no start date is provided, uses a system-wide one (beginning of time for cryptotrading)
If no end date is provided, 'now' is used
Stores results in CSV file on disk.
This function is called recursively to work around the limitations imposed by the provider API.
'''
def retrieve_trade_history(self, currencyPair, start=DT_START, end=DT_END, temp=None):
csv_fn = CSV_OUT_FOLDER + 'crypto_trades-' + currencyPair + '.csv'
'''
Check what data we already have on disk, reading first and last lines from file.
Data is stored on file from NEWEST to OLDEST.
''' '''
try: try:
with open(csv_fn, 'ab+') as f: with open(csv_fn, 'ab+') as f:
f.seek(0, os.SEEK_END) # First check file is not zero size f.seek(0, os.SEEK_END)
if(f.tell() > 2): if(f.tell() > 2): # First check file is not zero size
f.seek(-2, os.SEEK_END) # Jump to the second last byte. f.seek(0) # Go to the beginning to read first line
while f.read(1) != b"\n": # Until EOL is found... last_tradeID, end_file = self._retrieve_tradeID_date(f.readline())
f.seek(-2, os.SEEK_CUR) # ...jump back the read byte plus one more. f.seek(-2, os.SEEK_END) # Jump to the second last byte.
lastrow = f.readline() while f.read(1) != b"\n": # Until EOL is found...
return int(lastrow.split(',')[0]) + 300 f.seek(-2, os.SEEK_CUR) # ...jump back the read byte plus one more.
first_tradeID, start_file = self._retrieve_tradeID_date(f.readline())
if( first_tradeID == 1 and end_file + 3600 > DT_END ):
return
except Exception as e: except Exception as e:
log.error('Error opening file: %s' % csv_fn) log.error('Error opening file: %s' % csv_fn)
log.exception(e) log.exception(e)
return DT_START '''
Poloniex API limits querying TradeHistory to intervals smaller than 1 month,
so we make sure that start date is never more than 1 month apart from end date
'''
if( end - start > 2419200 ): # 60 s/min * 60 min/hr * 24 hr/day * 28 days
newstart = end - 2419200
else:
newstart = start
def get_data(self, currencyPair, start, end=9999999999, period=300): log.debug(currencyPair+': Retrieving from '+str(newstart)+' to '+str(end) +'\t '
url = self._api_path + 'command=returnChartData&currencyPair=' + currencyPair + '&start=' + str(start) + '&end=' + str(end) + '&period=' + str(period) + time.ctime(newstart) + ' - '+ time.ctime(end))
url = self._api_path + 'command=returnTradeHistory&currencyPair=' + currencyPair + '&start=' + str(newstart) + '&end=' + str(end)
try: try:
response = requests.get(url) response = requests.get(url)
except Exception as e: except Exception as e:
log.error('Failed to retrieve candlestick chart data for %s' % currencyPair) log.error('Failed to retrieve trade history data for %s' % currencyPair)
log.exception(e) log.exception(e)
return None return None
else:
if isinstance(response.json(), dict) and response.json()['error']:
log.error('Failed to to retrieve trade history data for %s: %s' % (currencyPair,response.json()['error']))
exit(1)
'''
If we get to transactionId == 1, and we already have that on disk,
we got to the end of TradeHistory for this coin.
'''
if('first_tradeID' in locals() and response.json()[-1]['tradeID'] == first_tradeID):
return
'''
There are primarily two scenarios:
a) There is newer data available that we need to add at the beginning
of the file. We'll retrieve all what we need until we get to what
we already have, writing it to a temporary file; and we will write
that at the beginning of our existing file.
b) We are going back in time, appending at the end of our existing
TradeHistory until the first transaction for this currencyPair
'''
try:
if( 'end_file' in locals() and end_file + 3600 < end):
if (temp is None):
temp = os.tmpfile()
tempcsv = csv.writer(temp)
for item in response.json():
if( item['tradeID'] <= last_tradeID ):
continue
tempcsv.writerow([
item['tradeID'],
item['date'],
item['type'],
item['rate'],
item['amount'],
item['total'],
item['globalTradeID']
])
if( response.json()[-1]['tradeID'] > last_tradeID ):
end = pd.to_datetime( response.json()[-1]['date'], infer_datetime_format=True).value // 10 ** 9
self.retrieve_trade_history(currencyPair, start, end, temp=temp)
else:
with open(csv_fn,'rb+') as f:
shutil.copyfileobj(f,temp)
f.seek(0)
temp.seek(0)
shutil.copyfileobj(temp,f)
temp.close()
end = start_file
else:
with open(csv_fn, 'ab') as csvfile:
csvwriter = csv.writer(csvfile)
for item in response.json():
if( 'first_tradeID' in locals() and item['tradeID'] >= first_tradeID ):
continue
csvwriter.writerow([
item['tradeID'],
item['date'],
item['type'],
item['rate'],
item['amount'],
item['total'],
item['globalTradeID']
])
end = pd.to_datetime( response.json()[-1]['date'], infer_datetime_format=True).value // 10 ** 9
except Exception as e:
log.error('Error opening %s' % csv_fn)
log.exception(e)
'''
If we got here, we aren't done yet. Call recursively with 'end' times
that go sequentially back in time.
'''
self.retrieve_trade_history(currencyPair, start, end)
return response.json()
''' '''
Pulls latest data for a single pair Generates OHLCV dataframe from a dataframe containing all TradeHistory
by resampling with 1-minute period
''' '''
def append_data_single_pair(self, currencyPair, repeat=0): def generate_ohlcv(self, df):
log.debug('Getting data for %s' % currencyPair) df.set_index('date', inplace=True) # Index by date
csv_fn = CSV_OUT_FOLDER + 'crypto_prices-' + currencyPair + '.csv' vol = df['total'].to_frame('volume') # Will deal with vol separately, as ohlc() messes it up
start = self._get_start_date(csv_fn) df.drop('total', axis=1, inplace=True) # Drop volume data from dataframe
# Only fetch data if more than 5min have passed since last fetch ohlc = df.resample('T').ohlc() # Resample OHLC in 1min bins
if (time.time() > start): ohlc.columns = ohlc.columns.map(lambda t: t[1]) # Raname columns by dropping 'rate'
data = self.get_data(currencyPair, start) closes = ohlc['close'].fillna(method='pad') # Pad forward missing 'close'
if data is not None: ohlc = ohlc.apply(lambda x: x.fillna(closes)) # Fill N/A with last close
try: vol = vol.resample('T').sum().fillna(0) # Add volumes by bin
with open(csv_fn, 'ab') as csvfile: ohlcv = pd.concat([ohlc,vol], axis=1) # Concatenate OHLC + Volume
csvwriter = csv.writer(csvfile) return ohlcv
for item in data:
if item['date'] == 0:
continue
csvwriter.writerow([
item['date'],
item['open'],
item['high'],
item['low'],
item['close'],
item['volume'],
])
except Exception as e:
log.error('Error opening %s' % csv_fn)
log.exception(e)
elif (repeat < CONN_RETRIES):
log.debug('Retrying: attemt %d' % (repeat+1) )
self.append_data_single_pair(currencyPair, repeat + 1)
''' '''
Pulls latest data for all currency pairs Generates OHLCV data file with 1minute bars from TradeHistory on disk
''' '''
def append_data(self): def write_ohlcv_file(self, currencyPair):
for currencyPair in self.currency_pairs: csv_trades = CSV_OUT_FOLDER + 'crypto_trades-' + currencyPair + '.csv'
self.append_data_single_pair(currencyPair) csv_1min = CSV_OUT_FOLDER + 'crypto_1min-' + currencyPair + '.csv'
# Rate limit is 6 calls per second, sleep 1sec/6 to be safe if( os.path.isfile(csv_1min) ):
time.sleep(0.17) log.debug(currencyPair+': 1min data already present. Delete the file if you want to rebuild it.')
else:
df = pd.read_csv(csv_trades, names=['tradeID','date','type','rate','amount','total','globalTradeID'],
dtype = {'tradeID': int, 'date': str, 'type': str, 'rate': float, 'amount': float, 'total': float, 'globalTradeID': int } )
df.drop(['tradeID','type','amount','globalTradeID'], axis=1, inplace=True)
df['date'] = pd.to_datetime(df['date'], infer_datetime_format=True)
ohlcv = self.generate_ohlcv(df)
try:
with open(csv_1min, 'ab') as csvfile:
csvwriter = csv.writer(csvfile)
for item in ohlcv.itertuples():
if item.Index == 0:
continue
csvwriter.writerow([
item.Index.value // 10 ** 9,
item.open,
item.high,
item.low,
item.close,
item.volume,
])
except Exception as e:
log.error('Error opening %s' % csv_fn)
log.exception(e)
log.debug(currencyPair+': Generated 1min OHLCV data.')
''' '''
Returns a data frame for all pairs, or for the requests currency pair. Returns a data frame for a given currencyPair from data on disk
Makes sure data is up to date
''' '''
def to_dataframe(self, start, end, currencyPair=None): def onemin_to_dataframe(self, currencyPair, start, end):
csv_fn = CSV_OUT_FOLDER + 'crypto_prices-' + currencyPair + '.csv' csv_fn = CSV_OUT_FOLDER + 'crypto_1min-' + currencyPair + '.csv'
last_date = self._get_start_date(csv_fn) df = pd.read_csv(csv_fn, names=['date', 'open', 'high', 'low', 'close', 'volume'])
if last_date + 300 < end or not os.path.exists(csv_fn): df['date'] = pd.to_datetime(df['date'],unit='s')
# get latest data
self.append_data_single_pair(currencyPair)
# CSV holds the latest snapshot
df = pd.read_csv(csv_fn, names=['date', 'open', 'high', 'low', 'close', 'volume'])
df['date']=pd.to_datetime(df['date'],unit='s')
df.set_index('date', inplace=True) df.set_index('date', inplace=True)
return df[start : end]
return df[datetime.fromtimestamp(start):datetime.fromtimestamp(end-1)]
if __name__ == '__main__': if __name__ == '__main__':
pc = PoloniexCurator() pc = PoloniexCurator()
pc.get_currency_pairs() pc.get_currency_pairs()
pc.append_data()
for currencyPair in pc.currency_pairs:
pc.retrieve_trade_history(currencyPair)
pc.write_ohlcv_file(currencyPair)
+2 -4
View File
@@ -215,13 +215,11 @@ cpdef _read_bcolz_data(ctable_t table,
else: else:
continue continue
if column_name in ['open', 'high', 'low', 'close']: if column_name in ['open', 'high', 'low', 'close', 'volume']:
where_nan = (outbuf == 0) where_nan = (outbuf == 0)
outbuf_as_float = outbuf.astype(float64) * .000001 outbuf_as_float = outbuf.astype(float64) * .000000001
outbuf_as_float[where_nan] = NAN outbuf_as_float[where_nan] = NAN
results.append(outbuf_as_float) results.append(outbuf_as_float)
elif column_name != 'volume':
results.append(outbuf.astype(uint32))
else: else:
results.append(outbuf) results.append(outbuf)
return results return results
+1 -1
View File
@@ -491,7 +491,7 @@ class BaseBundle(object):
data_frequency, data_frequency,
) )
raw_data.index = pd.to_datetime(raw_data.index, utc=True) raw_data.index = pd.to_datetime(raw_data.index, utc=True)
raw_data.index = raw_data.index.tz_localize('UTC') #raw_data.index = raw_data.index.tz_localize('UTC')
# Filter incoming data to fit start and end sessions. # Filter incoming data to fit start and end sessions.
raw_data = raw_data[ raw_data = raw_data[
+1
View File
@@ -24,6 +24,7 @@ class BasePricingBundle(BaseBundle):
('start_date', 'datetime64[ns]'), ('start_date', 'datetime64[ns]'),
('end_date', 'datetime64[ns]'), ('end_date', 'datetime64[ns]'),
('ac_date', 'datetime64[ns]'), ('ac_date', 'datetime64[ns]'),
('min_trade_size', 'float'),
] ]
@lazyval @lazyval
+31 -14
View File
@@ -13,6 +13,8 @@
# See the License for the specific language governing permissions and # See the License for the specific language governing permissions and
# limitations under the License. # limitations under the License.
import sys
from datetime import datetime from datetime import datetime
import pandas as pd import pandas as pd
@@ -23,6 +25,8 @@ from catalyst.data.bundles.core import register_bundle
from catalyst.data.bundles.base_pricing import BaseCryptoPricingBundle from catalyst.data.bundles.base_pricing import BaseCryptoPricingBundle
from catalyst.utils.memoize import lazyval from catalyst.utils.memoize import lazyval
from catalyst.curate.poloniex import PoloniexCurator
class PoloniexBundle(BaseCryptoPricingBundle): class PoloniexBundle(BaseCryptoPricingBundle):
@lazyval @lazyval
def name(self): def name(self):
@@ -36,7 +40,7 @@ class PoloniexBundle(BaseCryptoPricingBundle):
def frequencies(self): def frequencies(self):
return set(( return set((
'daily', 'daily',
#'5-minute', 'minute',
)) ))
@lazyval @lazyval
@@ -75,12 +79,14 @@ class PoloniexBundle(BaseCryptoPricingBundle):
start_date = sym_data.index[0] start_date = sym_data.index[0]
end_date = sym_data.index[-1] end_date = sym_data.index[-1]
ac_date = end_date + pd.Timedelta(days=1) ac_date = end_date + pd.Timedelta(days=1)
min_trade_size = 0.00000001
return ( return (
sym_md.symbol, sym_md.symbol,
start_date, start_date,
end_date, end_date,
ac_date, ac_date,
min_trade_size,
) )
def fetch_raw_symbol_frame(self, def fetch_raw_symbol_frame(self,
@@ -90,22 +96,28 @@ class PoloniexBundle(BaseCryptoPricingBundle):
start_date, start_date,
end_date, end_date,
frequency): frequency):
raw = pd.read_json(
self._format_data_url( if(frequency == 'minute'):
api_key, pc = PoloniexCurator()
symbol, raw = pc.onemin_to_dataframe(symbol, start_date, end_date)
start_date,
end_date, else:
frequency, raw = pd.read_json(
), self._format_data_url(
orient='records', api_key,
) symbol,
raw.set_index('date', inplace=True) start_date,
end_date,
frequency,
),
orient='records',
)
raw.set_index('date', inplace=True)
# BcolzDailyBarReader introduces a 1/1000 factor in the way pricing is stored # BcolzDailyBarReader introduces a 1/1000 factor in the way pricing is stored
# on disk, which we compensate here to get the right pricing amounts # on disk, which we compensate here to get the right pricing amounts
# ref: data/us_equity_pricing.py # ref: data/us_equity_pricing.py
scale = 1000 scale = 1
raw.loc[:, 'open'] /= scale raw.loc[:, 'open'] /= scale
raw.loc[:, 'high'] /= scale raw.loc[:, 'high'] /= scale
raw.loc[:, 'low'] /= scale raw.loc[:, 'low'] /= scale
@@ -166,4 +178,9 @@ register_bundle(PoloniexBundle, ['USDT_BTC',])
For a production environment make sure to use (to bundle all pairs): For a production environment make sure to use (to bundle all pairs):
register_bundle(PoloniexBundle) register_bundle(PoloniexBundle)
''' '''
register_bundle(PoloniexBundle, create_writers=False)
if 'ingest' in sys.argv and '-c' in sys.argv:
register_bundle(PoloniexBundle)
else:
register_bundle(PoloniexBundle, create_writers=False)
+4 -1
View File
@@ -18,6 +18,7 @@ from numpy import (
full, full,
nan, nan,
int64, int64,
float64,
zeros zeros
) )
from six import iteritems, with_metaclass from six import iteritems, with_metaclass
@@ -70,7 +71,9 @@ class AssetDispatchBarReader(with_metaclass(ABCMeta)):
return self._dt_window_size(start_dt, end_dt), num_sids return self._dt_window_size(start_dt, end_dt), num_sids
def _make_raw_array_out(self, field, shape): def _make_raw_array_out(self, field, shape):
if field != 'volume' and field != 'sid': if field == 'volume':
out = zeros(shape, dtype=float64)
elif field != 'sid':
out = full(shape, nan) out = full(shape, nan)
else: else:
out = zeros(shape, dtype=int64) out = zeros(shape, dtype=int64)
+1 -1
View File
@@ -38,7 +38,7 @@ from catalyst.utils.numpy_utils import float64_dtype
from catalyst.utils.pandas_utils import find_in_sorted_index from catalyst.utils.pandas_utils import find_in_sorted_index
# Default number of decimal places used for rounding asset prices. # Default number of decimal places used for rounding asset prices.
DEFAULT_ASSET_PRICE_DECIMALS = 3 DEFAULT_ASSET_PRICE_DECIMALS = 9
class HistoryCompatibleUSEquityAdjustmentReader(object): class HistoryCompatibleUSEquityAdjustmentReader(object):
+35 -35
View File
@@ -39,7 +39,7 @@ from catalyst.data._minute_bar_internal import (
from catalyst.gens.sim_engine import NANOS_IN_MINUTE from catalyst.gens.sim_engine import NANOS_IN_MINUTE
from catalyst.data.bar_reader import BarReader, NoDataOnDate from catalyst.data.bar_reader import BarReader, NoDataOnDate
from catalyst.data.us_equity_pricing import check_uint32_safe from catalyst.data.us_equity_pricing import check_uint64_safe
from catalyst.utils.calendars import get_calendar from catalyst.utils.calendars import get_calendar
from catalyst.utils.cli import maybe_show_progress from catalyst.utils.cli import maybe_show_progress
from catalyst.utils.memoize import lazyval from catalyst.utils.memoize import lazyval
@@ -52,7 +52,7 @@ FUTURES_MINUTES_PER_DAY = 1440
DEFAULT_EXPECTEDLEN = US_EQUITIES_MINUTES_PER_DAY * 252 * 15 DEFAULT_EXPECTEDLEN = US_EQUITIES_MINUTES_PER_DAY * 252 * 15
OHLC_RATIO = 1000 OHLC_RATIO = 100000000
class BcolzMinuteOverlappingData(Exception): class BcolzMinuteOverlappingData(Exception):
@@ -114,15 +114,15 @@ def _sid_subdir_path(sid):
def convert_cols(cols, scale_factor, sid, invalid_data_behavior): def convert_cols(cols, scale_factor, sid, invalid_data_behavior):
"""Adapt OHLCV columns into uint32 columns. """Adapt OHLCV columns into uint64 columns.
Parameters Parameters
---------- ----------
cols : dict cols : dict
A dict mapping each column name (open, high, low, close, volume) A dict mapping each column name (open, high, low, close, volume)
to a float column to convert to uint32. to a float column to convert to uint64.
scale_factor : int scale_factor : int
Factor to use to scale float values before converting to uint32. Factor to use to scale float values before converting to uint64.
sid : int sid : int
Sid of the relevant asset, for logging. Sid of the relevant asset, for logging.
invalid_data_behavior : str invalid_data_behavior : str
@@ -135,6 +135,7 @@ def convert_cols(cols, scale_factor, sid, invalid_data_behavior):
scaled_highs = np.nan_to_num(cols['high']) * scale_factor scaled_highs = np.nan_to_num(cols['high']) * scale_factor
scaled_lows = np.nan_to_num(cols['low']) * scale_factor scaled_lows = np.nan_to_num(cols['low']) * scale_factor
scaled_closes = np.nan_to_num(cols['close']) * scale_factor scaled_closes = np.nan_to_num(cols['close']) * scale_factor
scaled_volumes = np.nan_to_num(cols['volume']) * scale_factor
exclude_mask = np.zeros_like(scaled_opens, dtype=bool) exclude_mask = np.zeros_like(scaled_opens, dtype=bool)
@@ -143,11 +144,12 @@ def convert_cols(cols, scale_factor, sid, invalid_data_behavior):
('high', scaled_highs), ('high', scaled_highs),
('low', scaled_lows), ('low', scaled_lows),
('close', scaled_closes), ('close', scaled_closes),
('volume', scaled_volumes),
]: ]:
max_val = scaled_col.max() max_val = scaled_col.max()
try: try:
check_uint32_safe(max_val, col_name) check_uint64_safe(max_val, col_name)
except ValueError: except ValueError:
if invalid_data_behavior == 'raise': if invalid_data_behavior == 'raise':
raise raise
@@ -155,20 +157,20 @@ def convert_cols(cols, scale_factor, sid, invalid_data_behavior):
if invalid_data_behavior == 'warn': if invalid_data_behavior == 'warn':
logger.warn( logger.warn(
'Values for sid={}, col={} contain some too large for ' 'Values for sid={}, col={} contain some too large for '
'uint32 (max={}), filtering them out', 'uint64 (max={}), filtering them out',
sid, col_name, max_val, sid, col_name, max_val,
) )
# We want to exclude all rows that have an unsafe value in # We want to exclude all rows that have an unsafe value in
# this column. # this column.
exclude_mask &= (scaled_col >= np.iinfo(np.uint32).max) exclude_mask &= (scaled_col >= np.iinfo(np.uint64).max)
# Convert all cols to uint32. # Convert all cols to uint32.
opens = scaled_opens.astype(np.uint32) opens = scaled_opens.astype(np.uint64)
highs = scaled_highs.astype(np.uint32) highs = scaled_highs.astype(np.uint64)
lows = scaled_lows.astype(np.uint32) lows = scaled_lows.astype(np.uint64)
closes = scaled_closes.astype(np.uint32) closes = scaled_closes.astype(np.uint64)
volumes = cols['volume'].astype(np.uint32) volumes = scaled_volumes.astype(np.uint64)
# Exclude rows with unsafe values by setting to zero. # Exclude rows with unsafe values by setting to zero.
opens[exclude_mask] = 0 opens[exclude_mask] = 0
@@ -288,7 +290,7 @@ class BcolzMinuteBarMetadata(object):
ohlc_ratio : int ohlc_ratio : int
The default ratio by which to multiply the pricing data to The default ratio by which to multiply the pricing data to
convert the floats from floats to an integer to fit within convert the floats from floats to an integer to fit within
the np.uint32. If ohlc_ratios_per_sid is None or does not the np.uint64. If ohlc_ratios_per_sid is None or does not
contain a mapping for a given sid, this ratio is used. contain a mapping for a given sid, this ratio is used.
ohlc_ratios_per_sid : dict ohlc_ratios_per_sid : dict
A dict mapping each sid in the output to the factor by A dict mapping each sid in the output to the factor by
@@ -372,13 +374,13 @@ class BcolzMinuteBarWriter(object):
The last trading session in the data set. The last trading session in the data set.
default_ohlc_ratio : int, optional default_ohlc_ratio : int, optional
The default ratio by which to multiply the pricing data to The default ratio by which to multiply the pricing data to
convert from floats to integers that fit within np.uint32. If convert from floats to integers that fit within np.uint64. If
ohlc_ratios_per_sid is None or does not contain a mapping for a ohlc_ratios_per_sid is None or does not contain a mapping for a
given sid, this ratio is used. Default is OHLC_RATIO (1000). given sid, this ratio is used. Default is OHLC_RATIO (10^8).
ohlc_ratios_per_sid : dict, optional ohlc_ratios_per_sid : dict, optional
A dict mapping each sid in the output to the ratio by which to A dict mapping each sid in the output to the ratio by which to
multiply the pricing data to convert the floats from floats to multiply the pricing data to convert the floats from floats to
an integer to fit within the np.uint32. an integer to fit within the np.uint64.
expectedlen : int, optional expectedlen : int, optional
The expected length of the dataset, used when creating the initial The expected length of the dataset, used when creating the initial
bcolz ctable. bcolz ctable.
@@ -401,11 +403,9 @@ class BcolzMinuteBarWriter(object):
Each individual asset's data is stored as a bcolz table with a column for Each individual asset's data is stored as a bcolz table with a column for
each pricing field: (open, high, low, close, volume) each pricing field: (open, high, low, close, volume)
The open, high, low, and close columns are integers which are 1000 times The open, high, low, close and volume columns are integers which are 10^8 times
the quoted price, so that the data can represented and stored as an the quoted price, so that the data can represented and stored as an
np.uint32, supporting market prices quoted up to the thousands place. np.uint64, supporting market prices quoted up to the 1/10^8-th place.
volume is a np.uint32 with no mutation of the tens place.
The 'index' for each individual asset are a repeating period of minutes of The 'index' for each individual asset are a repeating period of minutes of
length `minutes_per_day` starting from each market open. length `minutes_per_day` starting from each market open.
@@ -573,7 +573,7 @@ class BcolzMinuteBarWriter(object):
if not os.path.exists(sid_containing_dirname): if not os.path.exists(sid_containing_dirname):
# Other sids may have already created the containing directory. # Other sids may have already created the containing directory.
os.makedirs(sid_containing_dirname) os.makedirs(sid_containing_dirname)
initial_array = np.empty(0, np.uint32) initial_array = np.empty(0, np.uint64)
table = ctable( table = ctable(
rootdir=path, rootdir=path,
columns=[ columns=[
@@ -610,7 +610,7 @@ class BcolzMinuteBarWriter(object):
minute_offset = len(table) % self._minutes_per_day minute_offset = len(table) % self._minutes_per_day
num_to_prepend = numdays * self._minutes_per_day - minute_offset num_to_prepend = numdays * self._minutes_per_day - minute_offset
prepend_array = np.zeros(num_to_prepend, np.uint32) prepend_array = np.zeros(num_to_prepend, np.uint64)
# Fill all OHLCV with zeros. # Fill all OHLCV with zeros.
table.append([prepend_array] * 5) table.append([prepend_array] * 5)
table.flush() table.flush()
@@ -815,11 +815,11 @@ class BcolzMinuteBarWriter(object):
minutes_count = all_minutes_in_window.size minutes_count = all_minutes_in_window.size
open_col = np.zeros(minutes_count, dtype=np.uint32) open_col = np.zeros(minutes_count, dtype=np.uint64)
high_col = np.zeros(minutes_count, dtype=np.uint32) high_col = np.zeros(minutes_count, dtype=np.uint64)
low_col = np.zeros(minutes_count, dtype=np.uint32) low_col = np.zeros(minutes_count, dtype=np.uint64)
close_col = np.zeros(minutes_count, dtype=np.uint32) close_col = np.zeros(minutes_count, dtype=np.uint64)
vol_col = np.zeros(minutes_count, dtype=np.uint32) vol_col = np.zeros(minutes_count, dtype=np.uint64)
dt_ixs = np.searchsorted(all_minutes_in_window.values, dt_ixs = np.searchsorted(all_minutes_in_window.values,
dts.astype('datetime64[ns]')) dts.astype('datetime64[ns]'))
@@ -1125,8 +1125,8 @@ class BcolzMinuteBarReader(MinuteBarReader):
else: else:
return np.nan return np.nan
if field != 'volume': #if field != 'volume':
value *= self._ohlc_ratio_inverse_for_sid(sid) value *= self._ohlc_ratio_inverse_for_sid(sid)
return value return value
def get_last_traded_dt(self, asset, dt): def get_last_traded_dt(self, asset, dt):
@@ -1248,7 +1248,7 @@ class BcolzMinuteBarReader(MinuteBarReader):
if field != 'volume': if field != 'volume':
out = np.full(shape, np.nan) out = np.full(shape, np.nan)
else: else:
out = np.zeros(shape, dtype=np.uint32) out = np.zeros(shape, dtype=np.float64)
for i, sid in enumerate(sids): for i, sid in enumerate(sids):
carray = self._open_minute_file(field, sid) carray = self._open_minute_file(field, sid)
@@ -1262,11 +1262,11 @@ class BcolzMinuteBarReader(MinuteBarReader):
where = values != 0 where = values != 0
# first slice down to len(where) because we might not have # first slice down to len(where) because we might not have
# written data for all the minutes requested # written data for all the minutes requested
if field != 'volume': #if field != 'volume':
out[:len(where), i][where] = ( out[:len(where), i][where] = (
values[where] * self._ohlc_ratio_inverse_for_sid(sid)) values[where] * self._ohlc_ratio_inverse_for_sid(sid))
else: #else:
out[:len(where), i][where] = values[where] # out[:len(where), i][where] = values[where]
results.append(out) results.append(out)
return results return results
+4 -1
View File
@@ -156,7 +156,10 @@ class DailyHistoryAggregator(object):
cache = self._caches[field] = (session, market_open, {}) cache = self._caches[field] = (session, market_open, {})
_, market_open, entries = cache _, market_open, entries = cache
market_open = market_open.tz_localize('UTC') try:
market_open = market_open.tz_localize('UTC')
except TypeError:
market_open = market_open.tz_convert('UTC')
if dt != market_open: if dt != market_open:
prev_dt = dt_value - self._one_min prev_dt = dt_value - self._one_min
else: else:
+12 -13
View File
@@ -11,6 +11,9 @@
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and # See the License for the specific language governing permissions and
# limitations under the License. # limitations under the License.
from __future__ import division # Python2 req to have division of ints yield float
from errno import ENOENT from errno import ENOENT
from functools import partial from functools import partial
from os import remove from os import remove
@@ -80,7 +83,6 @@ from catalyst.utils.cli import (
from ._equities import _compute_row_slices, _read_bcolz_data from ._equities import _compute_row_slices, _read_bcolz_data
from ._adjustments import load_adjustments_from_sqlite from ._adjustments import load_adjustments_from_sqlite
logger = logbook.Logger('UsEquityPricing') logger = logbook.Logger('UsEquityPricing')
OHLC = frozenset(['open', 'high', 'low', 'close']) OHLC = frozenset(['open', 'high', 'low', 'close'])
@@ -116,6 +118,8 @@ SQLITE_STOCK_DIVIDEND_PAYOUT_COLUMN_DTYPES = {
UINT32_MAX = iinfo(uint32).max UINT32_MAX = iinfo(uint32).max
UINT64_MAX = iinfo(uint64).max UINT64_MAX = iinfo(uint64).max
PRICE_ADJUSTMENT_FACTOR = 1000000000 # Provides 9 decimals resolution. Also affects _equities.pyx L220
def check_uint32_safe(value, colname): def check_uint32_safe(value, colname):
if value >= UINT32_MAX: if value >= UINT32_MAX:
@@ -433,11 +437,11 @@ class BcolzDailyBarWriter(object):
return raw_data return raw_data
winsorise_uint64(raw_data, invalid_data_behavior, 'volume', *OHLC) winsorise_uint64(raw_data, invalid_data_behavior, 'volume', *OHLC)
processed = (raw_data[list(OHLC)] * 1000000).astype('uint64') processed = (raw_data[list(OHLC)] * PRICE_ADJUSTMENT_FACTOR).astype('uint64')
dates = raw_data.index.values.astype('datetime64[s]') dates = raw_data.index.values.astype('datetime64[s]')
check_uint32_safe(dates.max().view(np.int64), 'day') check_uint32_safe(dates.max().view(np.int64), 'day')
processed['day'] = dates.astype('uint32') processed['day'] = dates.astype('uint32')
processed['volume'] = raw_data.volume.astype('uint64') processed['volume'] = (raw_data.volume * PRICE_ADJUSTMENT_FACTOR).astype('uint64')
return ctable.fromdataframe(processed) return ctable.fromdataframe(processed)
@@ -490,9 +494,8 @@ class BcolzDailyBarReader(SessionBarReader):
The data in these columns is interpreted as follows: The data in these columns is interpreted as follows:
- Price columns ('open', 'high', 'low', 'close') are interpreted as 1000 * - Price columns ('open', 'high', 'low', 'close') and Volume are interpreted
as-traded dollar value. as 10^9 * as-traded dollar value.
- Volume is interpreted as as-traded volume.
- Day is interpreted as seconds since midnight UTC, Jan 1, 1970. - Day is interpreted as seconds since midnight UTC, Jan 1, 1970.
- Id is the asset id of the row. - Id is the asset id of the row.
@@ -519,7 +522,6 @@ class BcolzDailyBarReader(SessionBarReader):
# Need to test keeping the entire array in memory for the course of a # Need to test keeping the entire array in memory for the course of a
# process first. # process first.
self._spot_cols = {} self._spot_cols = {}
self.PRICE_ADJUSTMENT_FACTOR = 0.001
self._read_all_threshold = read_all_threshold self._read_all_threshold = read_all_threshold
@lazyval @lazyval
@@ -759,13 +761,10 @@ class BcolzDailyBarReader(SessionBarReader):
""" """
ix = self.sid_day_index(sid, dt) ix = self.sid_day_index(sid, dt)
price = self._spot_col(field)[ix] price = self._spot_col(field)[ix]
if field != 'volume': if field != 'volume' and price == 0:
if price == 0: return nan
return nan
else:
return price * 0.001
else: else:
return price return price / PRICE_ADJUSTMENT_FACTOR
class PanelBarReader(SessionBarReader): class PanelBarReader(SessionBarReader):
+91 -10
View File
@@ -19,6 +19,7 @@ from time import sleep
from os import listdir from os import listdir
from os.path import isfile, join from os.path import isfile, join
from collections import deque from collections import deque
import numpy as np
import logbook import logbook
import pandas as pd import pandas as pd
@@ -28,14 +29,16 @@ from catalyst.algorithm import TradingAlgorithm
from catalyst.data.minute_bars import BcolzMinuteBarWriter, \ from catalyst.data.minute_bars import BcolzMinuteBarWriter, \
BcolzMinuteBarReader BcolzMinuteBarReader
from catalyst.errors import OrderInBeforeTradingStart from catalyst.errors import OrderInBeforeTradingStart
from catalyst.exchange.exchange_clock import ExchangeClock from catalyst.exchange.simple_clock import SimpleClock
from catalyst.exchange.live_graph_clock import LiveGraphClock
from catalyst.exchange.exchange_errors import ( from catalyst.exchange.exchange_errors import (
ExchangeRequestError, ExchangeRequestError,
ExchangePortfolioDataError, ExchangePortfolioDataError,
ExchangeTransactionError ExchangeTransactionError
) )
from catalyst.exchange.exchange_utils import get_exchange_minute_writer_root, \ from catalyst.exchange.exchange_utils import get_exchange_minute_writer_root, \
save_algo_object, get_algo_object, get_algo_folder save_algo_object, get_algo_object, get_algo_folder, get_algo_df, \
save_algo_df
from catalyst.exchange.stats_utils import get_pretty_stats from catalyst.exchange.stats_utils import get_pretty_stats
from catalyst.finance.performance.period import calc_period_stats from catalyst.finance.performance.period import calc_period_stats
from catalyst.gens.tradesimulation import AlgorithmSimulator from catalyst.gens.tradesimulation import AlgorithmSimulator
@@ -56,8 +59,19 @@ class ExchangeTradingAlgorithm(TradingAlgorithm):
def __init__(self, *args, **kwargs): def __init__(self, *args, **kwargs):
self.exchange = kwargs.pop('exchange', None) self.exchange = kwargs.pop('exchange', None)
self.algo_namespace = kwargs.pop('algo_namespace', None) self.algo_namespace = kwargs.pop('algo_namespace', None)
self.orders = {} self.live_graph = kwargs.pop('live_graph', None)
self._clock = None
self.minute_stats = deque(maxlen=60) self.minute_stats = deque(maxlen=60)
self.pnl_stats = get_algo_df(self.algo_namespace, 'pnl_stats')
self.custom_signals_stats = \
get_algo_df(self.algo_namespace, 'custom_signals_stats')
self.exposure_stats = \
get_algo_df(self.algo_namespace, 'exposure_stats')
self.is_running = True self.is_running = True
self.retry_check_open_orders = 5 self.retry_check_open_orders = 5
@@ -122,6 +136,13 @@ class ExchangeTradingAlgorithm(TradingAlgorithm):
sys.exit(0) sys.exit(0)
@property
def clock(self):
if self._clock is None:
return self._create_clock()
else:
return self._clock
def _create_clock(self): def _create_clock(self):
# The calendar's execution times are the minutes over which we actually # The calendar's execution times are the minutes over which we actually
@@ -137,10 +158,21 @@ class ExchangeTradingAlgorithm(TradingAlgorithm):
# This method is taken from TradingAlgorithm. # This method is taken from TradingAlgorithm.
# The clock has been replaced to use RealtimeClock # The clock has been replaced to use RealtimeClock
# TODO: should we apply a time skew? not sure to understand the utility. # TODO: should we apply a time skew? not sure to understand the utility.
return ExchangeClock(
self.sim_params.sessions, log.debug('creating clock')
time_skew=self.exchange.time_skew if self.live_graph:
) self._clock = LiveGraphClock(
self.sim_params.sessions,
time_skew=self.exchange.time_skew,
context=self
)
else:
self._clock = SimpleClock(
self.sim_params.sessions,
time_skew=self.exchange.time_skew
)
return self._clock
def _create_generator(self, sim_params): def _create_generator(self, sim_params):
if self.perf_tracker is None: if self.perf_tracker is None:
@@ -156,7 +188,7 @@ class ExchangeTradingAlgorithm(TradingAlgorithm):
self, self,
sim_params, sim_params,
self.data_portal, self.data_portal,
self._create_clock(), self.clock,
self._create_benchmark_source(), self._create_benchmark_source(),
self.restrictions, self.restrictions,
universe_func=self._calculate_universe universe_func=self._calculate_universe
@@ -222,6 +254,49 @@ class ExchangeTradingAlgorithm(TradingAlgorithm):
error=e error=e
) )
def add_pnl_stats(self, period_stats):
starting = period_stats['starting_cash']
current = period_stats['portfolio_value']
appreciation = (current / starting) - 1
perc = (appreciation * 100) if current != 0 else 0
log.debug('adding pnl stats: {:6f}%'.format(perc))
df = pd.DataFrame(
data=[dict(performance=perc)],
index=[period_stats['period_close']]
)
self.pnl_stats = pd.concat([self.pnl_stats, df])
save_algo_df(self.algo_namespace, 'pnl_stats', self.pnl_stats)
def add_custom_signals_stats(self, period_stats):
log.debug('adding custom signals stats: {}'.format(self.recorded_vars))
df = pd.DataFrame(
data=[self.recorded_vars],
index=[period_stats['period_close']],
)
self.custom_signals_stats = pd.concat([self.custom_signals_stats, df])
save_algo_df(self.algo_namespace, 'custom_signals_stats',
self.custom_signals_stats)
def add_exposure_stats(self, period_stats):
data = dict(
long_exposure=period_stats['long_exposure'],
base_currency=period_stats['ending_cash']
)
log.debug('adding exposure stats: {}'.format(data))
df = pd.DataFrame(
data=[data],
index=[period_stats['period_close']],
)
self.exposure_stats = pd.concat([self.exposure_stats, df])
save_algo_df(self.algo_namespace, 'exposure_stats',
self.exposure_stats)
def prepare_period_stats(self, start_dt, end_dt): def prepare_period_stats(self, start_dt, end_dt):
""" """
Creates a dictionary representing the state of the tracker. Creates a dictionary representing the state of the tracker.
@@ -314,9 +389,14 @@ class ExchangeTradingAlgorithm(TradingAlgorithm):
minute_stats = self.prepare_period_stats( minute_stats = self.prepare_period_stats(
data.current_dt, data.current_dt + timedelta(minutes=1)) data.current_dt, data.current_dt + timedelta(minutes=1))
# Saving the last hour in memory # Saving the last hour in memory
self.minute_stats.append(minute_stats) self.minute_stats.append(minute_stats)
self.add_pnl_stats(minute_stats)
self.add_custom_signals_stats(minute_stats)
self.add_exposure_stats(minute_stats)
print_df = pd.DataFrame(list(self.minute_stats)) print_df = pd.DataFrame(list(self.minute_stats))
log.debug( log.debug(
'statistics for the last {stats_minutes} minutes:\n{stats}'.format( 'statistics for the last {stats_minutes} minutes:\n{stats}'.format(
@@ -401,8 +481,9 @@ class ExchangeTradingAlgorithm(TradingAlgorithm):
if order_id is not None: if order_id is not None:
order = self.portfolio.open_orders[order_id] order = self.portfolio.open_orders[order_id]
self.perf_tracker.process_order(order) self.perf_tracker.process_order(order)
return order
return order else:
return None
def round_order(self, amount): def round_order(self, amount):
""" """
+12 -1
View File
@@ -119,7 +119,18 @@ class Bittrex(Exchange):
) )
return order return order
else: else:
raise CreateOrderError(exchange=self.name, error=order_status) if order_status == 'INSUFFICIENT_FUNDS':
log.warn('not enough funds to create order')
return None
elif order_status == 'DUST_TRADE_DISALLOWED_MIN_VALUE_50K_SAT':
log.warn('Your order is too small, order at least 50K'
' Satoshi')
return None
else:
raise CreateOrderError(
exchange=self.name,
error=order_status
)
else: else:
raise InvalidOrderStyle(exchange=self.name, raise InvalidOrderStyle(exchange=self.name,
style=style.__class__.__name__) style=style.__class__.__name__)
+5 -4
View File
@@ -531,10 +531,11 @@ class Exchange:
) )
) )
order = self.create_order(asset, amount, is_buy, style) order = self.create_order(asset, amount, is_buy, style)
if order:
self._portfolio.create_order(order) self._portfolio.create_order(order)
return order.id
return order.id else:
return None
@abstractmethod @abstractmethod
def get_open_orders(self, asset): def get_open_orders(self, asset):
+32
View File
@@ -3,6 +3,7 @@ import os
import pickle import pickle
import urllib import urllib
from datetime import date, datetime from datetime import date, datetime
import pandas as pd
from catalyst.exchange.exchange_errors import ExchangeAuthNotFound, \ from catalyst.exchange.exchange_errors import ExchangeAuthNotFound, \
ExchangeSymbolsNotFound ExchangeSymbolsNotFound
@@ -117,6 +118,37 @@ def append_algo_object(algo_name, key, obj, environ=None):
pickle.dump(obj, handle, protocol=pickle.HIGHEST_PROTOCOL) pickle.dump(obj, handle, protocol=pickle.HIGHEST_PROTOCOL)
def get_algo_df(algo_name, key, environ=None, rel_path=None):
folder = get_algo_folder(algo_name, environ)
if rel_path is not None:
folder = os.path.join(folder, rel_path)
filename = os.path.join(folder, key + '.csv')
if os.path.isfile(filename):
try:
with open(filename, 'rb') as handle:
return pd.read_csv(handle, index_col=0, parse_dates=True)
except IOError:
return pd.DataFrame()
else:
return pd.DataFrame()
def save_algo_df(algo_name, key, df, environ=None, rel_path=None):
folder = get_algo_folder(algo_name, environ)
if rel_path is not None:
folder = os.path.join(folder, rel_path)
ensure_directory(folder)
filename = os.path.join(folder, key + '.csv')
with open(filename, 'wb') as handle:
df.to_csv(handle)
def get_exchange_minute_writer_root(exchange_name, environ=None): def get_exchange_minute_writer_root(exchange_name, environ=None):
exchange_folder = get_exchange_folder(exchange_name, environ) exchange_folder = get_exchange_folder(exchange_name, environ)
+210
View File
@@ -0,0 +1,210 @@
#
# 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.
from datetime import timedelta
import matplotlib.dates as mdates
import pandas as pd
from catalyst.gens.sim_engine import (
BAR,
SESSION_START
)
from logbook import Logger
from matplotlib import pyplot as plt
from matplotlib import style
log = Logger('LiveGraphClock')
fmt = mdates.DateFormatter('%Y-%m-%d %H:%M')
class LiveGraphClock(object):
"""Realtime clock for live trading.
This class is a drop-in replacement for
:class:`zipline.gens.sim_engine.MinuteSimulationClock`.
This mixes the clock with a live graph.
Note
----
This seemingly awkward approach allows us to run the program using a single
thread. This is important because Matplotlib does not play nice with
multi-threaded environments. Zipline probably does not either.
Matplotlib has a pause() method which is a wrapper around time.sleep()
used in the SimpleClock. The key difference is that users
can still interact with the chart during the pause cycles. This is
what enables us to keep a single thread. This is also why we are not using
the 'animate' callback of Matplotlib. We need to direct access to the
__iter__ method in order to yield events to Zipline.
The :param:`time_skew` parameter represents the time difference between
the exchange and the live trading machine's clock. It's not used currently.
"""
def __init__(self, sessions, context, time_skew=pd.Timedelta('0s')):
self.sessions = sessions
self.time_skew = time_skew
self._last_emit = None
self._before_trading_start_bar_yielded = True
self.context = context
style.use('dark_background')
fig = plt.figure()
fig.canvas.set_window_title('Enigma Catalyst: {}'.format(
self.context.algo_namespace))
self.ax_pnl = fig.add_subplot(311)
self.ax_custom_signals = fig.add_subplot(312, sharex=self.ax_pnl)
self.ax_exposure = fig.add_subplot(313, sharex=self.ax_pnl)
if len(context.minute_stats) > 0:
self.draw_pnl()
self.draw_custom_signals()
self.draw_exposure()
# rotates and right aligns the x labels, and moves the bottom of the
# axes up to make room for them
fig.autofmt_xdate()
fig.subplots_adjust(hspace=0.5)
plt.tight_layout()
plt.ion()
plt.show()
def format_ax(self, ax):
"""
Trying to assign reasonable parameters to the time axis.
TODO: room for improvement
:param ax:
:return:
"""
ax.xaxis.set_major_locator(mdates.DayLocator(interval=1))
ax.xaxis.set_major_formatter(fmt)
locator = mdates.HourLocator(interval=4)
locator.MAXTICKS = 5000
ax.xaxis.set_minor_locator(locator)
datemin = pd.Timestamp.utcnow()
ax.set_xlim(datemin)
ax.grid(True)
def set_legend(self, ax):
ax.legend(loc='upper left', ncol=1, fontsize=10, numpoints=1)
def draw_pnl(self):
ax = self.ax_pnl
df = self.context.pnl_stats
ax.clear()
ax.set_title('Performance')
ax.plot(df.index, df['performance'], '-',
color='green',
linewidth=1.0,
label='Performance'
)
def perc(val):
return '{:2f}'.format(val)
ax.format_ydata = perc
self.set_legend(ax)
self.format_ax(ax)
def draw_custom_signals(self):
ax = self.ax_custom_signals
df = self.context.custom_signals_stats
colors = ['blue', 'green', 'red', 'black', 'orange', 'yellow', 'pink']
ax.clear()
ax.set_title('Custom Signals')
for index, column in enumerate(df.columns.values.tolist()):
ax.plot(df.index, df[column], '-',
color=colors[index],
linewidth=1.0,
label=column
)
self.set_legend(ax)
self.format_ax(ax)
def draw_exposure(self):
ax = self.ax_exposure
context = self.context
df = context.exposure_stats
ax.clear()
ax.set_title('Exposure')
ax.plot(df.index, df['base_currency'], '-',
color='green',
linewidth=1.0,
label='Base Currency: {}'.format(
context.exchange.base_currency.upper()
)
)
positions = context.exchange.portfolio.positions
symbols = []
for position in positions:
symbols.append(position.symbol)
ax.plot(df.index, df['long_exposure'], '-',
color='blue',
linewidth=1.0,
label='Long Exposure: {}'.format(
', '.join(symbols).upper()
)
)
self.set_legend(ax)
self.format_ax(ax)
def __iter__(self):
yield pd.Timestamp.utcnow(), SESSION_START
while True:
current_time = pd.Timestamp.utcnow()
current_minute = current_time.floor('1 min')
if self._last_emit is None or current_minute > self._last_emit:
log.debug('emitting minutely bar: {}'.format(current_minute))
self._last_emit = current_minute
yield current_minute, BAR
try:
self.draw_pnl()
self.draw_custom_signals()
self.draw_exposure()
plt.draw()
except Exception as e:
log.warn('Unable to update the graph: {}'.format(e))
else:
# I can't use the "animate" reactive approach here because
# I need to yield from the main loop.
# Workaround: https://stackoverflow.com/a/33050617/814633
plt.pause(1)
@@ -25,7 +25,7 @@ from logbook import Logger
log = Logger('ExchangeClock') log = Logger('ExchangeClock')
class ExchangeClock(object): class SimpleClock(object):
"""Realtime clock for live trading. """Realtime clock for live trading.
This class is a drop-in replacement for This class is a drop-in replacement for
+6 -3
View File
@@ -41,6 +41,7 @@ DEFAULT_EQUITY_VOLUME_SLIPPAGE_BAR_LIMIT = 0.025
DEFAULT_FUTURE_VOLUME_SLIPPAGE_BAR_LIMIT = 0.05 DEFAULT_FUTURE_VOLUME_SLIPPAGE_BAR_LIMIT = 0.05
class LiquidityExceeded(Exception): class LiquidityExceeded(Exception):
pass pass
@@ -205,20 +206,22 @@ class VolumeShareSlippage(SlippageModel):
def process_order(self, data, order): def process_order(self, data, order):
volume = data.current(order.asset, "volume") volume = data.current(order.asset, "volume")
min_trade_size = order.asset.min_trade_size
max_volume = self.volume_limit * volume max_volume = self.volume_limit * volume
# price impact accounts for the total volume of transactions # price impact accounts for the total volume of transactions
# created against the current minute bar # created against the current minute bar
remaining_volume = max_volume - self.volume_for_bar remaining_volume = max_volume - self.volume_for_bar
if remaining_volume < 1: if remaining_volume < min_trade_size:
# we can't fill any more transactions # we can't fill any more transactions
raise LiquidityExceeded() raise LiquidityExceeded()
# the current order amount will be the min of the # the current order amount will be the min of the
# volume available in the bar or the open amount. # volume available in the bar or the open amount.
cur_volume = int(min(remaining_volume, abs(order.open_amount))) cur_volume = min(remaining_volume, abs(order.open_amount))
if cur_volume < 1: if cur_volume < min_trade_size:
return None, None return None, None
# tally the current amount into our total amount ordered. # tally the current amount into our total amount ordered.
+1 -5
View File
@@ -65,14 +65,10 @@ def create_transaction(order, dt, price, amount):
# floor the amount to protect against non-whole number orders # floor the amount to protect against non-whole number orders
# TODO: Investigate whether we can add a robust check in blotter # TODO: Investigate whether we can add a robust check in blotter
# and/or tradesimulation, as well. # and/or tradesimulation, as well.
amount_magnitude = int(abs(amount))
if amount_magnitude < 1:
raise Exception("Transaction magnitude must be at least 1.")
transaction = Transaction( transaction = Transaction(
asset=order.asset, asset=order.asset,
amount=int(amount), amount=amount,
dt=dt, dt=dt,
price=price, price=price,
order_id=order.id order_id=order.id
+2
View File
@@ -17,6 +17,8 @@ import math
from numpy import isnan from numpy import isnan
def round_nearest(x, a):
return round(round(x / a) * a, -int(math.floor(math.log10(a))))
def tolerant_equals(a, b, atol=10e-7, rtol=10e-7, equal_nan=False): def tolerant_equals(a, b, atol=10e-7, rtol=10e-7, equal_nan=False):
"""Check if a and b are equal with some tolerance. """Check if a and b are equal with some tolerance.
+9 -5
View File
@@ -95,7 +95,8 @@ def _run(handle_data,
live, live,
exchange, exchange,
algo_namespace, algo_namespace,
base_currency): base_currency,
live_graph):
"""Run a backtest for the given algorithm. """Run a backtest for the given algorithm.
This is shared between the cli and :func:`catalyst.run_algo`. This is shared between the cli and :func:`catalyst.run_algo`.
@@ -277,7 +278,8 @@ def _run(handle_data,
) )
env = TradingEnvironment( env = TradingEnvironment(
load=partial(load_crypto_market_data, bundle=b, bundle_data=bundle_data, environ=environ), load=partial(load_crypto_market_data, bundle=b,
bundle_data=bundle_data, environ=environ),
bm_symbol='USDT_BTC', bm_symbol='USDT_BTC',
trading_calendar=open_calendar, trading_calendar=open_calendar,
asset_db_path=connstr, asset_db_path=connstr,
@@ -338,7 +340,7 @@ def _run(handle_data,
TradingAlgorithmClass = ( TradingAlgorithmClass = (
partial(ExchangeTradingAlgorithm, exchange=exchange, partial(ExchangeTradingAlgorithm, exchange=exchange,
algo_namespace=algo_namespace) algo_namespace=algo_namespace, live_graph=live_graph)
if live and exchange else TradingAlgorithm) if live and exchange else TradingAlgorithm)
perf = TradingAlgorithmClass( perf = TradingAlgorithmClass(
@@ -439,7 +441,8 @@ def run_algorithm(initialize,
live=False, live=False,
exchange_name=None, exchange_name=None,
base_currency=None, base_currency=None,
algo_namespace=None): algo_namespace=None,
live_graph=False):
"""Run a trading algorithm. """Run a trading algorithm.
Parameters Parameters
@@ -552,5 +555,6 @@ def run_algorithm(initialize,
live=live, live=live,
exchange=exchange_name, exchange=exchange_name,
algo_namespace=algo_namespace, algo_namespace=algo_namespace,
base_currency=base_currency base_currency=base_currency,
live_graph=live_graph
) )
+84
View File
@@ -0,0 +1,84 @@
name: catalyst
channels:
- statiskit
- defaults
dependencies:
- certifi=2016.2.28=py27_0
- coverage=4.4.1=py27_0
- nose=1.3.7=py27_1
- openssl=1.0.2l=0
- path.py=10.3.1=py27_0
- pip=9.0.1=py27_1
- python=2.7.13=0
- pyyaml=3.12=py27_0
- readline=6.2=2
- setuptools=36.4.0=py27_0
- six=1.10.0=py27_0
- sqlite=3.13.0=0
- tk=8.5.18=0
- wheel=0.29.0=py27_0
- yaml=0.1.6=0
- zlib=1.2.11=0
- libdev=1.0.0=py27_0
- python-dev=1.0.0=py27_0
- python-scons=3.0.0=py27_0
- pip:
- alembic==0.9.5
- backports.shutil-get-terminal-size==1.0.0
- bcolz==0.12.1
- bottleneck==1.2.1
- chardet==3.0.4
- click==6.7
- contextlib2==0.5.5
- cycler==0.10.0
- cyordereddict==1.0.0
- cython==0.26.1
- decorator==4.1.2
- empyrical==0.2.1
- enigma-catalyst>=0.2.dev2
- enum34==1.1.6
- functools32==3.2.3.post2
- idna==2.6
- intervaltree==2.1.0
- ipdb==0.10.3
- ipdbplugin==1.4.5
- ipython==5.5.0
- ipython-genutils==0.2.0
- logbook==1.1.0
- lru-dict==1.1.6
- mako==1.0.7
- markupsafe==1.0
- matplotlib==2.0.2
- multipledispatch==0.4.9
- networkx==1.11
- numexpr==2.6.4
- numpy==1.13.1
- pandas==0.19.2
- pandas-datareader==0.5.0
- pathlib2==2.3.0
- patsy==0.4.1
- pexpect==4.2.1
- pickleshare==0.7.4
- prompt-toolkit==1.0.15
- ptyprocess==0.5.2
- pygments==2.2.0
- pyparsing==2.2.0
- python-dateutil==2.6.1
- python-editor==1.0.3
- pytz==2017.2
- requests==2.18.4
- requests-file==1.4.2
- requests-ftp==0.3.1
- scandir==1.5
- scipy==0.19.1
- scons==3.0.0a20170821
- simplegeneric==0.8.1
- sortedcontainers==1.5.7
- sqlalchemy==1.1.14
- statsmodels==0.8.0
- subprocess32==3.2.7
- tables==3.4.2
- toolz==0.8.2
- traitlets==4.3.2
- urllib3==1.22
- wcwidth==0.1.7
+1 -1
View File
@@ -1,7 +1,7 @@
# Incompatible with earlier PIP versions # Incompatible with earlier PIP versions
pip>=7.1.0 pip>=7.1.0
# bcolz fails to install if this is not in the build_requires. # bcolz fails to install if this is not in the build_requires.
setuptools>18.0 setuptools>36.0
# Logging # Logging
Logbook==0.12.5 Logbook==0.12.5