mirror of
https://github.com/wassname/catalyst.git
synced 2026-07-22 12:40:30 +08:00
Compare commits
52
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2b85732e36 | ||
|
|
284c749bb5 | ||
|
|
d7f5e73f84 | ||
|
|
cde69da173 | ||
|
|
bcc75f6b00 | ||
|
|
f179381b64 | ||
|
|
10ba53b897 | ||
|
|
1cfe3b1bb2 | ||
|
|
268ff9c826 | ||
|
|
7eb184d946 | ||
|
|
1cc34a1485 | ||
|
|
aa2f2f3627 | ||
|
|
7e373e2f9c | ||
|
|
942e6f263c | ||
|
|
2e6d7d28ba | ||
|
|
3a823ea457 | ||
|
|
4daba6cfb4 | ||
|
|
fa018e2e0c | ||
|
|
315d25f7c0 | ||
|
|
cc7ffada96 | ||
|
|
5394c1bc91 | ||
|
|
b230b73829 | ||
|
|
930a68ab4a | ||
|
|
4e833981e4 | ||
|
|
2ea402ff10 | ||
|
|
da6b024edc | ||
|
|
565e9a3cea | ||
|
|
3c10d19a7e | ||
|
|
cf96e047cd | ||
|
|
6f6a8e1272 | ||
|
|
c2a02e7074 | ||
|
|
7d2cf97fbf | ||
|
|
195469897c | ||
|
|
c7b422d465 | ||
|
|
2dbace37bb | ||
|
|
2e903fd42c | ||
|
|
47a104b29c | ||
|
|
d248581523 | ||
|
|
48f6300e08 | ||
|
|
f7a143cb78 | ||
|
|
2f7cd97852 | ||
|
|
73eca75ed9 | ||
|
|
2ade2989e8 | ||
|
|
b1d5acf2ad | ||
|
|
5d5ec6b9be | ||
|
|
1b84023c5d | ||
|
|
97f3329c1b | ||
|
|
bdeb344999 | ||
|
|
52e1de954f | ||
|
|
7b9eafef4e | ||
|
|
8b141a0c28 | ||
|
|
7f602d7fcc |
+35
-1
@@ -509,6 +509,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 +632,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,
|
||||
|
||||
@@ -138,8 +138,9 @@ from catalyst.gens.sim_engine import MinuteSimulationClock
|
||||
from catalyst.sources.benchmark_source import BenchmarkSource
|
||||
from catalyst.catalyst_warnings import ZiplineDeprecationWarning
|
||||
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = logbook.Logger("ZiplineLog")
|
||||
log = logbook.Logger("CatalystLog", level=LOG_LEVEL)
|
||||
|
||||
|
||||
class TradingAlgorithm(object):
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -76,7 +76,9 @@ from catalyst.utils.numpy_utils import as_column
|
||||
from catalyst.utils.preprocess import preprocess
|
||||
from catalyst.utils.sqlite_utils import group_into_chunks, coerce_string_to_eng
|
||||
|
||||
log = Logger('assets.py')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = Logger('assets.py', level=LOG_LEVEL)
|
||||
|
||||
# A set of fields that need to be converted to strings before building an
|
||||
# Asset to avoid unicode fields
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
|
||||
import logbook
|
||||
|
||||
LOG_LEVEL = logbook.INFO
|
||||
+26
-26
@@ -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.')
|
||||
|
||||
|
||||
'''
|
||||
|
||||
@@ -215,7 +215,7 @@ cpdef _read_bcolz_data(ctable_t table,
|
||||
else:
|
||||
continue
|
||||
|
||||
if column_name in ['open', 'high', 'low', 'close']:
|
||||
if column_name in ['open', 'high', 'low', 'close', 'volume']:
|
||||
where_nan = (outbuf == 0)
|
||||
outbuf_as_float = outbuf.astype(float64) * .000000001
|
||||
outbuf_as_float[where_nan] = NAN
|
||||
|
||||
@@ -30,8 +30,10 @@ from catalyst.utils.cli import (
|
||||
)
|
||||
from catalyst.utils.memoize import lazyval
|
||||
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
logbook.StderrHandler().push_application()
|
||||
log = logbook.Logger(__name__)
|
||||
log = logbook.Logger(__name__, level=LOG_LEVEL)
|
||||
|
||||
DEFAULT_RETRIES = 5
|
||||
|
||||
|
||||
@@ -40,7 +40,9 @@ from catalyst.utils.cli import maybe_show_progress
|
||||
|
||||
from . import core as bundles
|
||||
|
||||
log = Logger(__name__)
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = Logger(__name__, level=LOG_LEVEL)
|
||||
seconds_per_call = (pd.Timedelta('10 minutes') / 2000).total_seconds()
|
||||
|
||||
class QuandlBundle(BaseEquityPricingBundle):
|
||||
|
||||
@@ -68,7 +68,9 @@ from catalyst.errors import (
|
||||
HistoryWindowStartsBeforeData,
|
||||
)
|
||||
|
||||
log = Logger('DataPortal')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = Logger('DataPortal', level=LOG_LEVEL)
|
||||
|
||||
BASE_FIELDS = frozenset([
|
||||
"open",
|
||||
|
||||
@@ -32,7 +32,9 @@ from ..utils.paths import (
|
||||
data_root,
|
||||
)
|
||||
|
||||
logger = logbook.Logger('Loader')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
logger = logbook.Logger('Loader', level=LOG_LEVEL)
|
||||
|
||||
# Mapping from index symbol to appropriate bond data
|
||||
INDEX_MAPPING = {
|
||||
|
||||
@@ -44,7 +44,9 @@ from catalyst.utils.calendars import get_calendar
|
||||
from catalyst.utils.cli import maybe_show_progress
|
||||
from catalyst.utils.memoize import lazyval
|
||||
|
||||
logger = logbook.Logger('MinuteBars')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
logger = logbook.Logger('MinuteBars', level=LOG_LEVEL)
|
||||
|
||||
US_EQUITIES_MINUTES_PER_DAY = 390
|
||||
FUTURES_MINUTES_PER_DAY = 1440
|
||||
|
||||
@@ -83,7 +83,9 @@ from catalyst.utils.cli import (
|
||||
from ._equities import _compute_row_slices, _read_bcolz_data
|
||||
from ._adjustments import load_adjustments_from_sqlite
|
||||
|
||||
logger = logbook.Logger('UsEquityPricing')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
logger = logbook.Logger('UsEquityPricing', level=LOG_LEVEL)
|
||||
|
||||
OHLC = frozenset(['open', 'high', 'low', 'close'])
|
||||
OHLCV = frozenset(['open', 'high', 'low', 'close', 'volume'])
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,8 +1,8 @@
|
||||
from catalyst.api import order, record, symbol
|
||||
|
||||
def initialize(context):
|
||||
context.asset = symbol('btc_usd')
|
||||
context.asset = symbol('btc_usd')
|
||||
|
||||
def handle_data(context, data):
|
||||
order(asset, 1)
|
||||
record(btc=data.current(context.asset, 'price'))
|
||||
order(context.asset, 1)
|
||||
record(btc = data.current(context.asset, 'price'))
|
||||
@@ -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
|
||||
|
||||
@@ -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('btc_usd')
|
||||
|
||||
|
||||
def handle_data(context, data):
|
||||
@@ -20,32 +20,32 @@ def handle_data(context, data):
|
||||
context.asset,
|
||||
fields='price',
|
||||
bar_count=15,
|
||||
frequency='1d'
|
||||
frequency='1m'
|
||||
)
|
||||
rsi = talib.RSI(prices.values, timeperiod=14)[-1]
|
||||
print('got rsi: {}'.format(rsi))
|
||||
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,
|
||||
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
|
||||
# )
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
from logbook import Logger
|
||||
|
||||
log = Logger('AssetFinderExchange')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = Logger('AssetFinderExchange', level=LOG_LEVEL)
|
||||
|
||||
|
||||
class AssetFinderExchange(object):
|
||||
@@ -41,9 +43,9 @@ class AssetFinderExchange(object):
|
||||
"""
|
||||
for sid in sids:
|
||||
if sid in self._asset_cache:
|
||||
log.info('got asset from cache: {}'.format(sid))
|
||||
log.debug('got asset from cache: {}'.format(sid))
|
||||
else:
|
||||
log.info('fetching asset: {}'.format(sid))
|
||||
log.debug('fetching asset: {}'.format(sid))
|
||||
return list()
|
||||
|
||||
def lookup_symbol(self, symbol, exchange, as_of_date=None, fuzzy=False):
|
||||
|
||||
@@ -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.finance.order import Order, ORDER_STATUS
|
||||
from catalyst.protocol import Account
|
||||
from catalyst.exchange.exchange_utils import get_exchange_symbols_filename, \
|
||||
download_exchange_symbols
|
||||
from catalyst.finance.order import Order, ORDER_STATUS
|
||||
from catalyst.protocol import Account
|
||||
|
||||
# Trying to account for REST api instability
|
||||
# https://stackoverflow.com/questions/15431044/can-i-set-max-retries-for-requests-request
|
||||
@@ -33,7 +33,9 @@ requests.adapters.DEFAULT_RETRIES = 20
|
||||
|
||||
BITFINEX_URL = 'https://api.bitfinex.com'
|
||||
|
||||
log = Logger('Bitfinex')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = Logger('Bitfinex', level=LOG_LEVEL)
|
||||
warning_logger = Logger('AlgoWarning')
|
||||
|
||||
|
||||
|
||||
@@ -5,25 +5,26 @@ 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.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.finance.execution import LimitOrder, StopLimitOrder
|
||||
from catalyst.finance.order import Order, ORDER_STATUS
|
||||
|
||||
log = Logger('Bittrex')
|
||||
log = Logger('Bittrex', level=LOG_LEVEL)
|
||||
|
||||
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
|
||||
@@ -64,10 +65,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)
|
||||
@@ -207,7 +208,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
|
||||
-------------------
|
||||
|
||||
@@ -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 = {}
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -19,17 +19,15 @@ 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)
|
||||
PricingDataNotLoadedError)
|
||||
|
||||
log = Logger('DataPortalExchange')
|
||||
log = Logger('DataPortalExchange', level=LOG_LEVEL)
|
||||
|
||||
|
||||
class DataPortalExchangeBase(DataPortal):
|
||||
@@ -82,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,
|
||||
@@ -167,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)
|
||||
|
||||
|
||||
@@ -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,7 +24,7 @@ from catalyst.exchange.exchange_utils import get_exchange_symbols
|
||||
from catalyst.finance.order import ORDER_STATUS
|
||||
from catalyst.finance.transaction import Transaction
|
||||
|
||||
log = Logger('Exchange')
|
||||
log = Logger('Exchange', level=LOG_LEVEL)
|
||||
|
||||
|
||||
class Exchange:
|
||||
@@ -87,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):
|
||||
@@ -167,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(),
|
||||
@@ -552,7 +554,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
|
||||
|
||||
@@ -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,10 @@ 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.utils.preprocess import preprocess
|
||||
|
||||
log = logbook.Logger('exchange_algorithm')
|
||||
log = logbook.Logger('exchange_algorithm', level=LOG_LEVEL)
|
||||
|
||||
|
||||
class ExchangeAlgorithmExecutor(AlgorithmSimulator):
|
||||
@@ -112,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]
|
||||
|
||||
@@ -523,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
|
||||
|
||||
@@ -555,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,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):
|
||||
|
||||
@@ -1,12 +1,13 @@
|
||||
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
|
||||
|
||||
log = Logger('exchange_blotter')
|
||||
log = Logger('exchange_blotter', level=LOG_LEVEL)
|
||||
|
||||
# It seems like we need to accept greater slippage risk in cryptos
|
||||
# Orders won't often close at Equity levels.
|
||||
|
||||
@@ -3,34 +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_bcolz_chunk, get_delta, get_month_start_end, \
|
||||
get_year_start_end, get_periods_range, 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
|
||||
|
||||
log = Logger('exchange_bundle', level=LOG_LEVEL)
|
||||
|
||||
BUNDLE_NAME_TEMPLATE = os.path.join('{root}', '{frequency}_bundle')
|
||||
|
||||
|
||||
def _cachpath(symbol, type_):
|
||||
return '-'.join([symbol, type_])
|
||||
|
||||
|
||||
BUNDLE_NAME_TEMPLATE = '{root}/{frequency}_bundle'
|
||||
log = Logger('exchange_bundle')
|
||||
log.level = INFO
|
||||
|
||||
|
||||
class ExchangeBundle:
|
||||
def __init__(self, exchange):
|
||||
self.exchange = exchange
|
||||
@@ -173,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)
|
||||
@@ -226,12 +223,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
|
||||
@@ -289,6 +292,7 @@ class ExchangeBundle:
|
||||
if not df.empty:
|
||||
df.sort_index(inplace=True)
|
||||
data.append((asset.sid, df))
|
||||
|
||||
self._write(data, writer, data_frequency)
|
||||
|
||||
if cleanup:
|
||||
@@ -298,6 +302,45 @@ class ExchangeBundle:
|
||||
|
||||
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) \
|
||||
and asset.start_date >= self.calendar.first_session:
|
||||
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 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
|
||||
@@ -314,23 +357,26 @@ 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:
|
||||
# If not, we continue to the next asset
|
||||
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]:
|
||||
@@ -344,29 +390,49 @@ 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)
|
||||
|
||||
# TODO: redundant gate, we are already filtering dates
|
||||
if first_trading_dt > period_start:
|
||||
dt += timedelta(days=1)
|
||||
continue
|
||||
|
||||
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)
|
||||
|
||||
# TODO: redundant gate, we are already filtering dates
|
||||
if first_trading_dt > period_start:
|
||||
dt += timedelta(days=1)
|
||||
continue
|
||||
|
||||
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_minute
|
||||
)
|
||||
if asset_end_year == period_end \
|
||||
and period_end > asset_end:
|
||||
period_end = asset_end
|
||||
and period_end > asset.end_minute:
|
||||
period_end = asset.end_minute
|
||||
|
||||
else:
|
||||
raise InvalidHistoryFrequencyError(
|
||||
frequency=data_frequency
|
||||
@@ -376,10 +442,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(
|
||||
@@ -393,6 +462,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
|
||||
@@ -407,13 +477,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,
|
||||
@@ -429,7 +510,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,
|
||||
@@ -447,7 +529,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,
|
||||
@@ -516,7 +600,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]),
|
||||
@@ -534,8 +618,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:
|
||||
@@ -543,7 +628,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]),
|
||||
@@ -554,8 +639,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
|
||||
@@ -601,3 +687,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))
|
||||
|
||||
@@ -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,9 +1,10 @@
|
||||
import numpy as np
|
||||
from logbook import Logger
|
||||
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
from catalyst.protocol import Portfolio, Positions, Position
|
||||
|
||||
log = Logger('ExchangePortfolio')
|
||||
log = Logger('ExchangePortfolio', level=LOG_LEVEL)
|
||||
|
||||
|
||||
class ExchangePortfolio(Portfolio):
|
||||
|
||||
@@ -1,14 +1,14 @@
|
||||
import json
|
||||
import os
|
||||
import pickle
|
||||
import urllib
|
||||
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 +33,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 +41,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):
|
||||
@@ -64,10 +66,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 +154,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 +166,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)
|
||||
|
||||
|
||||
@@ -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,11 @@ from catalyst.gens.sim_engine import (
|
||||
)
|
||||
from logbook import Logger
|
||||
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
from catalyst.exchange.exchange_errors import \
|
||||
MismatchingBaseCurrenciesExchanges
|
||||
|
||||
|
||||
log = Logger('LiveGraphClock')
|
||||
log = Logger('LiveGraphClock', level=LOG_LEVEL)
|
||||
|
||||
|
||||
class LiveGraphClock(object):
|
||||
|
||||
@@ -1,44 +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
|
||||
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.protocol import Account
|
||||
|
||||
log = Logger('Poloniex')
|
||||
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()
|
||||
@@ -124,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)
|
||||
|
||||
@@ -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')
|
||||
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -34,7 +34,9 @@ from catalyst.finance.commission import (
|
||||
from catalyst.finance.cancel_policy import NeverCancel
|
||||
from catalyst.utils.input_validation import expect_types
|
||||
|
||||
log = Logger('Blotter')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = Logger('Blotter', level=LOG_LEVEL)
|
||||
warning_logger = Logger('AlgoWarning')
|
||||
|
||||
|
||||
|
||||
@@ -24,7 +24,9 @@ from catalyst.errors import (
|
||||
TradingControlViolation,
|
||||
)
|
||||
|
||||
log = logbook.Logger('TradingControl')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = logbook.Logger('TradingControl', level=LOG_LEVEL)
|
||||
|
||||
|
||||
class TradingControl(with_metaclass(abc.ABCMeta)):
|
||||
|
||||
@@ -88,7 +88,10 @@ from six import itervalues, iteritems
|
||||
|
||||
import catalyst.protocol as zp
|
||||
|
||||
log = logbook.Logger('Performance')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = logbook.Logger('Performance', level=LOG_LEVEL)
|
||||
|
||||
TRADE_TYPE = zp.DATASOURCE_TYPE.TRADE
|
||||
|
||||
|
||||
|
||||
@@ -40,7 +40,9 @@ import logbook
|
||||
from catalyst.assets import Future, Asset
|
||||
from catalyst.utils.input_validation import expect_types
|
||||
|
||||
log = logbook.Logger('Performance')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = logbook.Logger('Performance', level=LOG_LEVEL)
|
||||
|
||||
|
||||
class Position(object):
|
||||
|
||||
@@ -32,7 +32,9 @@ from catalyst.assets import (
|
||||
)
|
||||
from . position import positiondict
|
||||
|
||||
log = logbook.Logger('Performance')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = logbook.Logger('Performance', level=LOG_LEVEL)
|
||||
|
||||
|
||||
PositionStats = namedtuple('PositionStats',
|
||||
|
||||
@@ -70,7 +70,9 @@ import catalyst.finance.risk as risk
|
||||
|
||||
from . position_tracker import PositionTracker
|
||||
|
||||
log = logbook.Logger('Performance')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = logbook.Logger('Performance', level=LOG_LEVEL)
|
||||
|
||||
|
||||
class PerformanceTracker(object):
|
||||
|
||||
@@ -38,7 +38,9 @@ from empyrical import (
|
||||
sortino_ratio,
|
||||
)
|
||||
|
||||
log = logbook.Logger('Risk Cumulative')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = logbook.Logger('Risk Cumulative', level=LOG_LEVEL)
|
||||
|
||||
|
||||
choose_treasury = functools.partial(choose_treasury, lambda *args: '10year',
|
||||
|
||||
@@ -36,7 +36,9 @@ from empyrical import (
|
||||
sortino_ratio
|
||||
)
|
||||
|
||||
log = logbook.Logger('Risk Period')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = logbook.Logger('Risk Period', level=LOG_LEVEL)
|
||||
|
||||
choose_treasury = functools.partial(risk.choose_treasury,
|
||||
risk.select_treasury_duration)
|
||||
|
||||
@@ -63,7 +63,9 @@ from dateutil.relativedelta import relativedelta
|
||||
|
||||
from . period import RiskMetricsPeriod
|
||||
|
||||
log = logbook.Logger('Risk Report')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = logbook.Logger('Risk Report', level=LOG_LEVEL)
|
||||
|
||||
|
||||
class RiskReport(object):
|
||||
|
||||
@@ -61,7 +61,9 @@ Risk Report
|
||||
import logbook
|
||||
import numpy as np
|
||||
|
||||
log = logbook.Logger('Risk')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = logbook.Logger('Risk', level=LOG_LEVEL)
|
||||
|
||||
|
||||
TREASURY_DURATIONS = [
|
||||
|
||||
@@ -26,7 +26,9 @@ from catalyst.data.loader import load_market_data
|
||||
from catalyst.utils.calendars import get_calendar
|
||||
from catalyst.utils.memoize import remember_last
|
||||
|
||||
log = logbook.Logger('Trading')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = logbook.Logger('Trading', level=LOG_LEVEL)
|
||||
|
||||
|
||||
DEFAULT_CAPITAL_BASE = 1e5
|
||||
|
||||
@@ -27,7 +27,9 @@ from catalyst.gens.sim_engine import (
|
||||
BEFORE_TRADING_START_BAR
|
||||
)
|
||||
|
||||
log = Logger('Trade Simulation')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = Logger('Trade Simulation', level=LOG_LEVEL)
|
||||
|
||||
|
||||
class AlgorithmSimulator(object):
|
||||
|
||||
@@ -23,7 +23,9 @@ from catalyst.protocol import (
|
||||
)
|
||||
from catalyst.assets import Equity
|
||||
|
||||
logger = Logger('Requests Source Logger')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
logger = Logger('Requests Source Logger', level=LOG_LEVEL)
|
||||
|
||||
|
||||
def roll_dts_to_midnight(dts, trading_day):
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -36,14 +36,16 @@ 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
|
||||
|
||||
log = Logger('run_algo')
|
||||
from catalyst.constants import LOG_LEVEL
|
||||
|
||||
log = Logger('run_algo', level=LOG_LEVEL)
|
||||
|
||||
|
||||
class _RunAlgoError(click.ClickException, ValueError):
|
||||
@@ -164,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'],
|
||||
@@ -235,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)
|
||||
|
||||
|
||||
@@ -52,7 +52,7 @@ My first algorithm
|
||||
~~~~~~~~~~~~~~~~~~
|
||||
|
||||
Lets take a look at a very simple algorithm from the ``examples``
|
||||
directory, ``buy_btc.py``:
|
||||
directory: `buy_btc_simple.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/buy_btc_simple.py>`_:
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
@@ -225,16 +225,16 @@ Thus, to execute our algorithm from above and save the results to
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
catalyst run -f buy_btc_simple.py -x bitfinex --start 2016-1-1 --end 2016-9-29 -o buy_simple_btc_out.pickle
|
||||
catalyst run -f buy_btc_simple.py -x bitfinex --start 2016-1-1 --end 2017-9-30 -o buy_btc_simple_out.pickle
|
||||
|
||||
|
||||
..
|
||||
.. parsed-literal
|
||||
.. parsed-literal::
|
||||
|
||||
.. AAPL
|
||||
.. [2015-11-04 22:45:32.820166] INFO: Performance: Simulated 3521 trading days out of 3521.
|
||||
.. [2015-11-04 22:45:32.820314] INFO: Performance: first open: 2000-01-03 14:31:00+00:00
|
||||
.. [2015-11-04 22:45:32.820401] INFO: Performance: last close: 2013-12-31 21:00:00+00:00
|
||||
INFO: run_algo: running algo in backtest mode
|
||||
INFO: exchange_algorithm: initialized trading algorithm in backtest mode
|
||||
INFO: Performance: Simulated 639 trading days out of 639.
|
||||
INFO: Performance: first open: 2016-01-01 00:00:00+00:00
|
||||
INFO: Performance: last close: 2017-09-30 23:59:00+00:00
|
||||
|
||||
|
||||
``run`` first calls the ``initialize()`` function, and then
|
||||
@@ -255,7 +255,7 @@ slippage model that ``catalyst`` uses).
|
||||
|
||||
Let's take a quick look at the performance ``DataFrame``. For this, we
|
||||
use ``pandas`` from inside the IPython Notebook and print the first ten
|
||||
rows. Note that ``catalyst`` makes heavy usage of
|
||||
rows. and print the first ten rows. Note that ``catalyst`` makes heavy usage of
|
||||
`pandas <http://pandas.pydata.org/>`_, especially for data input and
|
||||
outputting so it's worth spending some time to learn it.
|
||||
|
||||
@@ -265,17 +265,200 @@ outputting so it's worth spending some time to learn it.
|
||||
perf = pd.read_pickle('buy_btc_simple_out.pickle') # read in perf DataFrame
|
||||
perf.head()
|
||||
|
||||
.. raw:: html
|
||||
|
||||
<div style="max-height:1000px;max-width:1500px;overflow:auto;">
|
||||
<table border="1" class="dataframe">
|
||||
<thead>
|
||||
<tr style="text-align: right;">
|
||||
<th></th>
|
||||
<th>algo_volatility</th>
|
||||
<th>algorithm_period_return</th>
|
||||
<th>alpha</th>
|
||||
<th>benchmark_period_return</th>
|
||||
<th>benchmark_volatility</th>
|
||||
<th>beta</th>
|
||||
<th>btc</th>
|
||||
<th>capital_used</th>
|
||||
<th>ending_cash</th>
|
||||
<th>ending_exposure</th>
|
||||
<th>...</th>
|
||||
<th>short_exposure</th>
|
||||
<th>short_value</th>
|
||||
<th>shorts_count</th>
|
||||
<th>sortino</th>
|
||||
<th>starting_cash</th>
|
||||
<th>starting_exposure</th>
|
||||
<th>starting_value</th>
|
||||
<th>trading_days</th>
|
||||
<th>transactions</th>
|
||||
<th>treasury_period_return</th>
|
||||
</tr>
|
||||
</thead>
|
||||
<tbody>
|
||||
<tr>
|
||||
<th>2016-01-01 23:59:00+00:00</th>
|
||||
<td>NaN</td>
|
||||
<td>0.000000e+00</td>
|
||||
<td>NaN</td>
|
||||
<td>-0.010937</td>
|
||||
<td>NaN</td>
|
||||
<td>NaN</td>
|
||||
<td>433.979999</td>
|
||||
<td>0.000000</td>
|
||||
<td>1.000000e+07</td>
|
||||
<td>0.00</td>
|
||||
<td>...</td>
|
||||
<td>0</td>
|
||||
<td>0</td>
|
||||
<td>0</td>
|
||||
<td>NaN</td>
|
||||
<td>1.000000e+07</td>
|
||||
<td>0.00</td>
|
||||
<td>0.00</td>
|
||||
<td>1</td>
|
||||
<td>[]</td>
|
||||
<td>0.0227</td>
|
||||
</tr>
|
||||
<tr>
|
||||
<th>2016-01-02 23:59:00+00:00</th>
|
||||
<td>0.000011</td>
|
||||
<td>-9.536708e-07</td>
|
||||
<td>-0.000170</td>
|
||||
<td>-0.006480</td>
|
||||
<td>0.173338</td>
|
||||
<td>-0.000062</td>
|
||||
<td>432.700000</td>
|
||||
<td>-442.236708</td>
|
||||
<td>9.999558e+06</td>
|
||||
<td>432.70</td>
|
||||
<td>...</td>
|
||||
<td>0</td>
|
||||
<td>0</td>
|
||||
<td>0</td>
|
||||
<td>-11.224972</td>
|
||||
<td>1.000000e+07</td>
|
||||
<td>0.00</td>
|
||||
<td>0.00</td>
|
||||
<td>2</td>
|
||||
<td>[{u'order_id': u'7869f7828fa140328eb40477bb7de...</td>
|
||||
<td>0.0227</td>
|
||||
</tr>
|
||||
<tr>
|
||||
<th>2016-01-03 23:59:00+00:00</th>
|
||||
<td>0.000011</td>
|
||||
<td>-2.328842e-06</td>
|
||||
<td>-0.000176</td>
|
||||
<td>-0.026512</td>
|
||||
<td>0.197857</td>
|
||||
<td>0.000009</td>
|
||||
<td>428.390000</td>
|
||||
<td>-437.831716</td>
|
||||
<td>9.999120e+06</td>
|
||||
<td>856.78</td>
|
||||
<td>...</td>
|
||||
<td>0</td>
|
||||
<td>0</td>
|
||||
<td>0</td>
|
||||
<td>-12.754262</td>
|
||||
<td>9.999558e+06</td>
|
||||
<td>432.70</td>
|
||||
<td>432.70</td>
|
||||
<td>3</td>
|
||||
<td>[{u'order_id': u'be62ff77760c4599abaac43be9cc9...</td>
|
||||
<td>0.0227</td>
|
||||
</tr>
|
||||
<tr>
|
||||
<th>2016-01-04 23:59:00+00:00</th>
|
||||
<td>0.000011</td>
|
||||
<td>-2.380954e-06</td>
|
||||
<td>-0.000139</td>
|
||||
<td>-0.008640</td>
|
||||
<td>0.269790</td>
|
||||
<td>0.000020</td>
|
||||
<td>432.900000</td>
|
||||
<td>-442.441116</td>
|
||||
<td>9.998677e+06</td>
|
||||
<td>1298.70</td>
|
||||
<td>...</td>
|
||||
<td>0</td>
|
||||
<td>0</td>
|
||||
<td>0</td>
|
||||
<td>-11.287205</td>
|
||||
<td>9.999120e+06</td>
|
||||
<td>856.78</td>
|
||||
<td>856.78</td>
|
||||
<td>4</td>
|
||||
<td>[{u'order_id': u'd6dca79513214346a646079213526...</td>
|
||||
<td>0.0224</td>
|
||||
</tr>
|
||||
<tr>
|
||||
<th>2016-01-05 23:59:00+00:00</th>
|
||||
<td>0.000011</td>
|
||||
<td>-3.650729e-06</td>
|
||||
<td>-0.000158</td>
|
||||
<td>-0.021426</td>
|
||||
<td>0.245989</td>
|
||||
<td>0.000024</td>
|
||||
<td>431.840000</td>
|
||||
<td>-441.357754</td>
|
||||
<td>9.998236e+06</td>
|
||||
<td>1727.36</td>
|
||||
<td>...</td>
|
||||
<td>0</td>
|
||||
<td>0</td>
|
||||
<td>0</td>
|
||||
<td>-12.333847</td>
|
||||
<td>9.998677e+06</td>
|
||||
<td>1298.70</td>
|
||||
<td>1298.70</td>
|
||||
<td>5</td>
|
||||
<td>[{u'order_id': u'505275d6646a41f3856b22b16678d...</td>
|
||||
<td>0.0225</td>
|
||||
</tr>
|
||||
</tbody>
|
||||
</table>
|
||||
</div>
|
||||
|
||||
|
|
||||
There is a row for each trading day, starting on the first day of our
|
||||
simulation Jan 1st, 2016. In the columns you can find various
|
||||
information about the state of your algorithm. The very first column
|
||||
information about the state of your algorithm. The column
|
||||
``btc`` was placed there by the ``record()`` function mentioned earlier
|
||||
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.
|
||||
|
||||
Our algorithm performance as assessed by the
|
||||
``portfolio_value`` closely matches that of the bitcoin price. This
|
||||
is not surprising as our algorithm only bought bitcoin every chance it got.
|
||||
.. code-block:: python
|
||||
|
||||
%load_ext catalyst
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
%pylab inline
|
||||
figsize(12, 12)
|
||||
import matplotlib.pyplot as plt
|
||||
|
||||
ax1 = plt.subplot(211)
|
||||
perf.portfolio_value.plot(ax=ax1)
|
||||
ax1.set_ylabel('portfolio value')
|
||||
ax2 = plt.subplot(212, sharex=ax1)
|
||||
perf.btc.plot(ax=ax2)
|
||||
ax2.set_ylabel('bitcoin price')
|
||||
|
||||
.. parsed-literal::
|
||||
|
||||
Populating the interactive namespace from numpy and matplotlib
|
||||
|
||||
.. parsed-literal::
|
||||
|
||||
<matplotlib.text.Text at 0x10eaeadd0>
|
||||
|
||||
.. image:: https://s3.amazonaws.com/enigmaco-docs/github.io/buy_btc_simple_graph.png
|
||||
|
||||
Our algorithm performance as assessed by the ``portfolio_value`` closely
|
||||
matches that of the bitcoin price. This is not surprising as our algorithm
|
||||
only bought bitcoin every chance it got.
|
||||
|
||||
|
||||
Access to previous prices using ``history``
|
||||
@@ -305,23 +488,25 @@ a function we use in the ``handle_data()`` section:
|
||||
|
||||
.. code-block:: python
|
||||
|
||||
from catalyst.api import order, record, symbol
|
||||
%%catalyst --start 2016-4-1 --end 2017-9-30 -x bitfinex
|
||||
|
||||
from catalyst.api import order, record, symbol, order_target
|
||||
|
||||
def initialize(context):
|
||||
context.i = 0
|
||||
context.asset = symbol('btc_usd')
|
||||
|
||||
def handle_data(context, data):
|
||||
# Skip first 300 days to get full windows
|
||||
def handle_data(context, data):
|
||||
# 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:
|
||||
@@ -336,6 +521,46 @@ a function we use in the ``handle_data()`` section:
|
||||
short_mavg=short_mavg,
|
||||
long_mavg=long_mavg)
|
||||
|
||||
def analyze(context, perf):
|
||||
import matplotlib.pyplot as plt
|
||||
fig = plt.figure(figsize=(12,12))
|
||||
ax1 = fig.add_subplot(211)
|
||||
perf.portfolio_value.plot(ax=ax1)
|
||||
ax1.set_ylabel('portfolio value in $')
|
||||
|
||||
ax2 = fig.add_subplot(212)
|
||||
perf['btc'].plot(ax=ax2)
|
||||
perf[['short_mavg', 'long_mavg']].plot(ax=ax2)
|
||||
|
||||
perf_trans = perf.ix[[t != [] for t in perf.transactions]]
|
||||
buys = perf_trans.ix[[t[0]['amount'] > 0 for t in perf_trans.transactions]]
|
||||
sells = perf_trans.ix[
|
||||
[t[0]['amount'] < 0 for t in perf_trans.transactions]]
|
||||
ax2.plot(buys.index, perf.short_mavg.ix[buys.index],
|
||||
'^', markersize=10, color='m')
|
||||
ax2.plot(sells.index, perf.short_mavg.ix[sells.index],
|
||||
'v', markersize=10, color='k')
|
||||
ax2.set_ylabel('price in $')
|
||||
plt.legend(loc=0)
|
||||
plt.show()
|
||||
|
||||
Here we are explicitly defining an ``analyze()`` function that gets
|
||||
automatically called once the backtest is done.
|
||||
|
||||
Although it might not be directly apparent, the power of ``history()``
|
||||
(pun intended) can not be under-estimated as most algorithms make use of
|
||||
prior market developments in one form or another. You could easily
|
||||
devise a strategy that trains a classifier with
|
||||
`scikit-learn <http://scikit-learn.org/stable/>`__ which tries to
|
||||
predict future market movements based on past prices (note, that most of
|
||||
the ``scikit-learn`` functions require ``numpy.ndarray``\ s rather than
|
||||
``pandas.DataFrame``\ s, so you can simply pass the underlying
|
||||
``ndarray`` of a ``DataFrame`` via ``.values``).
|
||||
|
||||
We also used the ``order_target()`` function above. This and other
|
||||
functions like it can make order management and portfolio rebalancing
|
||||
much easier.
|
||||
|
||||
|
||||
Conclusions
|
||||
~~~~~~~~~~~
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
# Incompatible with earlier PIP versions
|
||||
pip>=7.1.0
|
||||
# bcolz fails to install if this is not in the build_requires.
|
||||
setuptools>18.0
|
||||
setuptools>36.0
|
||||
|
||||
# Logging
|
||||
Logbook==0.12.5
|
||||
|
||||
@@ -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')
|
||||
@@ -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')
|
||||
|
||||
@@ -7,7 +7,7 @@ from catalyst.exchange.exchange_utils import get_exchange_auth
|
||||
log = Logger('test_bittrex')
|
||||
|
||||
|
||||
class BittrexTestCase(BaseExchangeTestCase):
|
||||
class TestBittrexTestCase(BaseExchangeTestCase):
|
||||
@classmethod
|
||||
def setup(self):
|
||||
print ('creating bittrex object')
|
||||
|
||||
@@ -1,9 +1,10 @@
|
||||
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, \
|
||||
from catalyst.exchange.bundle_utils import get_bcolz_chunk, \
|
||||
get_periods_range
|
||||
from catalyst.exchange.exchange_bcolz import BcolzExchangeBarReader, \
|
||||
BcolzExchangeBarWriter
|
||||
@@ -13,10 +14,10 @@ from catalyst.exchange.exchange_utils import get_exchange_folder
|
||||
from catalyst.exchange.init_utils import get_exchange
|
||||
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'
|
||||
@@ -93,16 +94,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 +298,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 +308,10 @@ class ExchangeBundleTestCase:
|
||||
)
|
||||
|
||||
pass
|
||||
|
||||
def test_hash_symbol(self):
|
||||
symbol = 'etc_btc'
|
||||
sid = int(
|
||||
hashlib.sha256(symbol.encode('utf-8')).hexdigest(), 16
|
||||
) % 10 ** 6
|
||||
pass
|
||||
|
||||
@@ -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
|
||||
@@ -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')
|
||||
|
||||
@@ -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
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user