Compare commits

...
50 Commits
Author SHA1 Message Date
Victor Grau Serrat 3d88d6a2c7 Merge branch 'develop' - Release 0.3.3 2017-10-26 13:13:42 -06:00
Victor Grau Serrat 9c3a9e233b Merge branch 'develop' of github.com:enigmampc/catalyst into develop 2017-10-26 12:55:32 -06:00
Victor Grau Serrat c43509c28e catching missing -x in ingest-exchange 2017-10-26 12:55:18 -06:00
fredfortier 0e0bfc82b5 Fixed issues in the prepare_chunk logic 2017-10-26 14:04:30 -04:00
fredfortier 2f660db511 Fixed an issue with daily chunks end date 2017-10-26 13:52:37 -04:00
fredfortier fdc5a30060 Added data validation unit tests and minor fixes to the get_candles method of Poloniex. 2017-10-26 02:33:17 -04:00
fredfortier bb1d96ed5d Merge remote-tracking branch 'origin/develop' into develop 2017-10-25 19:44:05 -04:00
fredfortier 59501905ab Poloniex get_candles fix and created a unit test to validate data. 2017-10-25 19:43:57 -04:00
Victor Grau Serrat 2b85732e36 Merge branch 'develop' - Release 0.3.2 2017-10-24 21:59:53 -06:00
Victor Grau Serrat 284c749bb5 Merge branch 'develop' of github.com:enigmampc/catalyst into develop 2017-10-24 21:58:47 -06:00
VictorandGitHub d7f5e73f84 Merge pull request #43 from reinka/develop
[MIG] Migrated buy_and_hodl and buy_low_sell_high to version 0.3 to work with Poloniex exchange
2017-10-24 21:58:22 -06:00
Victor Grau Serrat cde69da173 Merge branch 'develop' of github.com:enigmampc/catalyst into develop 2017-10-24 21:55:25 -06:00
Victor Grau Serrat bcc75f6b00 FIX: Poloniex 1min curator 2017-10-24 21:54:59 -06:00
fredfortier f179381b64 Small python 3 fixes 2017-10-24 23:41:37 -04:00
fredfortier 10ba53b897 Merge remote-tracking branch 'origin/develop' into develop 2017-10-24 20:03:58 -04:00
fredfortier 1cfe3b1bb2 Fixed issues in the prepare_chunk logic 2017-10-24 20:03:50 -04:00
Victor Grau Serrat 268ff9c826 Merge branch 'develop' of github.com:enigmampc/catalyst into develop 2017-10-24 17:44:24 -06:00
Victor Grau Serrat 7eb184d946 exchange unit tests 2017-10-24 17:44:18 -06:00
fredfortier 1cc34a1485 Fixed urllib package for back compatibility 2017-10-24 19:01:30 -04:00
fredfortier aa2f2f3627 Filtered out starting dates before the calendar 2017-10-24 18:26:44 -04:00
fredfortier 7e373e2f9c Removing symbols.json in clean-exchange. 2017-10-24 18:02:58 -04:00
fredfortier 942e6f263c Fixed an issue with the bar reader. 2017-10-24 16:10:33 -04:00
fredfortier 2e6d7d28ba Fixed an issue with the bar reader. 2017-10-24 16:00:56 -04:00
fredfortier 3a823ea457 Python3 adjustments 2017-10-24 15:47:15 -04:00
fredfortier 4daba6cfb4 Added unit test 2017-10-24 15:46:14 -04:00
Victor Grau Serrat fa018e2e0c more bcolz unit tests 2017-10-24 13:42:53 -06:00
Victor Grau Serrat 315d25f7c0 Merge branch 'develop' of github.com:enigmampc/catalyst into develop 2017-10-24 12:32:00 -06:00
fredfortier cc7ffada96 Merge remote-tracking branch 'origin/develop' into develop 2017-10-24 14:23:57 -04:00
fredfortier 5394c1bc91 Fixed an issue with asset date in chunks 2017-10-24 14:23:47 -04:00
Victor Grau Serrat b230b73829 unit test bcolz writer 2017-10-24 11:28:31 -06:00
Victor Grau Serrat 930a68ab4a unit test for Bcolz writer expanded 2017-10-24 10:36:15 -06:00
Victor Grau Serrat 4e833981e4 unit test for Bcolz writer expanded 2017-10-24 09:55:37 -06:00
fredfortier 2ea402ff10 Modified bcolz unit test 2017-10-24 11:39:17 -04:00
Victor Grau Serrat da6b024edc unit test for Bcolz writer 2017-10-24 09:32:24 -06:00
Victor Grau Serrat 565e9a3cea Added param checking and help msg to clean bundle folders 2017-10-23 21:30:48 -06:00
fredfortier 3c10d19a7e Added method to clean bundle folders 2017-10-23 20:53:25 -04:00
fredfortier cf96e047cd Added method to clean bundle folders 2017-10-23 20:49:40 -04:00
fredfortier 6f6a8e1272 Merge remote-tracking branch 'origin/develop' into develop 2017-10-23 20:29:57 -04:00
fredfortier c2a02e7074 Fixed hash method to create sid numbers 2017-10-23 20:29:48 -04:00
Victor Grau Serrat 7d2cf97fbf FIX: Conda install for Windows 2017-10-23 16:02:28 -06:00
Victor Grau Serrat 195469897c FIX: Windows path 2017-10-23 14:43:56 -06:00
fredfortier c7b422d465 Fix to work around empty bundles 2017-10-22 18:14:35 -04:00
Victor Grau Serrat 2dbace37bb Merge branch 'develop' - Release 0.3.1
FIX: bundle start_dt cannot be earlier than asset_start
FIX: prior raise of AuthNotFound, now generates empty auth.json, and raises AuthEmpty when live
FIX: os.path.join to make BUNDLE_NAME_TEMPLATE compatible across OSes
2017-10-21 22:57:47 -06:00
Victor Grau Serrat 2e903fd42c FIX: bundle start_dt, empty auth, bundle_name_template->os.path.join 2017-10-21 22:56:22 -06:00
reinka 47a104b29c [MIG] Migrated to version 0.3 to work with Poloniex exchange. 2017-10-21 11:26:34 +02:00
fredfortier d248581523 Fixed an error message 2017-10-21 00:27:05 -04:00
fredfortier 48f6300e08 Optimized imports 2017-10-20 23:18:15 -04:00
VictorandGitHub f7a143cb78 Merge pull request #41 from abnera/patch-1
Fix issues with .yml file and incompatible packages.
2017-10-20 15:46:02 -06:00
Victor Grau Serrat 2f7cd97852 DOC: WIP fix tutorial 2017-10-20 15:37:04 -06:00
Abner Ayala-AcevedoandGitHub 73eca75ed9 Updated conda .yml file to work with enigma 0.3 or above.
Removed unnecessary libraries that were giving issues.
2017-10-20 14:30:06 -07:00
36 changed files with 968 additions and 478 deletions
+40 -2
View File
@@ -38,7 +38,7 @@ except NameError:
'--default-extension/--no-default-extension',
is_flag=True,
default=True,
help="Don't load the default catalyst extension.py file in $ZIPLINE_HOME.",
help="Don't load the default catalyst extension.py file in $CATALYST_HOME.",
)
@click.version_option()
def main(extension, strict_extensions, default_extension):
@@ -495,6 +495,10 @@ def ingest_exchange(exchange_name, data_frequency, start, end,
"""
Ingest data for the given exchange.
"""
if exchange_name is None:
ctx.fail("must specify an exchange name '-x'")
exchange = get_exchange(exchange_name)
exchange_bundle = ExchangeBundle(exchange)
@@ -509,6 +513,40 @@ def ingest_exchange(exchange_name, data_frequency, start, end,
)
@main.command(name='clean-exchange')
@click.option(
'-x',
'--exchange-name',
type=click.Choice({'bitfinex', 'bittrex', 'poloniex'}),
help='The name of the exchange bundle to ingest (supported: bitfinex,'
' bittrex, poloniex).',
)
@click.option(
'-f',
'--data-frequency',
type=click.Choice({'daily', 'minute'}),
default=None,
help='The bundle data frequency to remove. If not specified, it will '
'remove both daily and minute bundles.',
)
@click.pass_context
def clean_exchange(ctx, exchange_name, data_frequency):
"""Clean up bundles from 'ingest-exchange'.
"""
if exchange_name is None:
ctx.fail("must specify an exchange name '-x'")
exchange = get_exchange(exchange_name)
exchange_bundle = ExchangeBundle(exchange)
click.echo('Cleaning exchange bundle {}...'.format(exchange_name))
exchange_bundle.clean(
data_frequency=data_frequency,
)
click.echo('Done')
@main.command()
@click.option(
'-b',
@@ -598,7 +636,7 @@ def ingest(ctx, bundle, exchange_name, compile_locally, assets_version,
' This may not be passed with -e / --before or -a / --after',
)
def clean(bundle, before, after, keep_last):
"""Clean up data downloaded with the ingest command.
"""Clean up bundles from 'ingest'.
"""
bundles_module.clean(
bundle,
+7 -1
View File
@@ -17,6 +17,8 @@
"""
Cythonized Asset object.
"""
import hashlib
cimport cython
from cpython.number cimport PyNumber_Index
from cpython.object cimport (
@@ -501,7 +503,11 @@ cdef class TradingPair(Asset):
if sid == 0 or sid is None:
try:
sid = abs(hash(symbol)) % (10 ** 4)
# sid = abs(hash(symbol)) % (10 ** 4)
# TODO: try to encode the symbol in the main scope
sid = int(
hashlib.sha256(symbol.encode('utf-8')).hexdigest(), 16
) % 10 ** 6
except Exception as e:
raise SidHashError(symbol=symbol)
+26 -26
View File
@@ -212,32 +212,32 @@ class PoloniexCurator(object):
def write_ohlcv_file(self, currencyPair):
csv_trades = CSV_OUT_FOLDER + 'crypto_trades-' + currencyPair + '.csv'
csv_1min = CSV_OUT_FOLDER + 'crypto_1min-' + currencyPair + '.csv'
if( os.path.isfile(csv_1min) ):
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.')
#if( os.path.isfile(csv_1min) ):
# 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, 'w') 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.')
'''
+4 -4
View File
@@ -24,7 +24,7 @@ from catalyst.api import (
)
def initialize(context):
context.ASSET_NAME = 'USDT_BTC'
context.ASSET_NAME = 'BTC_USDT'
context.TARGET_HODL_RATIO = 0.8
context.RESERVE_RATIO = 1.0 - context.TARGET_HODL_RATIO
@@ -49,14 +49,14 @@ def handle_data(context, data):
orders = get_open_orders(context.asset) or []
for order in orders:
cancel_order(order)
# Stop buying after passing the reserve threshold
cash = context.portfolio.cash
if cash <= reserve_value:
context.is_buying = False
# Retrieve current asset price from pricing data
price = data[context.asset].price
price = data.current(context.asset, 'price')
# Check if still buying and could (approximately) afford another purchase
if context.is_buying and cash > price:
@@ -70,7 +70,7 @@ def handle_data(context, data):
record(
price=price,
volume=data[context.asset].volume,
volume=data.current(context.asset, 'volume'),
cash=cash,
starting_cash=context.portfolio.starting_cash,
leverage=context.account.leverage,
+10
View File
@@ -0,0 +1,10 @@
from catalyst.api import order, record, symbol
def initialize(context):
context.asset = symbol('btc_usd')
def handle_data(context, data):
order(context.asset, 1)
record(btc=data.current(context.asset, 'price'))
+1 -1
View File
@@ -27,7 +27,7 @@ log = Logger(algo_namespace)
def initialize(context):
log.info('initializing algo')
context.ASSET_NAME = 'XRP_USD'
context.ASSET_NAME = 'XRP_USDT'
context.asset = symbol(context.ASSET_NAME)
context.TARGET_POSITIONS = 5000
+18 -18
View File
@@ -1,13 +1,13 @@
import pandas as pd
import talib
import pandas as pd
from catalyst import run_algorithm
from catalyst.api import symbol
def initialize(context):
print('initializing')
context.asset = symbol('xrp_btc')
context.asset = symbol('burst_btc')
def handle_data(context, data):
@@ -27,25 +27,25 @@ def handle_data(context, data):
pass
# run_algorithm(
# capital_base=250,
# start=pd.to_datetime('2015-08-01', utc=True),
# end=pd.to_datetime('2017-9-30', utc=True),
# data_frequency='daily',
# initialize=initialize,
# handle_data=handle_data,
# analyze=None,
# exchange_name='poloniex',
# algo_namespace='simple_loop',
# base_currency='eth'
# )
run_algorithm(
capital_base=250,
start=pd.to_datetime('2017-08-01', utc=True),
end=pd.to_datetime('2017-9-30', utc=True),
data_frequency='minute',
initialize=initialize,
handle_data=handle_data,
analyze=None,
exchange_name='bitfinex',
live=True,
exchange_name='poloniex',
algo_namespace='simple_loop',
base_currency='eth',
live_graph=False
base_currency='btc'
)
# run_algorithm(
# initialize=initialize,
# handle_data=handle_data,
# analyze=None,
# exchange_name='bitfinex',
# live=True,
# algo_namespace='simple_loop',
# base_currency='eth',
# live_graph=False
# )
+13 -3
View File
@@ -1,10 +1,10 @@
import base64
import datetime
import hashlib
import hmac
import json
import re
import time
import datetime
import numpy as np
import pandas as pd
@@ -22,10 +22,10 @@ from catalyst.exchange.exchange_errors import (
InvalidOrderStyle, OrderCancelError)
from catalyst.exchange.exchange_execution import ExchangeLimitOrder, \
ExchangeStopLimitOrder, ExchangeStopOrder
from catalyst.exchange.exchange_utils import get_exchange_symbols_filename, \
download_exchange_symbols, get_symbols_string
from catalyst.finance.order import Order, ORDER_STATUS
from catalyst.protocol import Account
from catalyst.exchange.exchange_utils import get_exchange_symbols_filename, \
download_exchange_symbols
# Trying to account for REST api instability
# https://stackoverflow.com/questions/15431044/can-i-set-max-retries-for-requests-request
@@ -255,6 +255,16 @@ class Bitfinex(Exchange):
'1m', '5m', '15m', '30m', '1h', '3h', '6h', '12h', '1D', '7D', '14D',
'1M'
"""
log.debug(
'retrieving {bars} {freq} candles on {exchange} from '
'{end_dt} for markets {symbols}, '.format(
bars=bar_count,
freq=data_frequency,
exchange=self.name,
end_dt=end_dt,
symbols=get_symbols_string(assets)
)
)
freq_match = re.match(r'([0-9].*)(m|h|d)', data_frequency, re.M | re.I)
if freq_match:
+32 -11
View File
@@ -1,22 +1,24 @@
import json
import pandas as pd
import time
from catalyst.assets._assets import TradingPair
from logbook import Logger
from six.moves import urllib
from catalyst.constants import LOG_LEVEL
from catalyst.exchange.bittrex.bittrex_api import Bittrex_api
from catalyst.exchange.exchange import Exchange
from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.exchange_errors import InvalidHistoryFrequencyError, \
ExchangeRequestError, InvalidOrderStyle, OrderNotFound, OrderCancelError, \
CreateOrderError
from catalyst.exchange.exchange_utils import get_exchange_symbols_filename, \
download_exchange_symbols, get_symbols_string
from catalyst.finance.execution import LimitOrder, StopLimitOrder
from catalyst.finance.order import Order, ORDER_STATUS
from catalyst.exchange.exchange_utils import get_exchange_symbols_filename, \
download_exchange_symbols
from catalyst.constants import LOG_LEVEL
# TODO: consider using this: https://github.com/mondeja/bittrex_v2
log = Logger('Bittrex', level=LOG_LEVEL)
@@ -25,7 +27,7 @@ URL2 = 'https://bittrex.com/Api/v2.0'
class Bittrex(Exchange):
def __init__(self, key, secret, base_currency, portfolio=None):
self.api = Bittrex_api(key=key, secret=secret.encode('UTF-8'))
self.api = Bittrex_api(key=key, secret=secret)
self.name = 'bittrex'
self.color = 'blue'
self.base_currency = base_currency
@@ -66,10 +68,10 @@ class Bittrex(Exchange):
return exchange_symbol.lower()
def get_balances(self):
balances = self.api.getbalances()
try:
log.debug('retrieving wallet balances')
self.ask_request()
balances = self.api.getbalances()
except Exception as e:
raise ExchangeRequestError(error=e)
@@ -209,7 +211,7 @@ class Bittrex(Exchange):
)
def get_candles(self, data_frequency, assets, bar_count=None,
start_date=None):
start_dt=None, end_dt=None):
"""
Supported Intervals
-------------------
@@ -218,10 +220,27 @@ class Bittrex(Exchange):
:param data_frequency:
:param assets:
:param bar_count:
:param start_dt
:param end_dt
:return:
"""
log.info('retrieving candles')
# TODO: this has no effect at the moment
if end_dt is None:
end_dt = pd.Timestamp.utcnow()
log.debug(
'retrieving {bars} {freq} candles on {exchange} from '
'{end_dt} for markets {symbols}, '.format(
bars=bar_count,
freq=data_frequency,
exchange=self.name,
end_dt=end_dt,
symbols=get_symbols_string(assets)
)
)
data_frequency = data_frequency.lower()
if data_frequency == 'minute' or data_frequency == '1m':
frequency = 'oneMin'
elif data_frequency == '5m':
@@ -230,7 +249,7 @@ class Bittrex(Exchange):
frequency = 'thirtyMin'
elif data_frequency == '1h':
frequency = 'hour'
elif data_frequency == 'daily' or data_frequency == '1D':
elif data_frequency == 'daily' or data_frequency == '1d':
frequency = 'day'
else:
raise InvalidHistoryFrequencyError(
@@ -239,13 +258,14 @@ class Bittrex(Exchange):
# Making sure that assets are iterable
asset_list = [assets] if isinstance(assets, TradingPair) else assets
ohlc_map = dict()
for asset in asset_list:
end = int(time.mktime(end_dt.timetuple()))
url = '{url}/pub/market/GetTicks?marketName={symbol}' \
'&tickInterval={frequency}&_=1499127220008'.format(
'&tickInterval={frequency}&_={end}'.format(
url=URL2,
symbol=self.get_symbol(asset),
frequency=frequency
frequency=frequency,
end=end
)
try:
@@ -273,6 +293,7 @@ class Bittrex(Exchange):
return ohlc
ordered_candles = list(reversed(candles))
ohlc_map = dict()
if bar_count is None:
ohlc_map[asset] = ohlc_from_candle(ordered_candles[0])
else:
+5 -2
View File
@@ -4,10 +4,10 @@ import time
import hmac
import hashlib
from six.moves import urllib
# Workaround for backwards compatibility
# https://stackoverflow.com/questions/3745771/urllib-request-in-python-2-7
from six.moves import urllib
urlopen = urllib.request.urlopen
@@ -39,7 +39,10 @@ class Bittrex_api(object):
if method not in self.public:
url += '&apikey=' + self.key
url += '&nonce=' + str(int(time.time()))
signature = hmac.new(self.secret, url, hashlib.sha512).hexdigest()
signature = hmac.new(self.secret.encode('utf-8'),
url.encode('utf-8'),
hashlib.sha512).hexdigest()
headers = {'apisign': signature}
else:
headers = {}
+2 -38
View File
@@ -103,42 +103,6 @@ def get_start_dt(end_dt, bar_count, data_frequency):
return start_dt
def get_adj_dates(start, end, assets, data_frequency):
"""
Contains a date range to the trading availability of the specified pairs.
:param start:
:param end:
:param assets:
:param data_frequency:
:return:
"""
earliest_trade = None
last_entry = None
for asset in assets:
if earliest_trade is None or earliest_trade > asset.start_date:
earliest_trade = asset.start_date
end_asset = asset.end_minute if data_frequency == 'minute' else \
asset.end_daily
if end_asset is not None and \
(last_entry is None or end_asset > last_entry):
last_entry = end_asset
if start is None or earliest_trade > start:
start = earliest_trade
if end is None or (last_entry is not None and end > last_entry):
end = last_entry
if end is None or start >= end:
raise NoDataAvailableOnExchange(
exchange=asset.exchange.title(),
symbol=[asset.symbol.encode('utf-8')],
data_frequency=data_frequency,
)
return start, end
def get_month_start_end(dt):
@@ -243,12 +207,12 @@ def find_most_recent_time(bundle_name):
for folder in bundle_folders:
date = from_bundle_ingest_dirname(folder)
if not most_recent_bundle or date > \
most_recent_bundle[most_recent_bundle.keys()[0]]:
most_recent_bundle[list(most_recent_bundle.keys())[0]]:
most_recent_bundle = dict()
most_recent_bundle[folder] = date
if most_recent_bundle:
return most_recent_bundle.keys()[0]
return list(most_recent_bundle.keys())[0]
else:
return None
+5 -9
View File
@@ -19,17 +19,13 @@ import pandas as pd
from catalyst.assets._assets import TradingPair
from logbook import Logger
from catalyst.constants import LOG_LEVEL
from catalyst.data.data_portal import DataPortal
from catalyst.exchange.bundle_utils import get_start_dt
from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.exchange_errors import (
ExchangeRequestError,
ExchangeBarDataError,
PricingDataBeforeTradingError,
PricingDataNotLoadedError, InvalidHistoryFrequencyError,
BundleNotFoundError)
from catalyst.constants import LOG_LEVEL
PricingDataNotLoadedError)
log = Logger('DataPortalExchange', level=LOG_LEVEL)
@@ -84,7 +80,7 @@ class DataPortalExchangeBase(DataPortal):
return pd.concat(df_list)
else:
exchange = self.exchanges[exchange_assets.keys()[0]]
exchange = self.exchanges[list(exchange_assets.keys())[0]]
return self.get_exchange_history_window(
exchange,
assets,
@@ -169,8 +165,8 @@ class DataPortalExchangeBase(DataPortal):
exchange_assets[asset.exchange].append(asset)
if len(exchange_assets.keys()) == 1:
exchange = self.exchanges[exchange_assets.keys()[0]]
if len(list(exchange_assets.keys())) == 1:
exchange = self.exchanges[list(exchange_assets.keys())[0]]
return self.get_exchange_spot_value(
exchange, assets, field, dt, data_frequency)
+18 -12
View File
@@ -9,14 +9,14 @@ import pandas as pd
from catalyst.assets._assets import TradingPair
from logbook import Logger
from catalyst.constants import LOG_LEVEL
from catalyst.data.data_portal import BASE_FIELDS
from catalyst.exchange.bundle_utils import get_start_dt, \
get_delta, get_periods, get_adj_dates
get_delta, get_periods
from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.exchange_errors import MismatchingBaseCurrencies, \
InvalidOrderStyle, BaseCurrencyNotFoundError, SymbolNotFoundOnExchange, \
InvalidHistoryFrequencyError, MismatchingFrequencyError, \
BundleNotFoundError, NoDataAvailableOnExchange, PricingDataNotLoadedError
InvalidHistoryFrequencyError, PricingDataNotLoadedError
from catalyst.exchange.exchange_execution import ExchangeStopLimitOrder, \
ExchangeLimitOrder, ExchangeStopOrder
from catalyst.exchange.exchange_portfolio import ExchangePortfolio
@@ -24,8 +24,6 @@ from catalyst.exchange.exchange_utils import get_exchange_symbols
from catalyst.finance.order import ORDER_STATUS
from catalyst.finance.transaction import Transaction
from catalyst.constants import LOG_LEVEL
log = Logger('Exchange', level=LOG_LEVEL)
@@ -89,7 +87,7 @@ class Exchange:
self.request_cpt[now] = 0
return True
cpt_date = self.request_cpt.keys()[0]
cpt_date = list(self.request_cpt.keys())[0]
cpt = self.request_cpt[cpt_date]
if now > cpt_date + timedelta(minutes=1):
@@ -169,8 +167,10 @@ class Exchange:
asset = self.assets[key]
if not asset:
supported_symbols = [pair.symbol.encode('utf-8') for pair in
self.assets.values()]
supported_symbols = [
pair.symbol for pair in list(self.assets.values())
]
raise SymbolNotFoundOnExchange(
symbol=symbol,
exchange=self.name.title(),
@@ -373,7 +373,7 @@ class Exchange:
return value
def get_series_from_candles(self, candles, start_dt, end_dt,
field, previous_value=None):
data_frequency, field, previous_value=None):
"""
Get a series of field data for the specified candles.
@@ -388,9 +388,12 @@ class Exchange:
dates = [candle['last_traded'] for candle in candles]
values = [candle[field] for candle in candles]
periods = pd.date_range(start_dt, end_dt)
periods = self.bundle.get_calendar_periods_range(
start_dt, end_dt, data_frequency
)
series = pd.Series(values, index=dates)
#TODO: ensure that this working as expected, if not use fillna
series.reindex(periods, method='ffill', fill_value=previous_value)
return series
@@ -487,6 +490,7 @@ class Exchange:
data_frequency=data_frequency,
assets=asset,
bar_count=trailing_bar_count,
start_dt=start_dt,
end_dt=end_dt
)
@@ -497,6 +501,7 @@ class Exchange:
candles=candles,
start_dt=trailing_dt,
end_dt=end_dt,
data_frequency=data_frequency,
field=field,
previous_value=last_value
)
@@ -554,7 +559,7 @@ class Exchange:
portfolio.starting_cash = portfolio.cash
if portfolio.positions:
assets = portfolio.positions.keys()
assets = list(portfolio.positions.keys())
tickers = self.tickers(assets)
portfolio.positions_value = 0.0
@@ -784,13 +789,14 @@ class Exchange:
pass
@abc.abstractmethod
def get_orderbook(self, asset, order_type):
def get_orderbook(self, asset, order_type, limit):
"""
Retrieve the the orderbook for the given trading pair.
:param asset: TradingPair
:param order_type: str
The type of orders: bid, ask or all
:param limit
:return:
"""
+5 -5
View File
@@ -26,6 +26,7 @@ from catalyst.assets._assets import TradingPair
import catalyst.protocol as zp
from catalyst.algorithm import TradingAlgorithm
from catalyst.constants import LOG_LEVEL
from catalyst.data.minute_bars import BcolzMinuteBarWriter, \
BcolzMinuteBarReader
from catalyst.errors import OrderInBeforeTradingStart
@@ -51,10 +52,8 @@ from catalyst.utils.api_support import (
disallowed_in_before_trading_start)
from catalyst.utils.input_validation import error_keywords, ensure_upper_case, \
expect_types
from catalyst.utils.preprocess import preprocess
from catalyst.utils.math_utils import round_nearest
from catalyst.constants import LOG_LEVEL
from catalyst.utils.preprocess import preprocess
log = logbook.Logger('exchange_algorithm', level=LOG_LEVEL)
@@ -114,7 +113,7 @@ class ExchangeTradingAlgorithmBase(TradingAlgorithm):
else self.sim_params.end_session
if exchange_name is None:
exchange = self.exchanges.values()[0]
exchange = list(self.exchanges.values())[0]
else:
exchange = self.exchanges[exchange_name]
@@ -525,7 +524,7 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
self.add_pnl_stats(minute_stats)
if self.recorded_vars:
self.add_custom_signals_stats(minute_stats)
recorded_cols = self.recorded_vars.keys()
recorded_cols = list(self.recorded_vars.keys())
else:
recorded_cols = None
@@ -557,6 +556,7 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
except Exception as e:
log.warn('unable to calculate performance: {}'.format(e))
# TODO: pickle does not seem to work in python 3
try:
save_algo_object(
algo_name=self.algo_namespace,
+3 -3
View File
@@ -3,7 +3,6 @@ import numpy as np
from catalyst import get_calendar
from catalyst.data.minute_bars import BcolzMinuteBarReader, \
BcolzMinuteBarWriter
from catalyst.exchange.bundle_utils import get_periods, get_periods_range
class BcolzExchangeBarWriter(BcolzMinuteBarWriter):
@@ -17,7 +16,7 @@ class BcolzExchangeBarWriter(BcolzMinuteBarWriter):
end_session = end_session.floor('1d')
minutes_per_day = 1440 if self._data_frequency == 'minute' else 1
default_ohlc_ratio = kwargs.pop('default_ohlc_ratio', 1000000)
default_ohlc_ratio = kwargs.pop('default_ohlc_ratio', 100000000)
calendar = get_calendar('OPEN')
super(BcolzExchangeBarWriter, self) \
@@ -80,8 +79,9 @@ class BcolzExchangeBarReader(BcolzMinuteBarReader):
if mask is None:
mask = a != 0
inverse_ratio = self._ohlc_ratio_inverse_for_sid(sid)
out[:len(mask), i][mask] = (
a[mask] * self._ohlc_ratio_inverse_for_sid(sid)
a[mask] * inverse_ratio
)
if field in fields:
+1 -2
View File
@@ -1,13 +1,12 @@
from catalyst.assets._assets import TradingPair
from logbook import Logger
from catalyst.constants import LOG_LEVEL
from catalyst.finance.blotter import Blotter
from catalyst.finance.commission import CommissionModel
from catalyst.finance.slippage import SlippageModel
from catalyst.finance.transaction import Transaction
from catalyst.constants import LOG_LEVEL
log = Logger('exchange_blotter', level=LOG_LEVEL)
# It seems like we need to accept greater slippage risk in cryptos
+243 -105
View File
@@ -3,33 +3,34 @@ import shutil
from datetime import timedelta
import pandas as pd
from logbook import Logger, INFO
from logbook import Logger
from catalyst import get_calendar
from catalyst.constants import LOG_LEVEL
from catalyst.data.minute_bars import BcolzMinuteOverlappingData, \
BcolzMinuteBarMetadata
from catalyst.exchange.bundle_utils import range_in_bundle, \
get_bcolz_chunk, get_delta, get_adj_dates, get_month_start_end, \
get_year_start_end, get_periods_range, get_df_from_arrays, get_start_dt
get_bcolz_chunk, get_delta, get_month_start_end, \
get_year_start_end, get_df_from_arrays, get_start_dt
from catalyst.exchange.exchange_bcolz import BcolzExchangeBarReader, \
BcolzExchangeBarWriter
from catalyst.exchange.exchange_errors import EmptyValuesInBundleError, \
InvalidHistoryFrequencyError, PricingDataBeforeTradingError, \
TempBundleNotFoundError, NoDataAvailableOnExchange, \
InvalidHistoryFrequencyError, TempBundleNotFoundError, \
NoDataAvailableOnExchange, \
PricingDataNotLoadedError
from catalyst.exchange.exchange_utils import get_exchange_folder
from catalyst.utils.cli import maybe_show_progress
from catalyst.utils.paths import ensure_directory
from catalyst.constants import LOG_LEVEL
log = Logger('exchange_bundle', level=LOG_LEVEL)
BUNDLE_NAME_TEMPLATE = '{root}/{frequency}_bundle'
BUNDLE_NAME_TEMPLATE = os.path.join('{root}', '{frequency}_bundle')
def _cachpath(symbol, type_):
return '-'.join([symbol, type_])
class ExchangeBundle:
def __init__(self, exchange):
self.exchange = exchange
@@ -172,16 +173,13 @@ class ExchangeBundle:
invalid_data_behavior='raise'
)
except BcolzMinuteOverlappingData as e:
log.warn('chunk already exists: {}'.format(e))
log.debug('chunk already exists: {}'.format(e))
except Exception as e:
log.warn('error when writing data: {}, trying again'.format(e))
# This is workaround, there is an issue with empty
# session_label when using a newly created writer
key = writer._rootdir if data_frequency == 'minute' \
else writer._filename
del self._writers[key]
del self._writers[writer._rootdir]
writer = self.get_writer(writer._start_session,
writer._end_session, data_frequency)
@@ -196,6 +194,71 @@ class ExchangeBundle:
if data_frequency == 'minute' \
else self.calendar.sessions_in_range(start_dt, end_dt)
def ingest_df(self, ohlcv_df, data_frequency, asset, writer,
empty_rows_behavior='strip'):
"""
Ingest a DataFrame of OHLCV data for a given market.
:param ohlcv_df:
:param data_frequency:
:param asset:
:param writer:
:param path:
:param empty_rows_behavior:
:return:
"""
if empty_rows_behavior is not 'ignore':
nan_rows = ohlcv_df[ohlcv_df.isnull().T.any().T].index
if len(nan_rows) > 0:
dates = []
previous_date = None
for row_date in nan_rows.values:
row_date = pd.to_datetime(row_date)
if previous_date is None:
dates.append(row_date)
else:
seq_date = previous_date + get_delta(1, data_frequency)
if row_date > seq_date:
dates.append(previous_date)
dates.append(row_date)
previous_date = row_date
dates.append(pd.to_datetime(nan_rows.values[-1]))
name = '{} from {} to {}'.format(
asset.symbol, ohlcv_df.index[0], ohlcv_df.index[-1]
)
if empty_rows_behavior == 'warn':
log.warn(
'\n{name} with end minute {end_minute} has empty rows '
'in ranges: {dates}'.format(
name=name,
end_minute=asset.end_minute,
dates=dates
)
)
elif empty_rows_behavior == 'raise':
raise EmptyValuesInBundleError(
name=name,
end_minute=asset.end_minute,
dates=dates
)
else:
ohlcv_df.dropna(inplace=True)
data = []
if not ohlcv_df.empty:
ohlcv_df.sort_index(inplace=True)
data.append((asset.sid, ohlcv_df))
self._write(data, writer, data_frequency)
def ingest_ctable(self, asset, data_frequency, period, start_dt, end_dt,
writer, empty_rows_behavior='strip', cleanup=False):
"""
@@ -225,12 +288,18 @@ class ExchangeBundle:
if reader is None:
raise TempBundleNotFoundError(path=path)
arrays = reader.load_raw_arrays(
sids=[asset.sid],
fields=['open', 'high', 'low', 'close', 'volume'],
start_dt=start_dt,
end_dt=end_dt
)
arrays = None
try:
arrays = reader.load_raw_arrays(
sids=[asset.sid],
fields=['open', 'high', 'low', 'close', 'volume'],
start_dt=start_dt,
end_dt=end_dt
)
except Exception as e:
log.warn('skipping ctable for {} from {} to {}: {}'.format(
asset.symbol, start_dt, end_dt, e
))
if not arrays:
return path
@@ -238,65 +307,69 @@ class ExchangeBundle:
periods = self.get_calendar_periods_range(
start_dt, end_dt, data_frequency
)
df = get_df_from_arrays(arrays, periods)
if empty_rows_behavior is not 'ignore':
nan_rows = df[df.isnull().T.any().T].index
if len(nan_rows) > 0:
dates = []
previous_date = None
for row_date in nan_rows.values:
row_date = pd.to_datetime(row_date)
if previous_date is None:
dates.append(row_date)
else:
seq_date = previous_date + get_delta(1, data_frequency)
if row_date > seq_date:
dates.append(previous_date)
dates.append(row_date)
previous_date = row_date
dates.append(pd.to_datetime(nan_rows.values[-1]))
name = path.split('/')[-1]
if empty_rows_behavior == 'warn':
log.warn(
'\n{name} with end minute {end_minute} has empty rows '
'in ranges: {dates}'.format(
name=name,
end_minute=asset.end_minute,
dates=dates
)
)
elif empty_rows_behavior == 'raise':
raise EmptyValuesInBundleError(
name=name,
end_minute=asset.end_minute,
dates=dates
)
else:
df.dropna(inplace=True)
data = []
if not df.empty:
df.sort_index(inplace=True)
data.append((asset.sid, df))
self._write(data, writer, data_frequency)
self.ingest_df(
ohlcv_df=df,
data_frequency=data_frequency,
asset=asset,
writer=writer,
empty_rows_behavior=empty_rows_behavior
)
if cleanup:
log.debug('removing bundle folder following '
'ingestion: {}'.format(path))
log.debug(
'removing bundle folder following ingestion: {}'.format(path)
)
shutil.rmtree(path)
return path
def get_adj_dates(self, start, end, assets, data_frequency):
"""
Contains a date range to the trading availability of the specified pairs.
:param start:
:param end:
:param assets:
:param data_frequency:
:return:
"""
earliest_trade = None
last_entry = None
for asset in assets:
if earliest_trade is None or earliest_trade > asset.start_date:
if asset.start_date >= self.calendar.first_session:
earliest_trade = asset.start_date
else:
earliest_trade = self.calendar.first_session
end_asset = asset.end_minute if data_frequency == 'minute' else \
asset.end_daily
if end_asset is not None:
if last_entry is None or end_asset > last_entry:
last_entry = end_asset
else:
end = None
last_entry = None
if start is None or \
(earliest_trade is not None and earliest_trade > start):
start = earliest_trade
if end is None or (last_entry is not None and end > last_entry):
end = last_entry
if end is None or start is None or start >= end:
raise NoDataAvailableOnExchange(
exchange=asset.exchange.title(),
symbol=[asset.symbol],
data_frequency=data_frequency,
)
return start, end
def prepare_chunks(self, assets, data_frequency, start_dt, end_dt):
"""
Split a price data request into chunks corresponding to individual
@@ -313,23 +386,27 @@ class ExchangeBundle:
chunks = []
for asset in assets:
try:
asset_start, asset_end = \
get_adj_dates(start_dt, end_dt, [asset], data_frequency)
# Checking if the the asset has price data in the specified
# date range
adj_start, adj_end = self.get_adj_dates(
start_dt, end_dt, [asset], data_frequency
)
except NoDataAvailableOnExchange:
except NoDataAvailableOnExchange as e:
# If not, we continue to the next asset
log.debug('skipping {}: {}'.format(asset.symbol, e))
continue
# This is either the first trading day of the asset or the
# first session available in the calendar
first_trading_dt = asset.start_date \
if asset.start_date > self.calendar.first_session \
else self.calendar.first_session
# Aligning start / end dates with the daily calendar
sessions = get_periods_range(start_dt, end_dt, data_frequency) \
if data_frequency == 'minute' \
else self.calendar.sessions_in_range(start_dt, end_dt)
if asset_start < sessions[0]:
asset_start = sessions[0]
if asset_end > sessions[-1]:
asset_end = sessions[-1]
sessions = self.calendar.sessions_in_range(adj_start, adj_end)
# We loop through each session to create chunks for each period
chunk_labels = []
dt = sessions[0]
while dt <= sessions[-1]:
@@ -343,29 +420,39 @@ class ExchangeBundle:
# of the trading pair
if data_frequency == 'minute':
period_start, period_end = get_month_start_end(dt)
asset_start_month, _ = get_month_start_end(asset_start)
asset_start_month, _ = get_month_start_end(
first_trading_dt
)
if asset_start_month == period_start \
and period_start < asset_start:
period_start = asset_start
and period_start < first_trading_dt:
period_start = first_trading_dt
_, asset_end_month = get_month_start_end(asset_end)
# TODO: need to filter closed pairs?
_, asset_end_month = get_month_start_end(
asset.end_minute
)
if asset_end_month == period_end \
and period_end > asset_end:
period_end = asset_end
and period_end > asset.end_minute:
period_end = asset.end_minute
elif data_frequency == 'daily':
period_start, period_end = get_year_start_end(dt)
asset_start_year, _ = get_year_start_end(asset_start)
asset_start_year, _ = get_year_start_end(
first_trading_dt
)
if asset_start_year == period_start \
and period_start < asset_start:
period_start = asset_start
and period_start < first_trading_dt:
period_start = first_trading_dt
_, asset_end_year = get_year_start_end(asset_end)
_, asset_end_year = get_year_start_end(
asset.end_daily
)
if asset_end_year == period_end \
and period_end > asset_end:
period_end = asset_end
and period_end > asset.end_daily:
period_end = asset.end_daily
else:
raise InvalidHistoryFrequencyError(
frequency=data_frequency
@@ -375,10 +462,13 @@ class ExchangeBundle:
# Checking the last minute of the day instead.
range_start = period_start.replace(hour=23, minute=59) \
if data_frequency == 'minute' else period_start
# Checking if the data already exists in the bundle
# for the date range of the chunk. If not, we create
# a chunk for ingestion.
has_data = range_in_bundle(
asset, range_start, period_end, reader
)
if not has_data:
log.debug('adding period: {}'.format(label))
chunks.append(
@@ -392,6 +482,7 @@ class ExchangeBundle:
dt += timedelta(days=1)
# We sort the chunks by end date to ingest most recent data first
chunks.sort(key=lambda chunk: chunk['period_end'])
return chunks
@@ -406,13 +497,24 @@ class ExchangeBundle:
:param end_dt:
:return:
"""
writer = self.get_writer(start_dt, end_dt, data_frequency)
chunks = self.prepare_chunks(
assets=assets,
data_frequency=data_frequency,
start_dt=start_dt,
end_dt=end_dt
)
# Since chunks are either monthly or yearly, it is possible that
# our ingestion data range is greater than specified. We adjust
# the boundaries to ensure that the writer can write all data.
for chunk in chunks:
if chunk['period_start'] < start_dt:
start_dt = chunk['period_start']
if chunk['period_end'] > end_dt:
end_dt = chunk['period_end']
writer = self.get_writer(start_dt, end_dt, data_frequency)
with maybe_show_progress(
chunks,
show_progress,
@@ -428,7 +530,8 @@ class ExchangeBundle:
start_dt=chunk['period_start'],
end_dt=chunk['period_end'],
writer=writer,
empty_rows_behavior='strip'
empty_rows_behavior='strip',
cleanup=True
)
def ingest(self, data_frequency, include_symbols=None,
@@ -446,7 +549,9 @@ class ExchangeBundle:
:return:
"""
assets = self.get_assets(include_symbols, exclude_symbols)
start_dt, end_dt = get_adj_dates(start, end, assets, data_frequency)
start_dt, end_dt = self.get_adj_dates(
start, end, assets, data_frequency
)
for frequency in data_frequency.split(','):
self.ingest_assets(assets, start_dt, end_dt, frequency,
@@ -515,7 +620,7 @@ class ExchangeBundle:
return values
except Exception:
symbols = [asset.symbol.encode('utf-8') for asset in assets]
symbols = [asset.symbol for asset in assets]
raise PricingDataNotLoadedError(
field=field,
first_trading_day=min([asset.start_date for asset in assets]),
@@ -533,8 +638,9 @@ class ExchangeBundle:
data_frequency,
reset_reader=False):
start_dt = get_start_dt(end_dt, bar_count, data_frequency)
start_dt, end_dt = \
get_adj_dates(start_dt, end_dt, assets, data_frequency)
start_dt, end_dt = self.get_adj_dates(
start_dt, end_dt, assets, data_frequency
)
reader = self.get_reader(data_frequency)
if reset_reader:
@@ -542,7 +648,7 @@ class ExchangeBundle:
reader = self.get_reader(data_frequency)
if reader is None:
symbols = [asset.symbol.encode('utf-8') for asset in assets]
symbols = [asset.symbol for asset in assets]
raise PricingDataNotLoadedError(
field=field,
first_trading_day=min([asset.start_date for asset in assets]),
@@ -553,8 +659,9 @@ class ExchangeBundle:
)
for asset in assets:
asset_start_dt, asset_end_dt = \
get_adj_dates(start_dt, end_dt, assets, data_frequency)
asset_start_dt, asset_end_dt = self.get_adj_dates(
start_dt, end_dt, assets, data_frequency
)
in_bundle = range_in_bundle(
asset, asset_start_dt, asset_end_dt, reader
@@ -600,3 +707,34 @@ class ExchangeBundle:
series[asset] = value_series
return series
def clean(self, data_frequency):
log.debug('cleaning exchange {}, frequency {}'.format(
self.exchange.name, data_frequency
))
root = get_exchange_folder(self.exchange.name)
symbols = os.path.join(root, 'symbols.json')
if os.path.isfile(symbols):
os.remove(symbols)
temp_bundles = os.path.join(root, 'temp_bundles')
if os.path.isdir(temp_bundles):
log.debug('removing folder and content: {}'.format(temp_bundles))
shutil.rmtree(temp_bundles)
log.debug('{} removed'.format(temp_bundles))
frequencies = ['daily', 'minute'] if data_frequency is None \
else [data_frequency]
for frequency in frequencies:
label = '{}_bundle'.format(frequency)
frequency_bundle = os.path.join(root, label)
if os.path.isdir(frequency_bundle):
log.debug(
'removing folder and content: {}'.format(frequency_bundle)
)
shutil.rmtree(frequency_bundle)
log.debug('{} removed'.format(frequency_bundle))
+19 -7
View File
@@ -1,14 +1,17 @@
import sys, traceback
import sys
import traceback
from catalyst.errors import ZiplineError
def silent_except_hook(exctype, excvalue, exctraceback):
if exctype in [PricingDataBeforeTradingError, PricingDataNotLoadedError,
SymbolNotFoundOnExchange, NoDataAvailableOnExchange, ]:
SymbolNotFoundOnExchange, NoDataAvailableOnExchange,
ExchangeAuthEmpty]:
fn = traceback.extract_tb(exctraceback)[-1][0]
ln = traceback.extract_tb(exctraceback)[-1][1]
print "Error traceback: {1} (line {2})\n" \
"{0.__name__}: {3}".format(exctype, fn, ln, excvalue)
print("Error traceback: {1} (line {2})\n"
"{0.__name__}: {3}".format(exctype, fn, ln, excvalue))
else:
sys.__excepthook__(exctype, excvalue, exctraceback)
@@ -63,6 +66,13 @@ class ExchangeAuthNotFound(ZiplineError):
).strip()
class ExchangeAuthEmpty(ZiplineError):
msg = (
'Please enter your API token key and secret for exchange {exchange} '
'in the following file: {filename}'
).strip()
class ExchangeSymbolsNotFound(ZiplineError):
msg = (
'Unable to download or find a local copy of symbols.json for exchange '
@@ -204,7 +214,9 @@ class PricingDataNotLoadedError(ZiplineError):
class ApiCandlesError(ZiplineError):
msg = ('Unable to fetch candles from the remote API: {error}.').strip()
class NoDataAvailableOnExchange(ZiplineError):
msg = ('Requested data for trading pair {symbol} is not available on exchange {exchange} '
'in `{data_frequency}` frequency at this time. '
'Check `http://enigma.co/catalyst/status` for market coverage.').strip()
msg = (
'Requested data for trading pair {symbol} is not available on exchange {exchange} '
'in `{data_frequency}` frequency at this time. '
'Check `http://enigma.co/catalyst/status` for market coverage.').strip()
+1 -2
View File
@@ -1,9 +1,8 @@
import numpy as np
from logbook import Logger
from catalyst.protocol import Portfolio, Positions, Position
from catalyst.constants import LOG_LEVEL
from catalyst.protocol import Portfolio, Positions, Position
log = Logger('ExchangePortfolio', level=LOG_LEVEL)
+23 -12
View File
@@ -1,14 +1,16 @@
import json
import os
import pickle
import urllib
from catalyst.assets._assets import TradingPair
from six.moves.urllib import request
from datetime import date, datetime
import pandas as pd
from catalyst.exchange.exchange_errors import ExchangeAuthNotFound, \
ExchangeSymbolsNotFound
from catalyst.utils.paths import data_root, ensure_directory, last_modified_time
from catalyst.exchange.exchange_errors import ExchangeSymbolsNotFound
from catalyst.utils.paths import data_root, ensure_directory, \
last_modified_time
SYMBOLS_URL = 'https://s3.amazonaws.com/enigmaco/catalyst-exchanges/' \
'{exchange}/symbols.json'
@@ -33,7 +35,7 @@ def get_exchange_symbols_filename(exchange_name, environ=None):
def download_exchange_symbols(exchange_name, environ=None):
filename = get_exchange_symbols_filename(exchange_name)
url = SYMBOLS_URL.format(exchange=exchange_name)
response = urllib.urlretrieve(url=url, filename=filename)
response = request.urlretrieve(url=url, filename=filename)
return response
@@ -41,7 +43,9 @@ def get_exchange_symbols(exchange_name, environ=None):
filename = get_exchange_symbols_filename(exchange_name)
if not os.path.isfile(filename) or \
pd.Timedelta(pd.Timestamp('now', tz='UTC') - last_modified_time(filename)).days > 1:
pd.Timedelta(pd.Timestamp('now',
tz='UTC') - last_modified_time(
filename)).days > 1:
download_exchange_symbols(exchange_name, environ)
if os.path.isfile(filename):
@@ -55,6 +59,11 @@ def get_exchange_symbols(exchange_name, environ=None):
)
def get_symbols_string(assets):
array = [assets] if isinstance(assets, TradingPair) else assets
return ', '.join([asset.symbol for asset in array])
def get_exchange_auth(exchange_name, environ=None):
exchange_folder = get_exchange_folder(exchange_name, environ)
filename = os.path.join(exchange_folder, 'auth.json')
@@ -64,10 +73,11 @@ def get_exchange_auth(exchange_name, environ=None):
data = json.load(data_file)
return data
else:
raise ExchangeAuthNotFound(
exchange=exchange_name,
filename=filename
)
data = dict(name=exchange_name, key='', secret='')
with open(filename, 'w') as f:
json.dump(data, f, sort_keys=False, indent=2,
separators=(',', ':'))
return data
def get_algo_folder(algo_name, environ=None):
@@ -151,8 +161,8 @@ def save_algo_df(algo_name, key, df, environ=None, rel_path=None):
filename = os.path.join(folder, key + '.csv')
with open(filename, 'wb') as handle:
df.to_csv(handle)
with open(filename, 'wt') as handle:
df.to_csv(handle, encoding='UTF_8')
def get_exchange_minute_writer_root(exchange_name, environ=None):
@@ -163,6 +173,7 @@ def get_exchange_minute_writer_root(exchange_name, environ=None):
return minute_data_folder
def get_exchange_bundles_folder(exchange_name, environ=None):
exchange_folder = get_exchange_folder(exchange_name, environ)
+1 -3
View File
@@ -10,7 +10,6 @@
# 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 pandas as pd
from catalyst.gens.sim_engine import (
@@ -19,11 +18,10 @@ from catalyst.gens.sim_engine import (
)
from logbook import Logger
from catalyst.constants import LOG_LEVEL
from catalyst.exchange.exchange_errors import \
MismatchingBaseCurrenciesExchanges
from catalyst.constants import LOG_LEVEL
log = Logger('LiveGraphClock', level=LOG_LEVEL)
+37 -31
View File
@@ -1,46 +1,39 @@
import base64
import hashlib
import hmac
import json
import re
import json
import time
from collections import defaultdict
import numpy as np
import pandas as pd
import pytz
import requests
# import six
from six import iteritems
from catalyst.assets._assets import TradingPair
from logbook import Logger
# import six
from six import iteritems
from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.poloniex.poloniex_api import Poloniex_api
from catalyst.constants import LOG_LEVEL
# from websocket import create_connection
from catalyst.exchange.exchange import Exchange
from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.exchange_errors import (
ExchangeRequestError,
InvalidHistoryFrequencyError,
InvalidOrderStyle, OrderCancelError,
OrphanOrderReverseError)
InvalidOrderStyle, OrphanOrderReverseError)
from catalyst.exchange.exchange_execution import ExchangeLimitOrder, \
ExchangeStopLimitOrder, ExchangeStopOrder
from catalyst.finance.order import Order, ORDER_STATUS
from catalyst.protocol import Account
ExchangeStopLimitOrder
from catalyst.exchange.exchange_utils import get_exchange_symbols_filename, \
download_exchange_symbols
download_exchange_symbols, get_symbols_string
from catalyst.exchange.poloniex.poloniex_api import Poloniex_api
from catalyst.finance.order import Order, ORDER_STATUS
from catalyst.finance.transaction import Transaction
from catalyst.constants import LOG_LEVEL
from catalyst.protocol import Account
log = Logger('Poloniex', level=LOG_LEVEL)
class Poloniex(Exchange):
def __init__(self, key, secret, base_currency, portfolio=None):
self.api = Poloniex_api(key=key, secret=secret.encode('UTF-8'))
self.api = Poloniex_api(key=key, secret=secret)
self.name = 'poloniex'
self.assets = {}
self.load_assets()
@@ -126,9 +119,9 @@ class Poloniex(Exchange):
return order, executed_price
def get_balances(self):
log.debug('retrieving wallets balances')
balances = self.api.returnbalances()
try:
balances = self.api.returnbalances()
log.debug('retrieving wallets balances')
except Exception as e:
log.debug(e)
raise ExchangeRequestError(error=e)
@@ -193,22 +186,35 @@ class Poloniex(Exchange):
'5m', '15m', '30m', '2h', '4h', '1D'
"""
# TODO: implement end_dt and start_dt filters
if end_dt is None:
end_dt = pd.Timestamp.utcnow()
if (
data_frequency == '5m' or data_frequency == 'minute'): # TODO: Polo does not have '1m'
log.debug(
'retrieving {bars} {freq} candles on {exchange} from '
'{end_dt} for markets {symbols}, '.format(
bars=bar_count,
freq=data_frequency,
exchange=self.name,
end_dt=end_dt,
symbols=get_symbols_string(assets)
)
)
if data_frequency == '5m':
frequency = 300
elif (data_frequency == '15m'):
elif data_frequency == '15m':
frequency = 900
elif (data_frequency == '30m'):
elif data_frequency == '30m':
frequency = 1800
elif (data_frequency == '2h'):
elif data_frequency == '2h':
frequency = 7200
elif (data_frequency == '4h'):
elif data_frequency == '4h':
frequency = 14400
elif (data_frequency == '1D' or data_frequency == 'daily'):
elif data_frequency == '1D' or data_frequency == 'daily':
frequency = 86400
else:
# Poloniex does not offer 1m data candles
# It is likely to error out there frequently
raise InvalidHistoryFrequencyError(
frequency=data_frequency
)
@@ -219,8 +225,8 @@ class Poloniex(Exchange):
for asset in asset_list:
end = int(time.time())
if (bar_count is None):
end = int(time.mktime(end_dt.timetuple()))
if bar_count is None:
start = end - 2 * frequency
else:
start = end - bar_count * frequency
+81 -51
View File
@@ -19,19 +19,25 @@ class Poloniex_api(object):
self.max_requests_per_second = 6
self.request_cpt = dict()
self.public = ['returnTicker', 'return24Volume', 'returnOrderBook',
'returnTradeHistory', 'returnChartData',
'returnCurrencies', 'returnLoanOrders']
self.trading = ['returnBalances','returnCompleteBalances','returnDepositAddresses',
'generateNewAddress','returnDepositsWithdrawals','returnOpenOrders',
'returnTradeHistory','returnOrderTrades',
self.public = ['returnTicker', 'return24Volume', 'returnOrderBook',
'returnTradeHistory', 'returnChartData',
'returnCurrencies', 'returnLoanOrders']
self.trading = ['returnBalances', 'returnCompleteBalances',
'returnDepositAddresses',
'generateNewAddress', 'returnDepositsWithdrawals',
'returnOpenOrders',
'returnTradeHistory', 'returnOrderTrades',
'buy', 'sell', 'cancelOrder', 'moveOrder',
'withdraw', 'returnFeeInfo','returnAvailableAccountBalances',
'withdraw', 'returnFeeInfo',
'returnAvailableAccountBalances',
'returnTradableBalances', 'transferBalance',
'returnMarginAccountSummary','marginBuy','marginSell',
'getMarginPosition', 'closeMarginPosition','createLoanOffer',
'cancelLoanOffer','returnOpenLoanOffers','returnActiveLoans',
'returnLendingHistory','toggleAutoRenew']
'returnMarginAccountSummary', 'marginBuy',
'marginSell',
'getMarginPosition', 'closeMarginPosition',
'createLoanOffer',
'cancelLoanOffer', 'returnOpenLoanOffers',
'returnActiveLoans',
'returnLendingHistory', 'toggleAutoRenew']
def ask_request(self):
"""
@@ -50,7 +56,7 @@ class Poloniex_api(object):
self.request_cpt[now] = 0
return True
cpt_date = self.request_cpt.keys()[0]
cpt_date = list(self.request_cpt.keys())[0]
cpt = self.request_cpt[cpt_date]
if now > cpt_date + 1:
@@ -59,9 +65,8 @@ class Poloniex_api(object):
return True
if cpt >= self.max_requests_per_second:
log.debug('max requests 6 reached, sleeping for 1 seconds')
sleep(1)
time.sleep(1)
now = time.time()
self.request_cpt = dict()
@@ -73,21 +78,34 @@ class Poloniex_api(object):
def query(self, method, req={}):
if method in self.public:
url = 'https://poloniex.com/public?command=' + method + '&' + urllib.parse.urlencode(req)
url = 'https://poloniex.com/public?command=' + method + '&' + \
urllib.parse.urlencode(req)
headers = {}
post_data = None
elif method in self.trading:
url = 'https://poloniex.com/tradingApi'
req['command'] = method
req['nonce'] = int(time.time()*1000)
post_data = urllib.parse.urlencode(req)
signature = hmac.new(self.secret, post_data, hashlib.sha512).hexdigest()
headers = { 'Sign': signature, 'Key': self.key}
req['nonce'] = int(time.time() * 1000)
post_data = urllib.parse.urlencode(req)
signature = hmac.new(self.secret.encode('utf-8'),
post_data.encode('utf-8'),
hashlib.sha512).hexdigest()
headers = {'Sign': signature, 'Key': self.key}
post_data = post_data.encode('utf-8')
else:
raise ValueError('Method "' + method + '" not found in neither the Public API or Trading API endpoints')
raise ValueError(
'Method "' + method + '" not found in neither the Public API '
'or Trading API endpoints'
)
self.ask_request()
req = urllib.request.Request(url, data=post_data, headers=headers)
req = urllib.request.Request(
url,
data=post_data,
headers=headers
)
return json.loads(urlopen(req).read())
def returnticker(self):
@@ -100,15 +118,17 @@ class Poloniex_api(object):
return self.query('returnOrderBook', {'currencyPair': market})
def returntradehistory(self, market, start=None, end=None):
if(start is not None and end is not None):
return self.query('returntradehistory',
{'currencyPair': market, 'start': start, 'end': end })
if (start is not None and end is not None):
return self.query('returntradehistory',
{'currencyPair': market, 'start': start,
'end': end})
else:
return self.query('returntradehistory', {'currencyPair': market })
return self.query('returntradehistory', {'currencyPair': market})
def returnchartdata(self, market, period, start, end=9999999999):
return self.query('returnChartData', {'currencyPair': market, 'period': period,
'start': start, 'end': end})
return self.query('returnChartData',
{'currencyPair': market, 'period': period,
'start': start, 'end': end})
def returncurrencies(self):
return self.query('returnCurrencies', {})
@@ -120,7 +140,7 @@ class Poloniex_api(object):
return self.query('returnBalances')
def returncompletebalances(self, account):
if(account):
if (account):
return self.query('returnCompleteBalances', {'account': account})
else:
return self.query('returnCompleteBalances')
@@ -132,43 +152,54 @@ class Poloniex_api(object):
return self.query('generateNewAddress', {'currency': currency})
def returnDepositsWithdrawals(self, start, end):
return self.query('returnDepositsWithdrawals', {'start': start, 'end': end})
return self.query('returnDepositsWithdrawals',
{'start': start, 'end': end})
def returnopenorders(self, market):
return self.query('returnOpenOrders', {'currencyPair': market})
def returntradehistory(self, market):
#TODO: optional start and/or end and limit
# TODO: optional start and/or end and limit
return self.query('returnTradeHistory', {'currencyPair': market})
def returnordertrades(self, ordernumber):
return self.query('returnOrderTrades', {'orderNumber': ordernumber})
def buy(self, market, amount, rate, fillorkill=0, immediateorcancel=0, postonly=0):
if(fillorkill):
return self.query('buy', {'currencyPair': market, 'rate':rate, 'amount': amount,
def buy(self, market, amount, rate, fillorkill=0, immediateorcancel=0,
postonly=0):
if (fillorkill):
return self.query('buy', {'currencyPair': market, 'rate': rate,
'amount': amount,
'fillOrKill': fillorkill, })
elif(immediateorcancel):
return self.query('buy', {'currencyPair': market, 'rate':rate, 'amount': amount,
elif (immediateorcancel):
return self.query('buy', {'currencyPair': market, 'rate': rate,
'amount': amount,
'immediateOrCancel': immediateorcancel, })
elif(postonly):
return self.query('buy', {'currencyPair': market, 'rate':rate, 'amount': amount,
elif (postonly):
return self.query('buy', {'currencyPair': market, 'rate': rate,
'amount': amount,
'postOnly': postonly, })
else:
return self.query('buy', {'currencyPair': market, 'rate':rate, 'amount': amount, })
return self.query('buy', {'currencyPair': market, 'rate': rate,
'amount': amount, })
def sell(self, market, amount, rate, fillorkill=0, immediateorcancel=0, postonly=0):
if(fillorkill):
return self.query('sell', {'currencyPair': market, 'rate':rate, 'amount': amount,
'fillOrKill': fillorkill, })
elif(immediateorcancel):
return self.query('sell', {'currencyPair': market, 'rate':rate, 'amount': amount,
'immediateOrCancel': immediateorcancel, })
elif(postonly):
return self.query('sell', {'currencyPair': market, 'rate':rate, 'amount': amount,
'postOnly': postonly, })
def sell(self, market, amount, rate, fillorkill=0, immediateorcancel=0,
postonly=0):
if (fillorkill):
return self.query('sell', {'currencyPair': market, 'rate': rate,
'amount': amount,
'fillOrKill': fillorkill, })
elif (immediateorcancel):
return self.query('sell', {'currencyPair': market, 'rate': rate,
'amount': amount,
'immediateOrCancel': immediateorcancel, })
elif (postonly):
return self.query('sell', {'currencyPair': market, 'rate': rate,
'amount': amount,
'postOnly': postonly, })
else:
return self.query('sell', {'currencyPair': market, 'rate':rate, 'amount': amount, })
return self.query('sell', {'currencyPair': market, 'rate': rate,
'amount': amount, })
def cancelorder(self, ordernumber):
return self.query('cancelOrder', {'orderNumber': ordernumber})
@@ -180,4 +211,3 @@ class Poloniex_api(object):
def returnfeeinfo(self):
return self.query('returnFeeInfo')
+4 -4
View File
@@ -16,13 +16,13 @@ from time import sleep
import pandas as pd
from catalyst.gens.sim_engine import (
BAR,
SESSION_START,
MINUTE_END,
SESSION_END
SESSION_START
)
from logbook import Logger
log = Logger('ExchangeClock')
from catalyst.constants import LOG_LEVEL
log = Logger('ExchangeClock', level=LOG_LEVEL)
class SimpleClock(object):
+9
View File
@@ -49,3 +49,12 @@ def get_pretty_stats(stats_df, recorded_cols=None, num_rows=10):
columns=columns,
formatters=formatters
)
def df_to_string(df):
pd.set_option('display.expand_frame_repr', False)
pd.set_option('precision', 8)
pd.set_option('display.width', 1000)
pd.set_option('display.max_colwidth', 1000)
return df.to_string()
+1 -1
View File
@@ -126,7 +126,7 @@ def catalyst_root(environ=None):
root = environ.get('ZIPLINE_ROOT', None)
if root is None:
root = expanduser('~/.catalyst')
root = os.path.join(expanduser('~'),'.catalyst')
return root
+13 -4
View File
@@ -36,11 +36,11 @@ from catalyst.exchange.data_portal_exchange import DataPortalExchangeLive, \
from catalyst.exchange.asset_finder_exchange import AssetFinderExchange
from catalyst.exchange.exchange_portfolio import ExchangePortfolio
from catalyst.exchange.exchange_errors import (
ExchangeRequestError,
ExchangeRequestError, ExchangeAuthEmpty,
ExchangeRequestErrorTooManyAttempts,
BaseCurrencyNotFoundError, ExchangeNotFoundError)
from catalyst.exchange.exchange_utils import get_exchange_auth, \
get_algo_object
get_algo_object, get_exchange_folder
from logbook import Logger
from catalyst.constants import LOG_LEVEL
@@ -166,6 +166,12 @@ def _run(handle_data,
# This corresponds to the json file containing api token info
exchange_auth = get_exchange_auth(exchange_name)
if live and (exchange_auth['key'] == '' or exchange_auth['secret'] == ''):
raise ExchangeAuthEmpty(
exchange=exchange_name.title(),
filename=os.path.join(get_exchange_folder(exchange_name, environ), 'auth.json') )
if exchange_name == 'bitfinex':
exchanges[exchange_name] = Bitfinex(
key=exchange_auth['key'],
@@ -237,8 +243,11 @@ def _run(handle_data,
balances = exchange.get_balances()
except ExchangeRequestError as e:
if attempt_index < 20:
log.warn('exchange error when retrieving balances, {} '
'trying again in 5 seconds'.format(e))
log.warn(
'could not retrieve balances on {}: {}'.format(
exchange.name, e
)
)
sleep(5)
return fetch_capital_base(exchange, attempt_index + 1)
+11 -6
View File
@@ -429,6 +429,10 @@ and allows us to plot the price of bitcoin. For example, we could easily
examine now how our portfolio value changed over time compared to the
bitcoin price.
.. code-block:: python
%load_ext catalyst
.. code-block:: python
%pylab inline
@@ -484,7 +488,8 @@ a function we use in the ``handle_data()`` section:
.. code-block:: python
%%catalyst --start 2016-1-1 --end 2017-9-30 -x bitfinex -o dma.pickle
%%catalyst --start 2016-4-1 --end 2017-9-30 -x bitfinex
from catalyst.api import order, record, symbol, order_target
def initialize(context):
@@ -492,16 +497,16 @@ a function we use in the ``handle_data()`` section:
context.asset = symbol('btc_usd')
def handle_data(context, data):
# Skip first 300 days to get full windows
# Skip first 150 days to get full windows
context.i += 1
if context.i < 300:
if context.i < 150:
return
# Compute averages
# data.history() has to be called with the same params
# from above and returns a pandas dataframe.
short_mavg = data.history(context.asset, 'price', bar_count=100, frequency="1d").mean()
long_mavg = data.history(context.asset, 'price', bar_count=300, frequency="1d").mean()
short_mavg = data.history(context.asset, 'price', bar_count=50, frequency="1d").mean()
long_mavg = data.history(context.asset, 'price', bar_count=150, frequency="1d").mean()
# Trading logic
if short_mavg > long_mavg:
@@ -518,7 +523,7 @@ a function we use in the ``handle_data()`` section:
def analyze(context, perf):
import matplotlib.pyplot as plt
fig = plt.figure()
fig = plt.figure(figsize=(12,12))
ax1 = fig.add_subplot(211)
perf.portfolio_value.plot(ax=ax1)
ax1.set_ylabel('portfolio value in $')
+14 -40
View File
@@ -1,30 +1,22 @@
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
- mkl=2017.0.3=0
- numpy=1.13.1=py27_0
- openssl=1.0.2l
- 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
- scipy=0.19.1=np113py27_0
- setuptools=36.4.0=py27_1
- sqlite=3.13.0
- tk=8.5.18
- 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
- alembic==0.9.6
- backports.functools-lru-cache==1.4
- bcolz==0.12.1
- bottleneck==1.2.1
- chardet==3.0.4
@@ -32,36 +24,22 @@ dependencies:
- contextlib2==0.5.5
- cycler==0.10.0
- cyordereddict==1.0.0
- cython==0.26.1
- cython==0.27.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
- matplotlib==2.1.0
- multipledispatch==0.4.9
- networkx==1.11
- networkx==2.0
- 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
@@ -69,16 +47,12 @@ dependencies:
- 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
- six==1.11.0
- 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
- enigma-catalyst>=0.3
+150
View File
@@ -0,0 +1,150 @@
import shutil
import random
import tempfile
import pandas as pd
from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.exchange_bcolz import BcolzExchangeBarWriter, \
BcolzExchangeBarReader
from catalyst.exchange.bundle_utils import get_df_from_arrays
from nose.tools import assert_equals
class TestBcolzWriter(object):
@classmethod
def setup_class(cls):
cls.columns = ['open', 'high', 'low', 'close', 'volume']
def setUp(self):
self.root_dir = tempfile.mkdtemp() # Create a temporary directory
def tearDown(self):
shutil.rmtree(self.root_dir) # Remove the directory after the test
def generate_df(self, exchange_name, freq, start, end):
bundle = ExchangeBundle(exchange_name)
index = bundle.get_calendar_periods_range(start, end, freq)
df = pd.DataFrame(index=index, columns=self.columns)
df.fillna(random.random(), inplace=True)
return df
def test_bcolz_write_daily_past(self):
start = pd.to_datetime('2016-01-01')
end = pd.to_datetime('2016-12-31')
freq = 'daily'
df = self.generate_df('bitfinex', freq, start, end)
writer = BcolzExchangeBarWriter(
rootdir=self.root_dir,
start_session=start,
end_session=end,
data_frequency=freq,
write_metadata=True)
data = []
data.append((1, df))
writer.write(data)
pass
def test_bcolz_write_daily_present(self):
start = pd.to_datetime('2017-01-01')
end = pd.to_datetime('today')
freq = 'daily'
df = self.generate_df('bitfinex', freq, start, end)
writer = BcolzExchangeBarWriter(
rootdir=self.root_dir,
start_session=start,
end_session=end,
data_frequency=freq,
write_metadata=True)
data = []
data.append((1, df))
writer.write(data)
pass
def test_bcolz_write_minute_past(self):
start = pd.to_datetime('2015-04-01 00:00')
end = pd.to_datetime('2015-04-30 23:59')
freq = 'minute'
df = self.generate_df('bitfinex', freq, start, end)
writer = BcolzExchangeBarWriter(
rootdir=self.root_dir,
start_session=start,
end_session=end,
data_frequency=freq,
write_metadata=True)
data = []
data.append((1, df))
writer.write(data)
pass
def test_bcolz_write_minute_present(self):
start = pd.to_datetime('2017-10-01 00:00')
end = pd.to_datetime('today')
freq = 'minute'
df = self.generate_df('bitfinex', freq, start, end)
writer = BcolzExchangeBarWriter(
rootdir=self.root_dir,
start_session=start,
end_session=end,
data_frequency=freq,
write_metadata=True)
data = []
data.append((1, df))
writer.write(data)
pass
def bcolz_exchange_daily_write_read(self, exchange_name):
start = pd.to_datetime('2017-10-01 00:00')
end = pd.to_datetime('today')
freq = 'daily'
bundle = ExchangeBundle(exchange_name)
df = self.generate_df(exchange_name, freq, start, end)
print df.index[0],df.index[-1]
writer = BcolzExchangeBarWriter(
rootdir=self.root_dir,
start_session=df.index[0],
end_session=df.index[-1],
data_frequency=freq,
write_metadata=True)
data = []
data.append((1, df))
writer.write(data)
reader = BcolzExchangeBarReader(rootdir=self.root_dir,
data_frequency=freq)
arrays = reader.load_raw_arrays(self.columns, start, end, [1, ])
periods = bundle.get_calendar_periods_range(
start, end, freq
)
dx = get_df_from_arrays(arrays, periods)
assert_equals(df.equals(df), True)
pass
def test_bcolz_bitfinex_daily_write_read(self):
self.bcolz_exchange_daily_write_read('bitfinex')
def test_bcolz_poloniex_daily_write_read(self):
self.bcolz_exchange_daily_write_read('poloniex')
+1 -1
View File
@@ -8,7 +8,7 @@ from catalyst.finance.execution import (LimitOrder)
log = Logger('test_bitfinex')
class BitfinexTestCase(BaseExchangeTestCase):
class TestBitfinexTestCase(BaseExchangeTestCase):
@classmethod
def setup(self):
log.info('creating bitfinex object')
+10 -6
View File
@@ -1,3 +1,4 @@
import pandas as pd
from catalyst.exchange.bittrex.bittrex import Bittrex
from catalyst.finance.order import Order
from base import BaseExchangeTestCase
@@ -7,15 +8,15 @@ from catalyst.exchange.exchange_utils import get_exchange_auth
log = Logger('test_bittrex')
class BittrexTestCase(BaseExchangeTestCase):
class TestBittrex(BaseExchangeTestCase):
@classmethod
def setup(self):
print ('creating bittrex object')
auth = get_exchange_auth('bittrex')
self.exchange = Bittrex(
key=auth['key'],
secret=auth['secret'],
base_currency='btc'
base_currency=None,
portfolio=None
)
def test_order(self):
@@ -52,15 +53,18 @@ class BittrexTestCase(BaseExchangeTestCase):
log.info('retrieving candles')
ohlcv_neo = self.exchange.get_candles(
data_frequency='5m',
assets=self.exchange.get_asset('neo_btc')
assets=self.exchange.get_asset('neo_btc'),
bar_count=20,
end_dt=pd.to_datetime('2017-10-20', utc=True)
)
ohlcv_neo_ubq = self.exchange.get_candles(
data_frequency='5m',
data_frequency='1d',
assets=[
self.exchange.get_asset('neo_btc'),
self.exchange.get_asset('ubq_btc')
],
bar_count=14
bar_count=14,
end_dt=pd.to_datetime('2017-10-20', utc=True)
)
pass
+153 -11
View File
@@ -1,22 +1,24 @@
from logging import Logger
import hashlib
from logging import getLogger
import pandas as pd
from catalyst import get_calendar
from catalyst.exchange.bundle_utils import get_bcolz_chunk, get_periods, \
get_periods_range
from catalyst.exchange.bundle_utils import get_bcolz_chunk, \
get_periods_range, get_start_dt
from catalyst.exchange.exchange_bcolz import BcolzExchangeBarReader, \
BcolzExchangeBarWriter
from catalyst.exchange.exchange_bundle import ExchangeBundle, \
BUNDLE_NAME_TEMPLATE
from catalyst.exchange.exchange_utils import get_exchange_folder
from catalyst.exchange.init_utils import get_exchange
from catalyst.exchange.stats_utils import df_to_string
from catalyst.utils.paths import ensure_directory
log = Logger('test_exchange_bundle')
log = getLogger('test_exchange_bundle')
class ExchangeBundleTestCase:
class TestExchangeBundle:
def test_spot_value(self):
data_frequency = 'daily'
exchange_name = 'poloniex'
@@ -43,11 +45,11 @@ class ExchangeBundleTestCase:
exchange = get_exchange(exchange_name)
exchange_bundle = ExchangeBundle(exchange)
assets = [
exchange.get_asset('neo_eth')
exchange.get_asset('iot_btc')
]
# start = pd.to_datetime('2017-09-01', utc=True)
start = pd.to_datetime('2017-9-15', utc=True)
start = pd.to_datetime('2017-9-01', utc=True)
end = pd.to_datetime('2017-9-30', utc=True)
log.info('ingesting exchange bundle {}'.format(exchange_name))
@@ -93,16 +95,39 @@ class ExchangeBundleTestCase:
)
pass
def test_ingest_exchange(self):
# exchange_name = 'bitfinex'
# data_frequency = 'daily'
# include_symbols = 'neo_btc,bch_btc,eth_btc'
exchange_name = 'bitfinex'
data_frequency = 'minute'
exchange = get_exchange(exchange_name)
exchange_bundle = ExchangeBundle(exchange)
log.info('ingesting exchange bundle {}'.format(exchange_name))
exchange_bundle.ingest(
data_frequency=data_frequency,
include_symbols=None,
exclude_symbols=None,
start=None,
end=None,
show_progress=True
)
pass
def test_ingest_daily(self):
# exchange_name = 'bitfinex'
# data_frequency = 'daily'
# include_symbols = 'neo_btc,bch_btc,eth_btc'
exchange_name = 'poloniex'
exchange_name = 'bittrex'
data_frequency = 'daily'
include_symbols = 'btc_usdt'
include_symbols = 'wings_eth'
start = pd.to_datetime('2016-1-1', utc=True)
start = pd.to_datetime('2017-1-1', utc=True)
end = pd.to_datetime('2017-10-16', utc=True)
periods = get_periods_range(start, end, data_frequency)
@@ -274,7 +299,7 @@ class ExchangeBundleTestCase:
data_frequency = 'minute'
exchange = get_exchange(exchange_name)
asset = exchange.get_asset('neo_btc')
asset = exchange.get_asset('neos_btc')
path = get_bcolz_chunk(
exchange_name=exchange_name,
@@ -284,3 +309,120 @@ class ExchangeBundleTestCase:
)
pass
def test_hash_symbol(self):
symbol = 'etc_btc'
sid = int(
hashlib.sha256(symbol.encode('utf-8')).hexdigest(), 16
) % 10 ** 6
pass
def test_validate_data(self):
exchange_name = 'bitfinex'
data_frequency = 'minute'
exchange = get_exchange(exchange_name)
exchange_bundle = ExchangeBundle(exchange)
assets = [exchange.get_asset('iot_btc')]
end_dt = pd.to_datetime('2017-9-2 1:00', utc=True)
bar_count = 60
bundle_series = exchange_bundle.get_history_window_series(
assets=assets,
end_dt=end_dt,
bar_count=bar_count * 5,
field='close',
data_frequency='minute',
)
candles = exchange.get_candles(
assets=assets,
end_dt=end_dt,
bar_count=bar_count,
data_frequency='minute'
)
start_dt = get_start_dt(end_dt, bar_count, data_frequency)
frames = []
for asset in assets:
bundle_df = pd.DataFrame(
data=dict(bundle_price=bundle_series[asset]),
index=bundle_series[asset].index
)
exchange_series = exchange.get_series_from_candles(
candles=candles[asset],
start_dt=start_dt,
end_dt=end_dt,
data_frequency=data_frequency,
field='close'
)
exchange_df = pd.DataFrame(
data=dict(exchange_price=exchange_series),
index=exchange_series.index
)
df = exchange_df.join(bundle_df, how='left')
df['last_traded'] = df.index
df['asset'] = asset.symbol
df.set_index(['asset', 'last_traded'], inplace=True)
frames.append(df)
df = pd.concat(frames)
print('\n' + df_to_string(df))
pass
def test_ingest_candles(self):
exchange_name = 'bitfinex'
data_frequency = 'minute'
exchange = get_exchange(exchange_name)
bundle = ExchangeBundle(exchange)
assets = [exchange.get_asset('iot_btc')]
end_dt = pd.to_datetime('2017-10-20', utc=True)
bar_count = 100
start_dt = get_start_dt(end_dt, bar_count, data_frequency)
candles = exchange.get_candles(
assets=assets,
start_dt=start_dt,
end_dt=end_dt,
bar_count=bar_count,
data_frequency=data_frequency
)
writer = bundle.get_writer(start_dt, end_dt, data_frequency)
for asset in assets:
dates = [candle['last_traded'] for candle in candles[asset]]
values = dict()
for field in ['open', 'high', 'low', 'close', 'volume']:
values[field] = [candle[field] for candle in candles[asset]]
periods = bundle.get_calendar_periods_range(
start_dt, end_dt, data_frequency
)
df = pd.DataFrame(values, index=dates)
df = df.loc[periods].fillna(method='ffill')
# TODO: why do I get an extra bar?
bundle.ingest_df(
ohlcv_df=df,
data_frequency=data_frequency,
asset=asset,
writer=writer,
empty_rows_behavior='raise'
)
bundle_series = bundle.get_history_window_series(
assets=assets,
end_dt=end_dt,
bar_count=bar_count,
field='close',
data_frequency=data_frequency,
reset_reader=True
)
df = pd.DataFrame(bundle_series)
print('\n' + df_to_string(df))
pass
-50
View File
@@ -1,50 +0,0 @@
from unittest import TestCase
from logbook import Logger
from mock import patch, sentinel
from catalyst.exchange.simple_clock import SimpleClock
from catalyst.utils.calendars.trading_calendar import days_at_time
from datetime import time
from collections import defaultdict
from catalyst.utils.calendars import get_calendar
import pandas as pd
log = Logger('ExchangeClockTestCase')
class ExchangeClockTestCase(TestCase):
@classmethod
def setUpClass(cls):
cls.open_calendar = get_calendar("OPEN")
cls.sessions = pd.Timestamp.utcnow()
def setUp(self):
self.internal_clock = None
self.events = defaultdict(list)
def advance_clock(self, x):
"""Mock function for sleep. Advances the internal clock by 1 min"""
# The internal clock advance time must be 1 minute to match
# MinutesSimulationClock's update frequency
self.internal_clock += pd.Timedelta('1 min')
def get_clock(self, arg, *args, **kwargs):
"""Mock function for pandas.to_datetime which is used to query the
current time in RealtimeClock"""
assert arg == "now"
return self.internal_clock
def test_clock(self):
with patch('catalyst.exchange.simple_clock.pd.to_datetime') as to_dt, \
patch('catalyst.exchange.simple_clock.sleep') as sleep:
clock = SimpleClock(sessions=self.sessions)
to_dt.side_effect = self.get_clock
sleep.side_effect = self.advance_clock
start_time = pd.Timestamp.utcnow()
self.internal_clock = start_time
events = list(clock)
# Event 0 is SESSION_START which always happens at 00:00.
ts, event_type = events[1]
pass
+1 -1
View File
@@ -12,7 +12,7 @@ from catalyst.exchange.exchange_utils import get_exchange_auth
log = Logger('test_bitfinex')
class ExchangeDataPortalTestCase:
class TestExchangeDataPortalTestCase:
@classmethod
def setup(self):
log.info('creating bitfinex exchange')
+6 -6
View File
@@ -8,7 +8,7 @@ from catalyst.exchange.exchange_utils import get_exchange_auth
log = Logger('test_poloniex')
class PoloniexTestCase(BaseExchangeTestCase):
class TestPoloniexTestCase(BaseExchangeTestCase):
@classmethod
def setup(self):
print ('creating poloniex object')
@@ -21,7 +21,7 @@ class PoloniexTestCase(BaseExchangeTestCase):
def test_order(self):
log.info('creating order')
asset = self.exchange.get_asset('neo_btc')
asset = self.exchange.get_asset('neos_btc')
order_id = self.exchange.order(
asset=asset,
limit_price=0.0005,
@@ -33,7 +33,7 @@ class PoloniexTestCase(BaseExchangeTestCase):
def test_open_orders(self):
log.info('retrieving open orders')
asset = self.exchange.get_asset('neo_btc')
asset = self.exchange.get_asset('neos_btc')
orders = self.exchange.get_open_orders(asset)
pass
@@ -53,13 +53,13 @@ class PoloniexTestCase(BaseExchangeTestCase):
log.info('retrieving candles')
ohlcv_neo = self.exchange.get_candles(
data_frequency='5m',
assets=self.exchange.get_asset('neo_btc')
assets=self.exchange.get_asset('neos_btc')
)
ohlcv_neo_ubq = self.exchange.get_candles(
data_frequency='5m',
assets=[
self.exchange.get_asset('neo_btc'),
self.exchange.get_asset('ubq_btc')
self.exchange.get_asset('neos_btc'),
self.exchange.get_asset('via_btc')
],
bar_count=14
)