Compare commits

...
Author SHA1 Message Date
Frederic Fortier 6fc7dd6838 BLD: adjusting unit tests 2018-03-16 15:46:05 -04:00
Victor Grau Serrat 60e07c8a3c BUG: fix sanitize_df to min of int32 2018-03-12 22:56:09 +00:00
Frederic Fortier 437e5c5b80 BLD: fetching trades recursively 2018-03-10 21:37:31 -05:00
Frederic Fortier 8bf639b01a Merge remote-tracking branch 'remotes/origin/develop' into new_exchange_config
# Conflicts:
#	tests/exchange/test_bundle.py
2018-03-09 15:52:07 -05:00
Victor Grau Serrat 8f7d678170 BUG: [mktplace]: fix sanitize_df to handle 1-row DFs 2018-03-09 12:44:12 -07:00
Frederic Fortier 1de084a08f Merge remote-tracking branch 'origin/develop' into develop 2018-03-09 13:21:02 -05:00
Frederic Fortier 93ca4e990d BUG: fixed a dependency issue 2018-03-09 13:20:50 -05:00
Victor Grau Serrat 4694372496 BUG: [mktplace] ingest mismatch dest. folder 2018-03-09 09:56:09 -07:00
Victor Grau Serrat f0606b5ea4 BLD: [mktplace] clean --dataset param optional 2018-03-09 09:50:20 -07:00
AvishaiW b0dce13672 BUG: removed get_open_orders from buy_low_sell_high 2018-03-09 11:01:19 +02:00
Frederic Fortier a99a4f4e85 Merge remote-tracking branch 'remotes/origin/develop' into new_exchange_config 2018-03-08 17:28:41 -05:00
Frederic Fortier fc9837b678 BUG: fixed an issue with extracting bundles 2018-03-08 17:27:25 -05:00
Frederic Fortier a1228a8ea2 BLD: adjusted tests 2018-03-08 17:01:38 -05:00
Frederic Fortier 2b0830bafc Merge remote-tracking branch 'remotes/origin/develop' into new_exchange_config 2018-03-08 16:49:02 -05:00
Frederic Fortier 7ce382d52e BLD: getting config from the ohlcv directory on the server 2018-03-08 15:51:10 -05:00
Frederic Fortier 42aac37f34 BLD: specific parameters in unit test 2018-03-08 13:39:05 -05:00
Victor Grau Serrat 5de89549ee MAINT: [mktplace] ingest --dataset parameter now optional 2018-03-08 10:33:16 -07:00
AvishaiW 586d7f2954 DOC: added warning that its not possible to start & end at a specific time 2018-03-08 18:34:25 +02:00
Victor Grau Serrat 218dc0bafd BUG: marketplace os.rename -> shutil.move 2018-03-07 22:47:27 -07:00
Frederic Fortier 68f92dedd6 BLD: added a coordinator script and fixed an issue with candle calculation 2018-03-07 23:23:59 -05:00
Victor Grau Serrat 89c080fce7 BLD: marketplace: dataset param to subscribe optional, fix "from" tx 2018-03-07 16:05:38 -07:00
Frederic Fortier a30d498abf BLD: added symbol mapping when migrating bundles 2018-03-06 22:12:19 -05:00
Frederic Fortier ebdb9e423e BLD: one more attempt at fixing the dup candle issue 2018-03-06 20:13:59 -05:00
Victor Grau Serrat 4a794aa035 BLD: show catalyst version at runtime 2018-03-05 21:59:18 -07:00
Frederic Fortier 1f0c037d29 BLD: created a script which migrates existing bundles 2018-03-05 23:26:55 -05:00
Frederic Fortier 623ace1cbd BLD: updated CCXT 2018-03-05 23:26:22 -05:00
Frederic Fortier a004825a09 BLD: successful bundle comparison 2018-03-05 22:18:22 -05:00
Frederic Fortier cbc2ed2aaf BLD: fixes small issues when testing 2018-03-05 21:35:38 -05:00
Frederic Fortier 1de881e17f BLD: trying to pinpoint a duplicates issue 2018-03-05 19:41:57 -05:00
Frederic Fortier 6488ff5abe Merge remote-tracking branch 'remotes/origin/develop' into new_exchange_config
# Conflicts:
#	catalyst/exchange/exchange.py
#	catalyst/exchange/exchange_errors.py
2018-03-05 19:17:13 -05:00
AvishaiW cb4668f093 BUG: #243 added a function which reduces open orders amount from calculated target/amount for target orders 2018-03-06 00:05:29 +02:00
lenak25 6092471180 DOC: add some commented TODOs 2018-03-05 20:42:14 +02:00
lenak25 09068a4c37 BLD: adjust the example to Python 3 2018-03-04 18:20:14 +02:00
lenak25 ed406a30ff BLD: fix issue #260 - always request more data to avoid empty bars and always give the exact bar number 2018-03-04 17:45:08 +02:00
AvishaiW aa520d5a8b STY: pep8 change in exchange_blotter 2018-03-04 17:24:22 +02:00
lenak25 e59f46dfd6 BLD: improve periods calculation 2018-03-04 13:47:08 +02:00
Frederic Fortier 0d051a9496 BLD: made some adjustments to the client when testing the Binance data bundle 2018-03-03 23:09:59 -05:00
Frederic Fortier 964c90176b BLD: fixed float points in config files 2018-03-02 23:03:14 -05:00
Victor Grau Serrat 5197ab6cc2 BUG: isolated Python3 depedency to the marketplace 2018-03-02 12:50:46 -07:00
Victor Grau Serrat d7b6cb8490 BUG: marketplace typo 2018-03-02 12:18:48 -07:00
Victor Grau Serrat fa0e9332bf MAINT: CLI info on marketplace cmds 2018-03-02 11:43:59 -07:00
VictorandGitHub d9d6a4e52d Merge pull request #257 from mattbornski/master
fix incompatibility with web3==4.0.0b11 from prior web3 versions
2018-03-02 11:38:52 -07:00
VictorandGitHub b4e5b699bd BUG: fix2 incompatibility with web3==4.0.0b11 2018-03-02 11:37:22 -07:00
VictorandGitHub f990ecf14d BUG: fix incompatibility with web3==4.0.0b11 2018-03-02 11:34:18 -07:00
Frederic Fortier 40cfc65e02 BUG: fixed python2 syntax 2018-03-01 23:33:17 -05:00
Frederic Fortier abc48494c2 Merge branch 'develop' into new_exchange_config
# Conflicts:
#	catalyst/exchange/utils/exchange_utils.py
2018-03-01 21:48:30 -05:00
Frederic Fortier 73faa87269 Merge remote-tracking branch 'remotes/origin/develop' into new_exchange_config
# Conflicts:
#	catalyst/exchange/exchange.py
2018-03-01 21:46:30 -05:00
lenak25 5d3a1c2f8b BLD: cosmetics 2018-03-02 00:39:28 +02:00
lenak25 3159ec7dc2 BLD: fix periods calculations and bundle unit-test 2018-03-02 00:36:13 +02:00
lenak25 8be0626fc9 BLD:refine unit-test 2018-03-01 18:16:57 +02:00
lenak25 868f17fd9d BLD:flake fixes 2018-03-01 16:39:42 +02:00
lenak25 b95cf465fc BLD:updating the forward fill to set volume to zero and others values to the previous close value 2018-03-01 16:09:09 +02:00
AvishaiW c4b10bae39 DOC: added- creating a virtual env for 3.6 2018-03-01 09:47:18 +02:00
AvishaiW 6c4f7afaea BUG: fixed removing files- check the path, not the file 2018-03-01 09:42:07 +02:00
Victor Grau Serrat b85219d5b4 MAINT: undoing last 2 unwanted commits 2018-02-28 18:06:36 -07:00
Frederic Fortier c2f71cf852 BLD: testing adjusted scripts 2018-02-28 19:56:30 -05:00
AvishaiW 082b342d02 Merge remote-tracking branch 'origin/develop' into cloud_conn 2018-03-01 01:11:46 +02:00
AvishaiW 37c1057ab6 BLD: added a cmd for running on the cloud (WIP) 2018-03-01 01:09:35 +02:00
Avishai WeingartenandGitHub 46e1a87a3d DOC: added troubleshooting for python3 2018-02-27 18:47:43 +02:00
Avishai WeingartenandGitHub cfafafb8fc BUG #252 fixed utc time and file erased 2018-02-27 09:47:02 +02:00
Matt Bornski 497212383a Python 3 returns bytes, the parsing functions are looking for strings 2018-02-26 15:44:32 -08:00
AvishaiW e5870ea60a BUG: fix #252 #253 and split state into paper and live 2018-02-26 09:21:00 +02:00
AvishaiW 8587fee0ce BUG: revert previous changes #249 2018-02-25 11:29:26 +02:00
Victor Grau Serrat 388535b09c BUG: reverts changed introduced in 00f232e2d7 2018-02-22 22:14:03 -07:00
Victor Grau Serrat 25e9f0f58f BUG: reverts changed introduced in 00f232e2d7 2018-02-22 22:09:51 -07:00
AvishaiW b4bd557273 BUG: fixes for issues #204 #237
-modified parameters for cancel_orders
-update portfolio after any change in
 the orders before sync
2018-02-23 00:38:51 +02:00
Victor Grau Serrat 2577b53518 MAINT: conda environment updates 2018-02-22 12:55:19 -07:00
Victor Grau Serrat 8fe3ab344e MAINT: conda environment updates 2018-02-22 12:54:36 -07:00
Victor bfd7e4b2dd Update python3.6-environment.yml 2018-02-22 12:54:36 -07:00
lenak25 50310576f9 BUG:fix an issue with wrong timestamps seen at tests.exchange.test_suites.test_suite_bundle.TestSuiteBundle#test_validate_bundles (which issue #230 uncovered) 2018-02-22 17:42:14 +02:00
lenak25 fea2ed104e BUG: fix issue #236: handle properly empty candles received from exchanges 2018-02-22 16:50:48 +02:00
embaral 127878413e DOC: added an option "catalyst live --help" to the documentation. 2018-02-22 14:45:14 +02:00
embaral 6a6ccf5595 Merge remote-tracking branch 'origin/develop' into develop 2018-02-22 14:37:46 +02:00
embaral d40585f56e DOC: added an option "catalyst live --help" to the documentation. 2018-02-22 14:34:00 +02:00
Victor Grau Serrat 4337abd60a DOC: linking example_algo to their sources 2018-02-21 15:51:34 -07:00
AvishaiW 92e0a7bb88 Merge branch 'develop' of https://github.com/enigmampc/catalyst into develop 2018-02-21 20:45:39 +02:00
AvishaiW 20f8a75f4a BUG: for issue #237, update positions before checking balances 2018-02-21 20:42:39 +02:00
Victor Grau Serrat ec5fdecf91 DOC: marketplace code examples 2018-02-16 11:50:16 -07:00
VictorandGitHub 9956b5462d Update python3.6-environment.yml 2018-02-14 09:27:53 -07:00
Victor Grau Serrat 46f34d64a0 MAINT: conda env for Python3 2018-02-13 12:06:06 -07:00
Victor Grau Serrat bc8bf6941d MAINT: contract+abi pointing to master, not develop 2018-02-09 16:08:22 -08:00
Frederic Fortier 49b6792399 Merge branch 'develop' 2018-02-09 12:01:57 -05:00
Victor Grau Serrat 3c4c6c3dfd DOC: small edits, eliminating sphinx warnings 2018-02-08 22:10:57 -07:00
Frederic Fortier 403d7f9c29 Merge branch 'develop' 2018-02-08 17:32:47 -05:00
37 changed files with 939 additions and 1324 deletions
+16 -15
View File
@@ -767,12 +767,18 @@ def bundles():
@main.group() @main.group()
@click.pass_context @click.pass_context
def marketplace(ctx): def marketplace(ctx):
"""Access the Enigma Data Marketplace to:\n
- Register and Publish new datasets (seller-side)\n
- Subscribe and Ingest premium datasets (buyer-side)\n
"""
pass pass
@marketplace.command() @marketplace.command()
@click.pass_context @click.pass_context
def ls(ctx): def ls(ctx):
"""List all available datasets.
"""
click.echo('Listing of available data sources on the marketplace:', click.echo('Listing of available data sources on the marketplace:',
sys.stdout) sys.stdout)
marketplace = Marketplace() marketplace = Marketplace()
@@ -787,10 +793,8 @@ def ls(ctx):
) )
@click.pass_context @click.pass_context
def subscribe(ctx, dataset): def subscribe(ctx, dataset):
if dataset is None: """Subscribe to an exisiting dataset.
ctx.fail("must specify a dataset to subscribe to with '--dataset'\n" """
"List available dataset on the marketplace with "
"'catalyst marketplace ls'")
marketplace = Marketplace() marketplace = Marketplace()
marketplace.subscribe(dataset) marketplace.subscribe(dataset)
@@ -825,11 +829,8 @@ def subscribe(ctx, dataset):
) )
@click.pass_context @click.pass_context
def ingest(ctx, dataset, data_frequency, start, end): def ingest(ctx, dataset, data_frequency, start, end):
if dataset is None: """Ingest a dataset (requires subscription).
ctx.fail("must specify a dataset to clean with '--dataset'\n" """
"List available dataset on the marketplace with "
"'catalyst marketplace ls'")
click.echo('Ingesting data: {}'.format(dataset), sys.stdout)
marketplace = Marketplace() marketplace = Marketplace()
marketplace.ingest(dataset, data_frequency, start, end) marketplace.ingest(dataset, data_frequency, start, end)
@@ -842,19 +843,17 @@ def ingest(ctx, dataset, data_frequency, start, end):
) )
@click.pass_context @click.pass_context
def clean(ctx, dataset): def clean(ctx, dataset):
if dataset is None: """Clean/Remove local data for a given dataset.
ctx.fail("must specify a dataset to ingest with '--dataset'\n" """
"List available dataset on the marketplace with "
"'catalyst marketplace ls'")
click.echo('Cleaning data source: {}'.format(dataset), sys.stdout)
marketplace = Marketplace() marketplace = Marketplace()
marketplace.clean(dataset) marketplace.clean(dataset)
click.echo('Done', sys.stdout)
@marketplace.command() @marketplace.command()
@click.pass_context @click.pass_context
def register(ctx): def register(ctx):
"""Register a new dataset.
"""
marketplace = Marketplace() marketplace = Marketplace()
marketplace.register() marketplace.register()
@@ -878,6 +877,8 @@ def register(ctx):
) )
@click.pass_context @click.pass_context
def publish(ctx, dataset, datadir, watch): def publish(ctx, dataset, datadir, watch):
"""Publish data for a registered dataset.
"""
marketplace = Marketplace() marketplace = Marketplace()
if dataset is None: if dataset is None:
ctx.fail("must specify a dataset to publish data for " ctx.fail("must specify a dataset to publish data for "
+8 -6
View File
@@ -11,8 +11,10 @@ LOG_LEVEL = int(os.environ.get('CATALYST_LOG_LEVEL', logbook.INFO))
SYMBOLS_URL = 'https://s3.amazonaws.com/enigmaco/catalyst-exchanges/' \ SYMBOLS_URL = 'https://s3.amazonaws.com/enigmaco/catalyst-exchanges/' \
'{exchange}/symbols.json' '{exchange}/symbols.json'
EXCHANGE_CONFIG_URL = 'https://s3.amazonaws.com/enigmaco/exchanges/' \ EXCHANGE_CONFIG_URL = 'https://s3.amazonaws.com/enigmaco/ohlcv/' \
'{exchange}/config.json' '{exchange}/config.json'
BUNDLE_URL = 'https://s3.amazonaws.com/enigmaco/ohlcv/' \
'{exchange}/{data_frequency}/{name}.tar.gz'
DATE_TIME_FORMAT = '%Y-%m-%d %H:%M' DATE_TIME_FORMAT = '%Y-%m-%d %H:%M'
DATE_FORMAT = '%Y-%m-%d' DATE_FORMAT = '%Y-%m-%d'
@@ -28,20 +30,20 @@ AUTH_SERVER = 'https://data.enigma.co'
# TODO: switch to mainnet # TODO: switch to mainnet
ETH_REMOTE_NODE = 'https://ropsten.infura.io/' ETH_REMOTE_NODE = 'https://ropsten.infura.io/'
# TODO: move to MASTER branch on github
MARKETPLACE_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/' \ MARKETPLACE_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/develop/catalyst/marketplace/' \ 'catalyst/master/catalyst/marketplace/' \
'contract_marketplace_address.txt' 'contract_marketplace_address.txt'
MARKETPLACE_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \ MARKETPLACE_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/develop/catalyst/marketplace/' \ 'catalyst/master/catalyst/marketplace/' \
'contract_marketplace_abi.json' 'contract_marketplace_abi.json'
# TODO: switch to mainnet # TODO: switch to mainnet
ENIGMA_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/catalyst/' \ ENIGMA_CONTRACT = 'https://raw.githubusercontent.com/enigmampc/catalyst/' \
'develop/catalyst/marketplace/' \ 'master/catalyst/marketplace/' \
'contract_enigma_address.txt' 'contract_enigma_address.txt'
ENIGMA_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \ ENIGMA_CONTRACT_ABI = 'https://raw.githubusercontent.com/enigmampc/' \
'catalyst/develop/catalyst/marketplace/' \ 'catalyst/master/catalyst/marketplace/' \
'contract_enigma_abi.json' 'contract_enigma_abi.json'
-1
View File
@@ -7,7 +7,6 @@ from catalyst.api import (
order_target_percent, order_target_percent,
symbol, symbol,
record, record,
get_open_orders,
) )
from catalyst.exchange.utils.stats_utils import get_pretty_stats from catalyst.exchange.utils.stats_utils import get_pretty_stats
from catalyst.utils.run_algo import run_algorithm from catalyst.utils.run_algo import run_algorithm
+16 -28
View File
@@ -4,8 +4,7 @@ import pandas as pd
from logbook import Logger from logbook import Logger
from catalyst import run_algorithm from catalyst import run_algorithm
from catalyst.api import (record, symbol, order_target_percent, from catalyst.api import (record, symbol, order_target_percent,)
get_open_orders)
from catalyst.exchange.utils.stats_utils import extract_transactions from catalyst.exchange.utils.stats_utils import extract_transactions
NAMESPACE = 'dual_moving_average' NAMESPACE = 'dual_moving_average'
@@ -20,8 +19,8 @@ def initialize(context):
def handle_data(context, data): def handle_data(context, data):
# define the windows for the moving averages # define the windows for the moving averages
short_window = 2 short_window = 50
long_window = 3 long_window = 200
# Skip as many bars as long_window to properly compute the average # Skip as many bars as long_window to properly compute the average
context.i += 1 context.i += 1
@@ -63,7 +62,7 @@ def handle_data(context, data):
# Since we are using limit orders, some orders may not execute immediately # Since we are using limit orders, some orders may not execute immediately
# we wait until all orders are executed before considering more trades. # we wait until all orders are executed before considering more trades.
orders = get_open_orders(context.asset) orders = context.blotter.open_orders
if len(orders) > 0: if len(orders) > 0:
return return
@@ -150,27 +149,16 @@ def analyze(context, perf):
if __name__ == '__main__': if __name__ == '__main__':
run_algorithm( run_algorithm(
capital_base=1000, capital_base=1000,
data_frequency='minute', data_frequency='minute',
initialize=initialize, initialize=initialize,
handle_data=handle_data, handle_data=handle_data,
analyze=analyze, analyze=analyze,
exchange_name='bitfinex', exchange_name='bitfinex',
algo_namespace=NAMESPACE, algo_namespace=NAMESPACE,
base_currency='usd', base_currency='usd',
simulate_orders=True, start=pd.to_datetime('2017-9-22', utc=True),
live=True, end=pd.to_datetime('2017-9-23', utc=True),
) )
# run_algorithm(
# capital_base=1000,
# data_frequency='minute',
# initialize=initialize,
# handle_data=handle_data,
# analyze=analyze,
# exchange_name='bitfinex',
# algo_namespace=NAMESPACE,
# base_currency='usd',
# start=pd.to_datetime('2017-9-22', utc=True),
# end=pd.to_datetime('2017-9-23', utc=True),
# )
+1 -1
View File
@@ -66,7 +66,7 @@ def handle_data(context, data):
# Define portfolio optimization parameters # Define portfolio optimization parameters
n_portfolios = 50000 n_portfolios = 50000
results_array = np.zeros((3 + context.nassets, n_portfolios)) results_array = np.zeros((3 + context.nassets, n_portfolios))
for p in xrange(n_portfolios): for p in range(n_portfolios):
weights = np.random.random(context.nassets) weights = np.random.random(context.nassets)
weights /= np.sum(weights) weights /= np.sum(weights)
w = np.asmatrix(weights) w = np.asmatrix(weights)
+52 -23
View File
@@ -232,6 +232,21 @@ class CCXT(Exchange):
return frequencies return frequencies
def substitute_currency_code(self, currency, source='catalyst'):
if source == 'catalyst':
currency = currency.upper()
key = self.api.common_currency_code(currency).lower()
self._common_symbols[key] = currency.lower()
return key
else:
if currency in self._common_symbols:
return self._common_symbols[currency]
else:
return currency.lower()
def get_symbol(self, asset_or_symbol, source='catalyst'): def get_symbol(self, asset_or_symbol, source='catalyst'):
""" """
The CCXT symbol. The CCXT symbol.
@@ -386,7 +401,7 @@ class CCXT(Exchange):
) )
def get_candles(self, freq, assets, bar_count=1, start_dt=None, def get_candles(self, freq, assets, bar_count=1, start_dt=None,
end_dt=None): end_dt=None, floor_dates=True):
is_single = (isinstance(assets, TradingPair)) is_single = (isinstance(assets, TradingPair))
if is_single: if is_single:
assets = [assets] assets = [assets]
@@ -433,16 +448,20 @@ class CCXT(Exchange):
candles[asset] = [] candles[asset] = []
for ohlcv in ohlcvs: for ohlcv in ohlcvs:
candles[asset].append(dict( dt = pd.to_datetime(ohlcv[0], unit='ms', utc=True)
last_traded=pd.to_datetime( if floor_dates:
ohlcv[0], unit='ms', utc=True dt = dt.floor('1T')
),
open=ohlcv[1], candles[asset].append(
high=ohlcv[2], dict(
low=ohlcv[3], last_traded=dt,
close=ohlcv[4], open=ohlcv[1],
volume=ohlcv[5] high=ohlcv[2],
)) low=ohlcv[3],
close=ohlcv[4],
volume=ohlcv[5],
)
)
candles[asset] = sorted( candles[asset] = sorted(
candles[asset], key=lambda c: c['last_traded'] candles[asset], key=lambda c: c['last_traded']
) )
@@ -865,7 +884,8 @@ class CCXT(Exchange):
) )
raise ExchangeRequestError(error=e) raise ExchangeRequestError(error=e)
def cancel_order(self, order_param, asset_or_symbol=None): def cancel_order(self, order_param,
asset_or_symbol=None, params={}):
order_id = order_param.id \ order_id = order_param.id \
if isinstance(order_param, Order) else order_param if isinstance(order_param, Order) else order_param
@@ -877,7 +897,8 @@ class CCXT(Exchange):
try: try:
symbol = self.get_symbol(asset_or_symbol) \ symbol = self.get_symbol(asset_or_symbol) \
if asset_or_symbol is not None else None if asset_or_symbol is not None else None
self.api.cancel_order(id=order_id, symbol=symbol) self.api.cancel_order(id=order_id,
symbol=symbol, params=params)
except (ExchangeError, NetworkError) as e: except (ExchangeError, NetworkError) as e:
log.warn( log.warn(
@@ -995,19 +1016,27 @@ class CCXT(Exchange):
return result return result
def get_trades(self, asset, my_trades=True, start_dt=None, limit=100): def get_trades(self, asset, my_trades=True, start_dt=None, limit=100):
if not my_trades:
raise NotImplemented(
'get_trades only supports "my trades"'
)
# TODO: is it possible to sort this? Limit is useless otherwise. # TODO: is it possible to sort this? Limit is useless otherwise.
ccxt_symbol = self.get_symbol(asset) ccxt_symbol = self.get_symbol(asset)
if start_dt:
delta = start_dt - get_epoch()
since = int(delta.total_seconds()) * 1000
else:
since = None
try: try:
trades = self.api.fetch_my_trades( if my_trades:
symbol=ccxt_symbol, trades = self.api.fetch_my_trades(
since=start_dt, symbol=ccxt_symbol,
limit=limit, since=since,
) limit=limit,
)
else:
trades = self.api.fetch_trades(
symbol=ccxt_symbol,
since=since,
limit=limit,
)
except (ExchangeError, NetworkError) as e: except (ExchangeError, NetworkError) as e:
log.warn( log.warn(
'unable to fetch trades {} / {}: {}'.format( 'unable to fetch trades {} / {}: {}'.format(
+43 -30
View File
@@ -13,7 +13,8 @@ from catalyst.exchange.exchange_bundle import ExchangeBundle
from catalyst.exchange.exchange_errors import MismatchingBaseCurrencies, \ from catalyst.exchange.exchange_errors import MismatchingBaseCurrencies, \
SymbolNotFoundOnExchange, \ SymbolNotFoundOnExchange, \
PricingDataNotLoadedError, \ PricingDataNotLoadedError, \
NoDataAvailableOnExchange, NoValueForField, LastCandleTooEarlyError, \ NoDataAvailableOnExchange, NoValueForField, \
NoCandlesReceivedFromExchange, \
TickerNotFoundError, NotEnoughCashError TickerNotFoundError, NotEnoughCashError
from catalyst.exchange.utils.datetime_utils import get_delta, \ from catalyst.exchange.utils.datetime_utils import get_delta, \
get_periods_range, \ get_periods_range, \
@@ -256,9 +257,10 @@ class Exchange:
elif data_frequency is not None: elif data_frequency is not None:
applies = ( applies = (
( (
data_frequency == 'minute' and a.end_minute is not None) data_frequency == 'minute' and a.end_minute is not None
or ( ) or (
data_frequency == 'daily' and a.end_daily is not None) data_frequency == 'daily' and a.end_daily is not None
)
) )
else: else:
@@ -484,44 +486,52 @@ class Exchange:
freq, candle_size, unit, data_frequency = get_frequency( freq, candle_size, unit, data_frequency = get_frequency(
frequency, data_frequency, supported_freqs=['T', 'D', 'H'] frequency, data_frequency, supported_freqs=['T', 'D', 'H']
) )
# we want to avoid receiving empty candles
# so we request more than needed
# TODO: consider defining a const per asset
# and/or some retry mechanism (in each iteration request more data)
requested_bar_count = bar_count + 30
# The get_history method supports multiple asset # The get_history method supports multiple asset
candles = self.get_candles( candles = self.get_candles(
freq=freq, freq=freq,
assets=assets, assets=assets,
bar_count=bar_count, bar_count=requested_bar_count,
end_dt=end_dt if not is_current else None, end_dt=end_dt if not is_current else None,
) )
series = dict() # candles sanity check - verify no empty candles were received:
for asset in candles: for asset in candles:
first_candle = candles[asset][0] if not candles[asset]:
asset_series = self.get_series_from_candles( raise NoCandlesReceivedFromExchange(
candles=candles[asset], bar_count=requested_bar_count,
start_dt=first_candle['last_traded'], end_dt=end_dt,
end_dt=end_dt, asset=asset,
data_frequency=frequency, exchange=self.name)
field=field,
)
delta_candle_size = candle_size * 60 if unit == 'H' else candle_size series = get_candles_df(candles=candles,
# Checking to make sure that the dates match field=field,
delta = get_delta(delta_candle_size, data_frequency) freq=frequency,
adj_end_dt = end_dt - delta bar_count=requested_bar_count,
last_traded = asset_series.index[-1] end_dt=end_dt)
if last_traded < adj_end_dt: # TODO: consider how to approach this edge case
raise LastCandleTooEarlyError( # delta_candle_size = candle_size * 60 if unit == 'H' else candle_size
last_traded=last_traded, # Checking to make sure that the dates match
end_dt=adj_end_dt, # delta = get_delta(delta_candle_size, data_frequency)
exchange=self.name, # adj_end_dt = end_dt - delta
) # last_traded = asset_series.index[-1]
# if last_traded < adj_end_dt:
series[asset] = asset_series # raise LastCandleTooEarlyError(
# last_traded=last_traded,
# end_dt=adj_end_dt,
# exchange=self.name,
# )
df = pd.DataFrame(series) df = pd.DataFrame(series)
df.dropna(inplace=True) df.dropna(inplace=True)
return df return df.tail(bar_count)
def get_history_window_with_bundle(self, def get_history_window_with_bundle(self,
assets, assets,
@@ -569,7 +579,8 @@ class Exchange:
A dataframe containing the requested data. A dataframe containing the requested data.
""" """
# TODO: this function needs some work, we're currently using it just for benchmark data # TODO: this function needs some work,
# we're currently using it just for benchmark data
freq, candle_size, unit, data_frequency = get_frequency( freq, candle_size, unit, data_frequency = get_frequency(
frequency, data_frequency frequency, data_frequency
) )
@@ -906,7 +917,8 @@ class Exchange:
""" """
@abstractmethod @abstractmethod
def cancel_order(self, order_param, symbol_or_asset=None): def cancel_order(self, order_param,
symbol_or_asset=None, params={}):
"""Cancel an open order. """Cancel an open order.
Parameters Parameters
@@ -915,6 +927,7 @@ class Exchange:
The order_id or order object to cancel. The order_id or order object to cancel.
symbol_or_asset: str|TradingPair symbol_or_asset: str|TradingPair
The catalyst symbol, some exchanges need this The catalyst symbol, some exchanges need this
params:
""" """
pass pass
+73 -21
View File
@@ -164,6 +164,25 @@ class ExchangeTradingAlgorithmBase(TradingAlgorithm):
style) style)
return amount, style return amount, style
def _calculate_order_target_amount(self, asset, target):
"""
removes order amounts so we won't run into issues
when two orders are placed one after the other.
it then proceeds to removing positions amount at TradingAlgorithm
:param asset:
:param target:
:return: target
"""
if asset in self.blotter.open_orders:
for open_order in self.blotter.open_orders[asset]:
current_amount = open_order.amount
target -= current_amount
target = super(ExchangeTradingAlgorithmBase, self). \
_calculate_order_target_amount(asset, target)
return target
def round_order(self, amount, asset): def round_order(self, amount, asset):
""" """
We need fractions with cryptocurrencies We need fractions with cryptocurrencies
@@ -376,19 +395,30 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
if error: if error:
log.warning(error) log.warning(error)
self.pnl_stats = get_algo_df(self.algo_namespace, 'pnl_stats') # in order to save paper & live files separately
self.mode_name = 'paper' if kwargs['simulate_orders'] else 'live'
self.custom_signals_stats = \ self.pnl_stats = get_algo_df(
get_algo_df(self.algo_namespace, 'custom_signals_stats') self.algo_namespace,
'pnl_stats_{}'.format(self.mode_name),
)
self.exposure_stats = \ self.custom_signals_stats = get_algo_df(
get_algo_df(self.algo_namespace, 'exposure_stats') self.algo_namespace,
'custom_signals_stats_{}'.format(self.mode_name)
)
self.exposure_stats = get_algo_df(
self.algo_namespace,
'exposure_stats_{}'.format(self.mode_name)
)
self.is_running = True self.is_running = True
self.stats_minutes = 1 self.stats_minutes = 1
self._last_orders = [] self._last_orders = []
self._last_open_orders = []
self.trading_client = None self.trading_client = None
super(ExchangeTradingAlgorithmLive, self).__init__(*args, **kwargs) super(ExchangeTradingAlgorithmLive, self).__init__(*args, **kwargs)
@@ -515,7 +545,7 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
""" """
self.state = get_algo_object( self.state = get_algo_object(
algo_name=self.algo_namespace, algo_name=self.algo_namespace,
key='context.state', key='context.state_{}'.format(self.mode_name),
) )
if self.state is None: if self.state is None:
self.state = {} self.state = {}
@@ -538,7 +568,7 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
# Unpacking the perf_tracker and positions if available # Unpacking the perf_tracker and positions if available
cum_perf = get_algo_object( cum_perf = get_algo_object(
algo_name=self.algo_namespace, algo_name=self.algo_namespace,
key='cumulative_performance', key='cumulative_performance_{}'.format(self.mode_name),
) )
if cum_perf is not None: if cum_perf is not None:
tracker.cumulative_performance = cum_perf tracker.cumulative_performance = cum_perf
@@ -549,7 +579,7 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
todays_perf = get_algo_object( todays_perf = get_algo_object(
algo_name=self.algo_namespace, algo_name=self.algo_namespace,
key=today.strftime('%Y-%m-%d'), key=today.strftime('%Y-%m-%d'),
rel_path='daily_performance', rel_path='daily_performance_{}'.format(self.mode_name),
) )
if todays_perf is not None: if todays_perf is not None:
# Ensure single common position tracker # Ensure single common position tracker
@@ -687,7 +717,11 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
) )
self.pnl_stats = pd.concat([self.pnl_stats, df]) self.pnl_stats = pd.concat([self.pnl_stats, df])
save_algo_df(self.algo_namespace, 'pnl_stats', self.pnl_stats) save_algo_df(
self.algo_namespace,
'pnl_stats_{}'.format(self.mode_name),
self.pnl_stats,
)
def add_custom_signals_stats(self, period_stats): def add_custom_signals_stats(self, period_stats):
""" """
@@ -708,8 +742,11 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
) )
self.custom_signals_stats = pd.concat([self.custom_signals_stats, df]) self.custom_signals_stats = pd.concat([self.custom_signals_stats, df])
save_algo_df(self.algo_namespace, 'custom_signals_stats', save_algo_df(
self.custom_signals_stats) self.algo_namespace,
'custom_signals_stats_{}'.format(self.mode_name),
self.custom_signals_stats,
)
def add_exposure_stats(self, period_stats): def add_exposure_stats(self, period_stats):
""" """
@@ -736,7 +773,9 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
self.exposure_stats = pd.concat([self.exposure_stats, df]) self.exposure_stats = pd.concat([self.exposure_stats, df])
save_algo_df( save_algo_df(
self.algo_namespace, 'exposure_stats', self.exposure_stats self.algo_namespace,
'exposure_stats_{}'.format(self.mode_name),
self.exposure_stats
) )
def nullify_frame_stats(self, now): def nullify_frame_stats(self, now):
@@ -760,6 +799,7 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
obj=self.frame_stats, obj=self.frame_stats,
rel_path='frame_stats' rel_path='frame_stats'
) )
error = remove_old_files( error = remove_old_files(
algo_name=self.algo_namespace, algo_name=self.algo_namespace,
today=now, today=now,
@@ -792,12 +832,17 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
self.nullify_frame_stats(now=data.current_dt) self.nullify_frame_stats(now=data.current_dt)
self.performance_needs_update = False self.performance_needs_update = False
orders = list(self.perf_tracker.todays_performance.orders_by_id.keys()) last_orders_list = list(self.blotter.orders.keys())
if orders != self._last_orders: open_orders_list = list(self.blotter.open_orders.keys())
if last_orders_list != self._last_orders or \
open_orders_list != self._last_open_orders:
self.performance_needs_update = True self.performance_needs_update = True
# Saving current orders to detect changes in the next frame # Saving current order positions
self._last_orders = copy.deepcopy(orders) # to detect changes in the next frame
self._last_orders = copy.deepcopy(last_orders_list)
self._last_open_orders = copy.deepcopy(open_orders_list)
if self.performance_needs_update: if self.performance_needs_update:
self.perf_tracker.update_performance() self.perf_tracker.update_performance()
@@ -839,7 +884,7 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
log.debug('saving cumulative performance object') log.debug('saving cumulative performance object')
save_algo_object( save_algo_object(
algo_name=self.algo_namespace, algo_name=self.algo_namespace,
key='cumulative_performance', key='cumulative_performance_{}'.format(self.mode_name),
obj=self.perf_tracker.cumulative_performance, obj=self.perf_tracker.cumulative_performance,
) )
log.debug('saving todays performance object') log.debug('saving todays performance object')
@@ -847,12 +892,12 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
algo_name=self.algo_namespace, algo_name=self.algo_namespace,
key=today.strftime('%Y-%m-%d'), key=today.strftime('%Y-%m-%d'),
obj=self.perf_tracker.todays_performance, obj=self.perf_tracker.todays_performance,
rel_path='daily_performance' rel_path='daily_performance_{}'.format(self.mode_name)
) )
log.debug('saving context.state object') log.debug('saving context.state object')
save_algo_object( save_algo_object(
algo_name=self.algo_namespace, algo_name=self.algo_namespace,
key='context.state', key='context.state_{}'.format(self.mode_name),
obj=self.state) obj=self.state)
def _process_stats(self, data): def _process_stats(self, data):
@@ -908,6 +953,7 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
csv_bytes = stats_to_algo_folder( csv_bytes = stats_to_algo_folder(
stats=self.frame_stats, stats=self.frame_stats,
algo_namespace=self.algo_namespace, algo_namespace=self.algo_namespace,
folder_name='stats_{}'.format(self.mode_name),
recorded_cols=recorded_cols, recorded_cols=recorded_cols,
) )
except Exception as e: except Exception as e:
@@ -1012,13 +1058,19 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
args=(order_id,)) args=(order_id,))
@api_method @api_method
def cancel_order(self, order_param, exchange_name): def cancel_order(self, order_param, exchange_name,
symbol=None, params={}):
"""Cancel an open order. """Cancel an open order.
Parameters Parameters
---------- ----------
order_param : str or Order order_param : str or Order
The order_id or order object to cancel. The order_id or order object to cancel.
exchange_name: name of exchange from
which you want to cancel the order
symbol:
params:
""" """
exchange = self.exchanges[exchange_name] exchange = self.exchanges[exchange_name]
@@ -1032,4 +1084,4 @@ class ExchangeTradingAlgorithmLive(ExchangeTradingAlgorithmBase):
sleeptime=self.attempts['retry_sleeptime'], sleeptime=self.attempts['retry_sleeptime'],
retry_exceptions=(ExchangeRequestError,), retry_exceptions=(ExchangeRequestError,),
cleanup=lambda: log.warn('cancelling order again.'), cleanup=lambda: log.warn('cancelling order again.'),
args=(order_id,)) args=(order_id, symbol, params))
+1 -1
View File
@@ -68,7 +68,7 @@ class TradingPairFeeSchedule(CommissionModel):
multiplier = maker \ multiplier = maker \
if ((order.amount > 0 and order.limit < transaction.price) if ((order.amount > 0 and order.limit < transaction.price)
or (order.amount < 0 and order.limit > transaction.price)) \ or (order.amount < 0 and order.limit > transaction.price)) \
and order.limit_reached else taker and order.limit_reached else taker
fee = cost * multiplier fee = cost * multiplier
return fee return fee
+20 -13
View File
@@ -28,7 +28,8 @@ from catalyst.exchange.exchange_errors import EmptyValuesInBundleError, \
from catalyst.exchange.utils.bundle_utils import range_in_bundle, \ from catalyst.exchange.utils.bundle_utils import range_in_bundle, \
get_bcolz_chunk, get_df_from_arrays, get_assets get_bcolz_chunk, get_df_from_arrays, get_assets
from catalyst.exchange.utils.datetime_utils import get_start_dt, \ from catalyst.exchange.utils.datetime_utils import get_start_dt, \
get_period_label, get_month_start_end, get_year_start_end get_period_label, get_month_start_end, get_year_start_end, get_period, \
timestr_to_dt
from catalyst.exchange.utils.exchange_utils import get_exchange_folder from catalyst.exchange.utils.exchange_utils import get_exchange_folder
from catalyst.utils.cli import maybe_show_progress from catalyst.utils.cli import maybe_show_progress
from catalyst.utils.paths import ensure_directory from catalyst.utils.paths import ensure_directory
@@ -513,8 +514,8 @@ class ExchangeBundle:
continue continue
dates = pd.date_range( dates = pd.date_range(
start=get_period_label(adj_start, data_frequency), start=get_period(adj_start, data_frequency),
end=get_period_label(adj_end, data_frequency), end=get_period(adj_end, data_frequency),
freq='MS' if data_frequency == 'minute' else 'AS', freq='MS' if data_frequency == 'minute' else 'AS',
tz=UTC tz=UTC
) )
@@ -553,7 +554,9 @@ class ExchangeBundle:
# We sort the chunks by end date to ingest most recent data first # We sort the chunks by end date to ingest most recent data first
chunks[asset].sort( chunks[asset].sort(
key=lambda chunk: pd.to_datetime(chunk['period']) key=lambda chunk: timestr_to_dt(
chunk['period'], data_frequency
)
) )
return chunks return chunks
@@ -608,7 +611,8 @@ class ExchangeBundle:
exchange=self.exchange_name, exchange=self.exchange_name,
frequency=data_frequency, frequency=data_frequency,
symbol=asset.symbol symbol=asset.symbol
)) as it: )
) as it:
for chunk in it: for chunk in it:
problems += self.ingest_ctable( problems += self.ingest_ctable(
asset=chunk['asset'], asset=chunk['asset'],
@@ -623,16 +627,19 @@ class ExchangeBundle:
# We sort the chunks by end date to ingest most recent data first # We sort the chunks by end date to ingest most recent data first
all_chunks.sort( all_chunks.sort(
key=lambda chunk: pd.to_datetime(chunk['period']) key=lambda chunk: timestr_to_dt(
chunk['period'], data_frequency
)
) )
with maybe_show_progress( with maybe_show_progress(
all_chunks, all_chunks,
show_progress, show_progress,
label='Ingesting {frequency} price data on ' label='Ingesting {frequency} price data on '
'{exchange}'.format( '{exchange}'.format(
exchange=self.exchange_name, exchange=self.exchange_name,
frequency=data_frequency, frequency=data_frequency,
)) as it: )
) as it:
for chunk in it: for chunk in it:
problems += self.ingest_ctable( problems += self.ingest_ctable(
asset=chunk['asset'], asset=chunk['asset'],
+7
View File
@@ -324,6 +324,13 @@ class BalanceTooLowError(ZiplineError):
).strip() ).strip()
class NoCandlesReceivedFromExchange(ZiplineError):
msg = (
'Although requesting {bar_count} candles until {end_dt} of asset {asset}, '
'an empty list of candles was received for {exchange}.'
).strip()
class MarketsNotFoundError(ZiplineError): class MarketsNotFoundError(ZiplineError):
msg = ( msg = (
'Exchange {exchange} contains no valid market so it is unusable in ' 'Exchange {exchange} contains no valid market so it is unusable in '
+7 -5
View File
@@ -5,6 +5,7 @@ from datetime import datetime
import numpy as np import numpy as np
import pandas as pd import pandas as pd
from catalyst.constants import BUNDLE_URL
from catalyst.data.bundles.core import download_without_progress from catalyst.data.bundles.core import download_without_progress
from catalyst.exchange.utils.exchange_utils import get_exchange_bundles_folder from catalyst.exchange.utils.exchange_utils import get_exchange_bundles_folder
import os import os
@@ -48,10 +49,11 @@ def get_bcolz_chunk(exchange_name, symbol, data_frequency, period):
path = os.path.join(root, name) path = os.path.join(root, name)
if not os.path.isdir(path): if not os.path.isdir(path):
url = 'https://s3.amazonaws.com/enigmaco/catalyst-bundles/' \ url = BUNDLE_URL.format(
'exchange-{exchange}/{name}.tar.gz'.format(
exchange=exchange_name, exchange=exchange_name,
name=name) data_frequency=data_frequency,
name=name,
)
bytes = download_without_progress(url) bytes = download_without_progress(url)
with tarfile.open('r', fileobj=bytes) as tar: with tarfile.open('r', fileobj=bytes) as tar:
@@ -75,14 +77,14 @@ def get_df_from_arrays(arrays, periods):
""" """
ohlcv = dict() ohlcv = dict()
for index, field in enumerate( for index, field in enumerate(['open', 'high', 'low', 'close', 'volume']):
['open', 'high', 'low', 'close', 'volume']):
ohlcv[field] = arrays[index].flatten() ohlcv[field] = arrays[index].flatten()
df = pd.DataFrame( df = pd.DataFrame(
data=ohlcv, data=ohlcv,
index=periods index=periods
) )
df.index.name = 'last_traded'
return df return df
+1 -1
View File
@@ -71,7 +71,7 @@ def scan_exchange_configs(features=None, history=None, is_authenticated=False,
def get_exchange_config(exchange_name, path=None, environ=None, def get_exchange_config(exchange_name, path=None, environ=None,
expiry='2H'): expiry='1H'):
""" """
The de-serialized content of the exchange's config.json. The de-serialized content of the exchange's config.json.
Parameters Parameters
+26
View File
@@ -164,6 +164,12 @@ def get_start_dt(end_dt, bar_count, data_frequency, include_first=True):
return start_dt return start_dt
def timestr_to_dt(timestr, data_frequency):
dt_format = '%Y' if data_frequency == 'daily' else '%Y%m'
dt = pd.to_datetime(timestr, format=dt_format, utc=True)
return dt
def get_period_label(dt, data_frequency): def get_period_label(dt, data_frequency):
""" """
The period label for the specified date and frequency. The period label for the specified date and frequency.
@@ -177,6 +183,26 @@ def get_period_label(dt, data_frequency):
------- -------
str str
"""
if data_frequency == 'minute':
return '{}{:02d}'.format(dt.year, dt.month)
else:
return '{}'.format(dt.year)
def get_period(dt, data_frequency):
"""
The period label for the specified date and frequency.
Parameters
----------
dt: datetime
data_frequency: str
Returns
-------
str
""" """
if data_frequency == 'minute': if data_frequency == 'minute':
return '{}-{:02d}'.format(dt.year, dt.month) return '{}-{:02d}'.format(dt.year, dt.month)
+79 -32
View File
@@ -1,24 +1,21 @@
import hashlib import hashlib
import json
import os import os
import pickle
import shutil import shutil
from datetime import date, datetime
import json
import pandas as pd import pandas as pd
import pickle
from catalyst.assets._assets import TradingPair from catalyst.assets._assets import TradingPair
from datetime import date, datetime
from six import string_types from six import string_types
from six.moves.urllib import request from six.moves.urllib import request
from catalyst.constants import DATE_FORMAT, SYMBOLS_URL from catalyst.constants import EXCHANGE_CONFIG_URL
from catalyst.exchange.exchange_errors import ExchangeSymbolsNotFound, \
InvalidHistoryFrequencyError, InvalidHistoryFrequencyAlias
from catalyst.exchange.utils.serialization_utils import ExchangeJSONEncoder, \ from catalyst.exchange.utils.serialization_utils import ExchangeJSONEncoder, \
ExchangeJSONDecoder, ConfigJSONEncoder ExchangeJSONDecoder, ConfigJSONEncoder
from catalyst.utils.deprecate import deprecated
from catalyst.utils.paths import data_root, ensure_directory, \ from catalyst.utils.paths import data_root, ensure_directory, \
last_modified_time last_modified_time
from six import string_types
from six.moves.urllib import request
def get_sid(symbol): def get_sid(symbol):
@@ -109,6 +106,7 @@ def download_exchange_config(exchange_name, filename, environ=None):
request.urlretrieve(url=url, filename=filename) request.urlretrieve(url=url, filename=filename)
@deprecated
def get_exchange_config(exchange_name, filename=None, environ=None): def get_exchange_config(exchange_name, filename=None, environ=None):
""" """
The de-serialized content of the exchange's config.json. The de-serialized content of the exchange's config.json.
@@ -144,6 +142,7 @@ def get_exchange_config(exchange_name, filename=None, environ=None):
except ValueError: except ValueError:
return dict() return dict()
def save_exchange_config(exchange_name, config, filename=None, environ=None): def save_exchange_config(exchange_name, config, filename=None, environ=None):
""" """
Save assets into an exchange_config file. Save assets into an exchange_config file.
@@ -413,7 +412,7 @@ def clear_frame_stats_directory(algo_name):
return error return error
def remove_old_files(algo_name, today, rel_path): def remove_old_files(algo_name, today, rel_path, environ=None):
""" """
remove old files from a directory remove old files from a directory
to avoid overloading the disk to avoid overloading the disk
@@ -423,27 +422,31 @@ def remove_old_files(algo_name, today, rel_path):
algo_name: str algo_name: str
today: Timestamp today: Timestamp
rel_path: str rel_path: str
environ:
Returns Returns
------- -------
error: str error: str
""" """
error = None error = None
algo_folder = get_algo_folder(algo_name) algo_folder = get_algo_folder(algo_name, environ)
folder = os.path.join(algo_folder, rel_path) folder = os.path.join(algo_folder, rel_path)
ensure_directory(folder)
# run on all files in the folder # run on all files in the folder
for f in os.listdir(folder): for f in os.listdir(folder):
creation_unix = os.path.getctime(f) try:
creation_time = pd.to_datetime(creation_unix, unit='s', ) file_path = os.path.join(folder, f)
creation_unix = os.path.getctime(file_path)
creation_time = pd.to_datetime(creation_unix, unit='s', utc=True)
# if the file is older than 30 days erase it # if the file is older than 30 days erase it
if today - pd.DateOffset(30) > creation_time: if today - pd.DateOffset(30) > creation_time:
try: os.unlink(file_path)
os.unlink(f) except OSError:
except OSError: error = 'unable to erase files in {}'.format(folder)
error = 'unable to erase files in {}'.format(folder)
return error return error
@@ -652,25 +655,69 @@ def save_asset_data(folder, df, decimals=8):
) )
def get_candles_df(candles, field, freq, bar_count, end_dt, def forward_fill_df_if_needed(df, periods):
previous_value=None): df = df.reindex(periods)
# volume should always be 0 (if there were no trades in this interval)
df['volume'] = df['volume'].fillna(0.0)
# ie pull the last close into this close
df['close'] = df.fillna(method='pad')
# now copy the close that was pulled down from the last timestep
# into this row, across into o/h/l
df['open'] = df['open'].fillna(df['close'])
df['low'] = df['low'].fillna(df['close'])
df['high'] = df['high'].fillna(df['close'])
return df
def transform_candles_to_df(candles):
return pd.DataFrame(candles).set_index('last_traded')
def get_candles_df(candles, field, freq, bar_count, end_dt):
all_series = dict() all_series = dict()
for asset in candles: for asset in candles:
periods = pd.date_range(end=end_dt, periods=bar_count, freq=freq) asset_df = transform_candles_to_df(candles[asset])
rounded_end_dt = end_dt.floor(freq)
periods = pd.date_range(end=rounded_end_dt,
periods=bar_count,
freq=freq)
asset_df = forward_fill_df_if_needed(asset_df, periods)
dates = [candle['last_traded'] for candle in candles[asset]] all_series[asset] = pd.Series(asset_df[field])
values = [candle[field] for candle in candles[asset]]
series = pd.Series(values, index=dates)
series = series.reindex(
periods,
method='ffill',
fill_value=previous_value,
)
series.sort_index(inplace=True)
all_series[asset] = series
df = pd.DataFrame(all_series) df = pd.DataFrame(all_series)
df.dropna(inplace=True) df.dropna(inplace=True)
return df return df
def get_trades_df(trades):
df = pd.DataFrame(trades)
df.index = pd.to_datetime(df.pop('datetime'))
df.index = df.index.tz_localize('UTC')
return df
def candles_from_trades(trades_df, freq):
"""
Calculate OHLCV from candles.
Parameters
----------
trades_df
freq
Returns
-------
"""
df = trades_df['price'].resample(freq).ohlc() # type: pd.DataFrame
df['volume'] = trades_df['amount'].resample(freq).sum()
df.dropna(axis=0, how='all', inplace=True)
df.sort_index(inplace=True, ascending=False)
return df
+3
View File
@@ -17,6 +17,9 @@ def get_exchange(exchange_name, base_currency=None, must_authenticate=False,
skip_init=False, auth_alias=None, config=None): skip_init=False, auth_alias=None, config=None):
key = (exchange_name, base_currency) key = (exchange_name, base_currency)
if key in exchange_cache: if key in exchange_cache:
if not skip_init:
exchange_cache[key].init()
return exchange_cache[key] return exchange_cache[key]
exchange_auth = get_exchange_auth(exchange_name, alias=auth_alias) exchange_auth = get_exchange_auth(exchange_name, alias=auth_alias)
@@ -37,7 +37,13 @@ class ExchangeJSONEncoder(json.JSONEncoder):
return obj.strftime(DATE_TIME_FORMAT) return obj.strftime(DATE_TIME_FORMAT)
elif isinstance(obj, TradingPair): elif isinstance(obj, TradingPair):
return obj.to_dict() asset = obj.to_dict()
asset['maker'] = round(asset['maker'], asset['decimals'])
asset['taker'] = round(asset['taker'], asset['decimals'])
asset['lot'] = round(asset['lot'], 4)
asset['min_trade_size'] = round(asset['min_trade_size'], 4)
asset['max_trade_size'] = round(asset['max_trade_size'], 4)
return asset
# Let the base class default method raise the TypeError # Let the base class default method raise the TypeError
return JSONEncoder.default(self, obj) return JSONEncoder.default(self, obj)
+4 -2
View File
@@ -396,7 +396,8 @@ def email_error(algo_name, dt, e, environ=None):
)}) )})
def stats_to_algo_folder(stats, algo_namespace, recorded_cols=None): def stats_to_algo_folder(stats, algo_namespace,
folder_name, recorded_cols=None):
""" """
Saves the performance stats to the algo local folder. Saves the performance stats to the algo local folder.
@@ -404,6 +405,7 @@ def stats_to_algo_folder(stats, algo_namespace, recorded_cols=None):
---------- ----------
stats: list[Object] stats: list[Object]
algo_namespace: str algo_namespace: str
folder_name: str
recorded_cols: list[str] recorded_cols: list[str]
Returns Returns
@@ -416,7 +418,7 @@ def stats_to_algo_folder(stats, algo_namespace, recorded_cols=None):
timestr = time.strftime('%Y%m%d') timestr = time.strftime('%Y%m%d')
folder = get_algo_folder(algo_namespace) folder = get_algo_folder(algo_namespace)
stats_folder = os.path.join(folder, 'stats') stats_folder = os.path.join(folder, folder_name)
ensure_directory(stats_folder) ensure_directory(stats_folder)
filename = os.path.join(stats_folder, '{}.csv'.format(timestr)) filename = os.path.join(stats_folder, '{}.csv'.format(timestr))
+132 -38
View File
@@ -7,6 +7,7 @@ import re
import shutil import shutil
import sys import sys
import time import time
import webbrowser
import bcolz import bcolz
import logbook import logbook
@@ -23,7 +24,7 @@ from catalyst.exchange.utils.stats_utils import set_print_settings
from catalyst.marketplace.marketplace_errors import ( from catalyst.marketplace.marketplace_errors import (
MarketplacePubAddressEmpty, MarketplaceDatasetNotFound, MarketplacePubAddressEmpty, MarketplaceDatasetNotFound,
MarketplaceNoAddressMatch, MarketplaceHTTPRequest, MarketplaceNoAddressMatch, MarketplaceHTTPRequest,
MarketplaceNoCSVFiles) MarketplaceNoCSVFiles, MarketplaceRequiresPython3)
from catalyst.marketplace.utils.auth_utils import get_key_secret, \ from catalyst.marketplace.utils.auth_utils import get_key_secret, \
get_signed_headers get_signed_headers
from catalyst.marketplace.utils.bundle_utils import merge_bundles from catalyst.marketplace.utils.bundle_utils import merge_bundles
@@ -32,6 +33,7 @@ from catalyst.marketplace.utils.eth_utils import bin_hex, from_grains, \
from catalyst.marketplace.utils.path_utils import get_bundle_folder, \ from catalyst.marketplace.utils.path_utils import get_bundle_folder, \
get_data_source_folder, get_marketplace_folder, \ get_data_source_folder, get_marketplace_folder, \
get_user_pubaddr, get_temp_bundles_folder, extract_bundle get_user_pubaddr, get_temp_bundles_folder, extract_bundle
from catalyst.utils.paths import ensure_directory
if sys.version_info.major < 3: if sys.version_info.major < 3:
import urllib import urllib
@@ -44,7 +46,10 @@ log = logbook.Logger('Marketplace', level=LOG_LEVEL)
class Marketplace: class Marketplace:
def __init__(self): def __init__(self):
global Web3 global Web3
from web3 import Web3, HTTPProvider try:
from web3 import Web3, HTTPProvider
except ImportError:
raise MarketplaceRequiresPython3()
self.addresses = get_user_pubaddr() self.addresses = get_user_pubaddr()
@@ -60,7 +65,8 @@ class Marketplace:
contract_url = urllib.urlopen(MARKETPLACE_CONTRACT) contract_url = urllib.urlopen(MARKETPLACE_CONTRACT)
self.mkt_contract_address = Web3.toChecksumAddress( self.mkt_contract_address = Web3.toChecksumAddress(
contract_url.readline().strip()) contract_url.readline().decode(
contract_url.info().get_content_charset()).strip())
abi_url = urllib.urlopen(MARKETPLACE_CONTRACT_ABI) abi_url = urllib.urlopen(MARKETPLACE_CONTRACT_ABI)
abi = json.load(abi_url) abi = json.load(abi_url)
@@ -73,7 +79,8 @@ class Marketplace:
contract_url = urllib.urlopen(ENIGMA_CONTRACT) contract_url = urllib.urlopen(ENIGMA_CONTRACT)
self.eng_contract_address = Web3.toChecksumAddress( self.eng_contract_address = Web3.toChecksumAddress(
contract_url.readline().strip()) contract_url.readline().decode(
contract_url.info().get_content_charset()).strip())
abi_url = urllib.urlopen(ENIGMA_CONTRACT_ABI) abi_url = urllib.urlopen(ENIGMA_CONTRACT_ABI)
abi = json.load(abi_url) abi = json.load(abi_url)
@@ -136,10 +143,10 @@ class Marketplace:
return address, address_i return address, address_i
def sign_transaction(self, from_address, tx): def sign_transaction(self, tx):
print('\nVisit https://www.myetherwallet.com/#offline-transaction and ' url = 'https://www.myetherwallet.com/#offline-transaction'
'enter the following parameters:\n\n' print('\nVisit {url} and enter the following parameters:\n\n'
'From Address:\t\t{_from}\n' 'From Address:\t\t{_from}\n'
'\n\tClick the "Generate Information" button\n\n' '\n\tClick the "Generate Information" button\n\n'
'To Address:\t\t{to}\n' 'To Address:\t\t{to}\n'
@@ -148,13 +155,16 @@ class Marketplace:
'Gas Price:\t\t[Accept the default value]\n' 'Gas Price:\t\t[Accept the default value]\n'
'Nonce:\t\t\t{nonce}\n' 'Nonce:\t\t\t{nonce}\n'
'Data:\t\t\t{data}\n'.format( 'Data:\t\t\t{data}\n'.format(
_from=from_address, url=url,
to=tx['to'], _from=tx['from'],
value=tx['value'], to=tx['to'],
gas=tx['gas'], value=tx['value'],
nonce=tx['nonce'], gas=tx['gas'],
data=tx['data'], ) nonce=tx['nonce'],
) data=tx['data'], )
)
webbrowser.open_new(url)
signed_tx = input('Copy and Paste the "Signed Transaction" ' signed_tx = input('Copy and Paste the "Signed Transaction" '
'field here:\n') 'field here:\n')
@@ -175,8 +185,7 @@ class Marketplace:
print('\nYou can check the outcome of your transaction here:\n' print('\nYou can check the outcome of your transaction here:\n'
'{}\n\n'.format(etherscan)) '{}\n\n'.format(etherscan))
def list(self): def _list(self):
data_sources = self.mkt_contract.functions.getAllProviders().call() data_sources = self.mkt_contract.functions.getAllProviders().call()
data = [] data = []
@@ -188,15 +197,44 @@ class Marketplace:
dataset=self.to_text(data_source) dataset=self.to_text(data_source)
) )
) )
return pd.DataFrame(data)
def list(self):
df = self._list()
df = pd.DataFrame(data)
set_print_settings() set_print_settings()
if df.empty: if df.empty:
print('There are no datasets available yet.') print('There are no datasets available yet.')
else: else:
print(df) print(df)
def subscribe(self, dataset): def subscribe(self, dataset=None):
if dataset is None:
df_sets = self._list()
if df_sets.empty:
print('There are no datasets available yet.')
return
set_print_settings()
while True:
print(df_sets)
dataset_num = input('Choose the dataset you want to '
'subscribe to [0..{}]: '.format(
df_sets.size - 1))
try:
dataset_num = int(dataset_num)
except ValueError:
print('Enter a number between 0 and {}'.format(
df_sets.size - 1))
else:
if dataset_num not in range(0, df_sets.size):
print('Enter a number between 0 and {}'.format(
df_sets.size - 1))
else:
dataset = df_sets.iloc[dataset_num]['dataset']
break
dataset = dataset.lower() dataset = dataset.lower()
@@ -259,14 +297,14 @@ class Marketplace:
'buy: {} ENG. Get enough ENG to cover the costs of the ' 'buy: {} ENG. Get enough ENG to cover the costs of the '
'monthly\nsubscription for what you are trying to buy, ' 'monthly\nsubscription for what you are trying to buy, '
'and try again.'.format( 'and try again.'.format(
address, from_grains(balance), price)) address, from_grains(balance), price))
return return
while True: while True:
agree_pay = input('Please confirm that you agree to pay {} ENG ' agree_pay = input('Please confirm that you agree to pay {} ENG '
'for a monthly subscription to the dataset "{}" ' 'for a monthly subscription to the dataset "{}" '
'starting today. [default: Y] '.format( 'starting today. [default: Y] '.format(
price, dataset)) or 'y' price, dataset)) or 'y'
if agree_pay.lower() not in ('y', 'n'): if agree_pay.lower() not in ('y', 'n'):
print("Please answer Y or N.") print("Please answer Y or N.")
else: else:
@@ -287,13 +325,14 @@ class Marketplace:
self.mkt_contract_address, self.mkt_contract_address,
grains, grains,
).buildTransaction( ).buildTransaction(
{'nonce': self.web3.eth.getTransactionCount(address)} {'from': address,
'nonce': self.web3.eth.getTransactionCount(address)}
) )
if 'ropsten' in ETH_REMOTE_NODE: if 'ropsten' in ETH_REMOTE_NODE:
tx['gas'] = min(int(tx['gas'] * 1.5), 4700000) tx['gas'] = min(int(tx['gas'] * 1.5), 4700000)
signed_tx = self.sign_transaction(address, tx) signed_tx = self.sign_transaction(tx)
try: try:
tx_hash = '0x{}'.format( tx_hash = '0x{}'.format(
bin_hex(self.web3.eth.sendRawTransaction(signed_tx)) bin_hex(self.web3.eth.sendRawTransaction(signed_tx))
@@ -328,13 +367,14 @@ class Marketplace:
tx = self.mkt_contract.functions.subscribe( tx = self.mkt_contract.functions.subscribe(
Web3.toHex(dataset), Web3.toHex(dataset),
).buildTransaction( ).buildTransaction({
{'nonce': self.web3.eth.getTransactionCount(address)}) 'from': address,
'nonce': self.web3.eth.getTransactionCount(address)})
if 'ropsten' in ETH_REMOTE_NODE: if 'ropsten' in ETH_REMOTE_NODE:
tx['gas'] = min(int(tx['gas'] * 1.5), 4700000) tx['gas'] = min(int(tx['gas'] * 1.5), 4700000)
signed_tx = self.sign_transaction(address, tx) signed_tx = self.sign_transaction(tx)
try: try:
tx_hash = '0x{}'.format(bin_hex( tx_hash = '0x{}'.format(bin_hex(
@@ -369,7 +409,7 @@ class Marketplace:
'You can now ingest this dataset anytime during the ' 'You can now ingest this dataset anytime during the '
'next month by running the following command:\n' 'next month by running the following command:\n'
'catalyst marketplace ingest --dataset={}'.format( 'catalyst marketplace ingest --dataset={}'.format(
dataset, address, dataset)) dataset, address, dataset))
def process_temp_bundle(self, ds_name, path): def process_temp_bundle(self, ds_name, path):
""" """
@@ -387,6 +427,7 @@ class Marketplace:
""" """
tmp_bundle = extract_bundle(path) tmp_bundle = extract_bundle(path)
bundle_folder = get_data_source_folder(ds_name) bundle_folder = get_data_source_folder(ds_name)
ensure_directory(bundle_folder)
if os.listdir(bundle_folder): if os.listdir(bundle_folder):
zsource = bcolz.ctable(rootdir=tmp_bundle, mode='r') zsource = bcolz.ctable(rootdir=tmp_bundle, mode='r')
ztarget = bcolz.ctable(rootdir=bundle_folder, mode='r') ztarget = bcolz.ctable(rootdir=bundle_folder, mode='r')
@@ -397,7 +438,33 @@ class Marketplace:
pass pass
def ingest(self, ds_name, start=None, end=None, force_download=False): def ingest(self, ds_name=None, start=None, end=None, force_download=False):
if ds_name is None:
df_sets = self._list()
if df_sets.empty:
print('There are no datasets available yet.')
return
set_print_settings()
while True:
print(df_sets)
dataset_num = input('Choose the dataset you want to '
'ingest [0..{}]: '.format(
df_sets.size - 1))
try:
dataset_num = int(dataset_num)
except ValueError:
print('Enter a number between 0 and {}'.format(
df_sets.size - 1))
else:
if dataset_num not in range(0, df_sets.size):
print('Enter a number between 0 and {}'.format(
df_sets.size - 1))
else:
ds_name = df_sets.iloc[dataset_num]['dataset']
break
# ds_name = ds_name.lower() # ds_name = ds_name.lower()
@@ -426,10 +493,10 @@ class Marketplace:
print('Your subscription to dataset "{}" expired on {} UTC.' print('Your subscription to dataset "{}" expired on {} UTC.'
'Please renew your subscription by running:\n' 'Please renew your subscription by running:\n'
'catalyst marketplace subscribe --dataset={}'.format( 'catalyst marketplace subscribe --dataset={}'.format(
ds_name, ds_name,
pd.to_datetime(check_sub[4], unit='s', utc=True), pd.to_datetime(check_sub[4], unit='s', utc=True),
ds_name) ds_name)
) )
if 'key' in self.addresses[address_i]: if 'key' in self.addresses[address_i]:
key = self.addresses[address_i]['key'] key = self.addresses[address_i]['key']
@@ -493,14 +560,40 @@ class Marketplace:
return df return df
def clean(self, data_source_name, data_frequency=None): def clean(self, ds_name=None, data_frequency=None):
data_source_name = data_source_name.lower()
if ds_name is None:
mktplace_root = get_marketplace_folder()
folders = [os.path.basename(f.rstrip('/'))
for f in glob.glob('{}/*/'.format(mktplace_root))
if 'temp_bundles' not in f]
while True:
for idx, f in enumerate(folders):
print('{}\t{}'.format(idx, f))
dataset_num = input('Choose the dataset you want to '
'clean [0..{}]: '.format(
len(folders) - 1))
try:
dataset_num = int(dataset_num)
except ValueError:
print('Enter a number between 0 and {}'.format(
len(folders) - 1))
else:
if dataset_num not in range(0, len(folders)):
print('Enter a number between 0 and {}'.format(
len(folders) - 1))
else:
ds_name = folders[dataset_num]
break
ds_name = ds_name.lower()
if data_frequency is None: if data_frequency is None:
folder = get_data_source_folder(data_source_name) folder = get_data_source_folder(ds_name)
else: else:
folder = get_bundle_folder(data_source_name, data_frequency) folder = get_bundle_folder(ds_name, data_frequency)
shutil.rmtree(folder) shutil.rmtree(folder)
pass pass
@@ -604,13 +697,14 @@ class Marketplace:
grains, grains,
address, address,
).buildTransaction( ).buildTransaction(
{'nonce': self.web3.eth.getTransactionCount(address)} {'from': address,
'nonce': self.web3.eth.getTransactionCount(address)}
) )
if 'ropsten' in ETH_REMOTE_NODE: if 'ropsten' in ETH_REMOTE_NODE:
tx['gas'] = min(int(tx['gas'] * 1.5), 4700000) tx['gas'] = min(int(tx['gas'] * 1.5), 4700000)
signed_tx = self.sign_transaction(address, tx) signed_tx = self.sign_transaction(tx)
try: try:
tx_hash = '0x{}'.format( tx_hash = '0x{}'.format(
@@ -621,7 +715,7 @@ class Marketplace:
) )
except Exception as e: except Exception as e:
print('Unable to subscribe to data source: {}'.format(e)) print('Unable to register the requested dataset: {}'.format(e))
return return
self.check_transaction(tx_hash) self.check_transaction(tx_hash)
+10 -1
View File
@@ -9,7 +9,8 @@ def silent_except_hook(exctype, excvalue, exctraceback):
MarketplaceNoAddressMatch, MarketplaceHTTPRequest, MarketplaceNoAddressMatch, MarketplaceHTTPRequest,
MarketplaceNoCSVFiles, MarketplaceContractDataNoMatch, MarketplaceNoCSVFiles, MarketplaceContractDataNoMatch,
MarketplaceSubscriptionExpired, MarketplaceJSONError, MarketplaceSubscriptionExpired, MarketplaceJSONError,
MarketplaceWalletNotSupported, MarketplaceEmptySignature]: MarketplaceWalletNotSupported, MarketplaceEmptySignature,
MarketplaceRequiresPython3]:
fn = traceback.extract_tb(exctraceback)[-1][0] fn = traceback.extract_tb(exctraceback)[-1][0]
ln = traceback.extract_tb(exctraceback)[-1][1] ln = traceback.extract_tb(exctraceback)[-1][1]
print("Error traceback: {1} (line {2})\n" print("Error traceback: {1} (line {2})\n"
@@ -86,3 +87,11 @@ class MarketplaceJSONError(ZiplineError):
'The configuration file {file} is malformed. Please correct ' 'The configuration file {file} is malformed. Please correct '
'the following error:\n{error}' 'the following error:\n{error}'
) )
class MarketplaceRequiresPython3(ZiplineError):
msg = (
'\nCatalyst requires Python3 to access the Enigma Data Marketplace.\n'
'If you want to use the Data Marketplace, you need to reinstall '
'Catalyst\nwith Python3. See the documentation website for additional '
'information.')
+59 -1
View File
@@ -1,8 +1,12 @@
import os import os
import random
import re
import shutil import shutil
import bcolz import bcolz
import numpy as np
import pandas as pd import pandas as pd
from six import string_types
def merge_bundles(zsource, ztarget): def merge_bundles(zsource, ztarget):
@@ -27,10 +31,64 @@ def merge_bundles(zsource, ztarget):
df.drop_duplicates(inplace=True) df.drop_duplicates(inplace=True)
df.set_index(['date', 'symbol'], drop=False, inplace=True) df.set_index(['date', 'symbol'], drop=False, inplace=True)
sanitize_df(df)
dirname = os.path.basename(ztarget.rootdir) dirname = os.path.basename(ztarget.rootdir)
bak_dir = ztarget.rootdir.replace(dirname, '.{}'.format(dirname)) bak_dir = ztarget.rootdir.replace(dirname, '.{}'.format(dirname))
os.rename(ztarget.rootdir, bak_dir) shutil.move(ztarget.rootdir, bak_dir)
z = bcolz.ctable.fromdataframe(df=df, rootdir=ztarget.rootdir) z = bcolz.ctable.fromdataframe(df=df, rootdir=ztarget.rootdir)
shutil.rmtree(bak_dir) shutil.rmtree(bak_dir)
return z return z
def sanitize_df(df):
# Using a sampling method to identify dates for efficiency with
# large datasets
if len(df) > 100:
indexes = random.sample(range(0, len(df) - 1), 100)
elif len(df) > 1:
indexes = range(0, len(df) - 1)
else:
indexes = [0, ]
for column in df.columns:
is_date = False
for index in indexes:
value = df[column].iloc[index]
if not isinstance(value, string_types):
continue
# TODO: assuming that the date is at least daily
exp = re.compile(r'^\d{4}-\d{2}-\d{2}.*$')
matches = exp.findall(value)
if matches:
is_date = True
break
if is_date:
df[column] = pd.to_datetime(df[column])
else:
try:
ser = safely_reduce_dtype(df[column])
df[column] = ser
except Exception:
pass
return df
def safely_reduce_dtype(ser): # pandas.Series or numpy.array
orig_dtype = "".join(
[x for x in ser.dtype.name if x.isalpha()]) # float/int
mx = 1
for val in ser.values:
new_itemsize = np.min_scalar_type(val).itemsize
if mx < new_itemsize:
mx = new_itemsize
if orig_dtype == 'int':
mx = max(mx, 4)
new_dtype = orig_dtype + str(mx * 8)
return ser.astype(new_dtype)
+34
View File
@@ -0,0 +1,34 @@
from catalyst.api import symbol
from catalyst.utils.run_algo import run_algorithm
coins = ['dash', 'btc', 'dash', 'etc', 'eth', 'ltc', 'nxt', 'rep', 'str',
'xmr', 'xrp', 'zec']
symbols = None
def initialize(context):
pass
def _handle_data(context, data):
global symbols
if symbols is None: symbols = [symbol(c + '_usdt') for c in coins]
print('getting history for: %s' % [s.symbol for s in symbols])
history = data.history(symbols,
['close', 'volume'],
bar_count=1, # EXCEPTION, Change to 2
frequency='5T')
# print 'history: %s' % history.shape
run_algorithm(initialize=initialize,
handle_data=_handle_data,
analyze=lambda _, results: True,
exchange_name='poloniex',
base_currency='usdt',
algo_namespace='issue-236',
live=True,
data_frequency='minute',
capital_base=3000,
simulate_orders=True)
+12 -1
View File
@@ -10,6 +10,7 @@ import click
import pandas as pd import pandas as pd
from six import string_types from six import string_types
import catalyst
from catalyst.data.bundles import load from catalyst.data.bundles import load
from catalyst.data.data_portal import DataPortal from catalyst.data.data_portal import DataPortal
from catalyst.exchange.exchange_pricing_loader import ExchangePricingLoader, \ from catalyst.exchange.exchange_pricing_loader import ExchangePricingLoader, \
@@ -23,7 +24,7 @@ try:
from pygments.formatters import TerminalFormatter from pygments.formatters import TerminalFormatter
PYGMENTS = True PYGMENTS = True
except: except ImportError:
PYGMENTS = False PYGMENTS = False
from toolz import valfilter, concatv from toolz import valfilter, concatv
from functools import partial from functools import partial
@@ -151,6 +152,7 @@ def _run(handle_data,
'We encourage you to report any issue on GitHub: ' 'We encourage you to report any issue on GitHub: '
'https://github.com/enigmampc/catalyst/issues' 'https://github.com/enigmampc/catalyst/issues'
) )
log.info('Catalyst version {}'.format(catalyst.__version__))
sleep(3) sleep(3)
if live: if live:
@@ -261,6 +263,15 @@ def _run(handle_data,
# We still need to support bundles for other misc data, but we # We still need to support bundles for other misc data, but we
# can handle this later. # can handle this later.
if start != pd.tslib.normalize_date(start) or \
end != pd.tslib.normalize_date(end):
# todo: add to Sim_Params the option to start & end at specific times
log.warn(
"Catalyst currently starts and ends on the start and "
"end of the dates specified, respectively. We hope to "
"Modify this and support specific times in a future release."
)
data = DataPortalExchangeBacktest( data = DataPortalExchangeBacktest(
exchange_names=[exchange_name for exchange_name in exchanges], exchange_names=[exchange_name for exchange_name in exchanges],
asset_finder=None, asset_finder=None,
+2 -156
View File
@@ -580,162 +580,8 @@ which you can skim through for now. A copy of this algorithm is available in
the ``examples`` directory: the ``examples`` directory:
`dual_moving_average.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/dual_moving_average.py>`_. `dual_moving_average.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/dual_moving_average.py>`_.
.. code-block:: python .. literalinclude:: ../../catalyst/examples/dual_moving_average.py
:language: python
import numpy as np
import pandas as pd
from logbook import Logger
import matplotlib.pyplot as plt
from catalyst import run_algorithm
from catalyst.api import (order, record, symbol, order_target_percent,
get_open_orders)
from catalyst.exchange.utils.stats_utils import extract_transactions
NAMESPACE = 'dual_moving_average'
log = Logger(NAMESPACE)
def initialize(context):
context.i = 0
context.asset = symbol('ltc_usd')
context.base_price = None
def handle_data(context, data):
# define the windows for the moving averages
short_window = 50
long_window = 200
# Skip as many bars as long_window to properly compute the average
context.i += 1
if context.i < long_window:
return
# Compute moving averages calling data.history() for each
# moving average with the appropriate parameters. We choose to use
# minute bars for this simulation -> freq="1m"
# Returns a pandas dataframe.
short_mavg = data.history(context.asset, 'price',
bar_count=short_window, frequency="1m").mean()
long_mavg = data.history(context.asset, 'price',
bar_count=long_window, frequency="1m").mean()
# Let's keep the price of our asset in a more handy variable
price = data.current(context.asset, 'price')
# If base_price is not set, we use the current value. This is the
# price at the first bar which we reference to calculate price_change.
if context.base_price is None:
context.base_price = price
price_change = (price - context.base_price) / context.base_price
# Save values for later inspection
record(price=price,
cash=context.portfolio.cash,
price_change=price_change,
short_mavg=short_mavg,
long_mavg=long_mavg)
# Since we are using limit orders, some orders may not execute immediately
# we wait until all orders are executed before considering more trades.
orders = get_open_orders(context.asset)
if len(orders) > 0:
return
# Exit if we cannot trade
if not data.can_trade(context.asset):
return
# We check what's our position on our portfolio and trade accordingly
pos_amount = context.portfolio.positions[context.asset].amount
# Trading logic
if short_mavg > long_mavg and pos_amount == 0:
# we buy 100% of our portfolio for this asset
order_target_percent(context.asset, 1)
elif short_mavg < long_mavg and pos_amount > 0:
# we sell all our positions for this asset
order_target_percent(context.asset, 0)
def analyze(context, perf):
# Get the base_currency that was passed as a parameter to the simulation
exchange = list(context.exchanges.values())[0]
base_currency = exchange.base_currency.upper()
# First chart: Plot portfolio value using base_currency
ax1 = plt.subplot(411)
perf.loc[:, ['portfolio_value']].plot(ax=ax1)
ax1.legend_.remove()
ax1.set_ylabel('Portfolio Value\n({})'.format(base_currency))
start, end = ax1.get_ylim()
ax1.yaxis.set_ticks(np.arange(start, end, (end-start)/5))
# Second chart: Plot asset price, moving averages and buys/sells
ax2 = plt.subplot(412, sharex=ax1)
perf.loc[:, ['price','short_mavg','long_mavg']].plot(ax=ax2, label='Price')
ax2.legend_.remove()
ax2.set_ylabel('{asset}\n({base})'.format(
asset = context.asset.symbol,
base = base_currency
))
start, end = ax2.get_ylim()
ax2.yaxis.set_ticks(np.arange(start, end, (end-start)/5))
transaction_df = extract_transactions(perf)
if not transaction_df.empty:
buy_df = transaction_df[transaction_df['amount'] > 0]
sell_df = transaction_df[transaction_df['amount'] < 0]
ax2.scatter(
buy_df.index.to_pydatetime(),
perf.loc[buy_df.index, 'price'],
marker='^',
s=100,
c='green',
label=''
)
ax2.scatter(
sell_df.index.to_pydatetime(),
perf.loc[sell_df.index, 'price'],
marker='v',
s=100,
c='red',
label=''
)
# Third chart: Compare percentage change between our portfolio
# and the price of the asset
ax3 = plt.subplot(413, sharex=ax1)
perf.loc[:, ['algorithm_period_return', 'price_change']].plot(ax=ax3)
ax3.legend_.remove()
ax3.set_ylabel('Percent Change')
start, end = ax3.get_ylim()
ax3.yaxis.set_ticks(np.arange(start, end, (end-start)/5))
# Fourth chart: Plot our cash
ax4 = plt.subplot(414, sharex=ax1)
perf.cash.plot(ax=ax4)
ax4.set_ylabel('Cash\n({})'.format(base_currency))
start, end = ax4.get_ylim()
ax4.yaxis.set_ticks(np.arange(0, end, end/5))
plt.show()
if __name__ == '__main__':
run_algorithm(
capital_base=1000,
data_frequency='minute',
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='bitfinex',
algo_namespace=NAMESPACE,
base_currency='usd',
start=pd.to_datetime('2017-9-22', utc=True),
end=pd.to_datetime('2017-9-23', utc=True),
)
In order to run the code above, you have to ingest the needed data first: In order to run the code above, you have to ingest the needed data first:
+17 -881
View File
@@ -52,35 +52,8 @@ Buy BTC Simple Algorithm
Source code: `examples/buy_btc_simple.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/buy_btc_simple.py>`_ Source code: `examples/buy_btc_simple.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/buy_btc_simple.py>`_
.. code-block:: python .. literalinclude:: ../../catalyst/examples/buy_btc_simple.py
:language: python
'''
Run this example, by executing the following from your terminal:
catalyst ingest-exchange -x bitfinex -f daily -i btc_usdt
catalyst run -f buy_btc_simple.py -x bitfinex --start 2016-1-1 --end 2017-9-30 -o buy_btc_simple_out.pickle
If you want to run this code using another exchange, make sure that
the asset is available on that exchange. For example, if you were to run
it for exchange Poloniex, you would need to edit the following line:
context.asset = symbol('btc_usdt') # note 'usdt' instead of 'usd'
and specify exchange poloniex as follows:
catalyst ingest-exchange -x poloniex -f daily -i btc_usdt
catalyst run -f buy_btc_simple.py -x poloniex --start 2016-1-1 --end 2017-9-30 -o buy_btc_simple_out.pickle
To see which assets are available on each exchange, visit:
https://www.enigma.co/catalyst/status
'''
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'))
This simple algorithm does not produce any output nor displays any chart. This simple algorithm does not produce any output nor displays any chart.
@@ -90,8 +63,6 @@ This simple algorithm does not produce any output nor displays any chart.
Buy and Hodl Algorithm Buy and Hodl Algorithm
~~~~~~~~~~~~~~~~~~~~~~ ~~~~~~~~~~~~~~~~~~~~~~
Source code: `examples/buy_and_hodl.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/buy_and_hodl.py>`_
First ingest the historical pricing data needed to run this algorithm: First ingest the historical pricing data needed to run this algorithm:
.. code-block:: bash .. code-block:: bash
@@ -119,157 +90,10 @@ that 2015-3-1 is the earliest date that Catalyst supports (if you choose an
earlier date, you'll get an error), and the most recent date you can choose is earlier date, you'll get an error), and the most recent date you can choose is
one day prior to the current date. one day prior to the current date.
Source code: `examples/buy_and_hodl.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/buy_and_hodl.py>`_
.. code-block:: python .. literalinclude:: ../../catalyst/examples/buy_and_hodl.py
:language: python
#!/usr/bin/env python
#
# Copyright 2017 Enigma MPC, Inc.
# Copyright 2015 Quantopian, Inc.
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
import pandas as pd
import matplotlib.pyplot as plt
from catalyst import run_algorithm
from catalyst.api import (order_target_value, symbol, record,
cancel_order, get_open_orders, )
def initialize(context):
context.ASSET_NAME = 'btc_usd'
context.TARGET_HODL_RATIO = 0.8
context.RESERVE_RATIO = 1.0 - context.TARGET_HODL_RATIO
context.is_buying = True
context.asset = symbol(context.ASSET_NAME)
context.i = 0
def handle_data(context, data):
context.i += 1
starting_cash = context.portfolio.starting_cash
target_hodl_value = context.TARGET_HODL_RATIO * starting_cash
reserve_value = context.RESERVE_RATIO * starting_cash
# Cancel any outstanding orders
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.current(context.asset, 'price')
# Check if still buying and could (approximately) afford another purchase
if context.is_buying and cash > price:
print('buying')
# Place order to make position in asset equal to target_hodl_value
order_target_value(
context.asset,
target_hodl_value,
limit_price=price * 1.1,
)
record(
price=price,
volume=data.current(context.asset, 'volume'),
cash=cash,
starting_cash=context.portfolio.starting_cash,
leverage=context.account.leverage,
)
def analyze(context=None, results=None):
# Plot the portfolio and asset data.
ax1 = plt.subplot(611)
results[['portfolio_value']].plot(ax=ax1)
ax1.set_ylabel('Portfolio Value (USD)')
ax2 = plt.subplot(612, sharex=ax1)
ax2.set_ylabel('{asset} (USD)'.format(asset=context.ASSET_NAME))
results[['price']].plot(ax=ax2)
trans = results.ix[[t != [] for t in results.transactions]]
buys = trans.ix[
[t[0]['amount'] > 0 for t in trans.transactions]
]
ax2.scatter(
buys.index.to_pydatetime(),
results.price[buys.index],
marker='^',
s=100,
c='g',
label=''
)
ax3 = plt.subplot(613, sharex=ax1)
results[['leverage', 'alpha', 'beta']].plot(ax=ax3)
ax3.set_ylabel('Leverage ')
ax4 = plt.subplot(614, sharex=ax1)
results[['starting_cash', 'cash']].plot(ax=ax4)
ax4.set_ylabel('Cash (USD)')
results[[
'treasury',
'algorithm',
'benchmark',
]] = results[[
'treasury_period_return',
'algorithm_period_return',
'benchmark_period_return',
]]
ax5 = plt.subplot(615, sharex=ax1)
results[[
'treasury',
'algorithm',
'benchmark',
]].plot(ax=ax5)
ax5.set_ylabel('Percent Change')
ax6 = plt.subplot(616, sharex=ax1)
results[['volume']].plot(ax=ax6)
ax6.set_ylabel('Volume (mCoins/5min)')
plt.legend(loc=3)
# Show the plot.
plt.gcf().set_size_inches(18, 8)
plt.show()
if __name__ == '__main__':
run_algorithm(
capital_base=10000,
data_frequency='daily',
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='bitfinex',
algo_namespace='buy_and_hodl',
base_currency='usd',
start=pd.to_datetime('2015-03-01', utc=True),
end=pd.to_datetime('2017-10-31', utc=True),
)
.. image:: https://s3.amazonaws.com/enigmaco-docs/github.io/example_buy_and_hodl.png .. image:: https://s3.amazonaws.com/enigmaco-docs/github.io/example_buy_and_hodl.png
@@ -278,166 +102,13 @@ one day prior to the current date.
Dual Moving Average Crossover Dual Moving Average Crossover
~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ ~~~~~~~~~~~~~~~~~~~~~~~~~~~~~
Source Code: `examples/dual_moving_average.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/dual_moving_average.py>`_
This strategy is covered in detail in the last part of This strategy is covered in detail in the last part of
`this tutorial <beginner-tutorial.html#history>`_. `this tutorial <beginner-tutorial.html#history>`_.
.. code-block:: python Source Code: `examples/dual_moving_average.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/dual_moving_average.py>`_
import numpy as np .. literalinclude:: ../../catalyst/examples/dual_moving_average.py
import pandas as pd :language: python
from logbook import Logger
import matplotlib.pyplot as plt
from catalyst import run_algorithm
from catalyst.api import (order, record, symbol, order_target_percent,
get_open_orders)
from catalyst.exchange.stats_utils import extract_transactions
NAMESPACE = 'dual_moving_average'
log = Logger(NAMESPACE)
def initialize(context):
context.i = 0
context.asset = symbol('ltc_usd')
context.base_price = None
def handle_data(context, data):
# define the windows for the moving averages
short_window = 50
long_window = 200
# Skip as many bars as long_window to properly compute the average
context.i += 1
if context.i < long_window:
return
# Compute moving averages calling data.history() for each
# moving average with the appropriate parameters. We choose to use
# minute bars for this simulation -> freq="1m"
# Returns a pandas dataframe.
short_mavg = data.history(context.asset, 'price',
bar_count=short_window, frequency="1m").mean()
long_mavg = data.history(context.asset, 'price',
bar_count=long_window, frequency="1m").mean()
# Let's keep the price of our asset in a more handy variable
price = data.current(context.asset, 'price')
# If base_price is not set, we use the current value. This is the
# price at the first bar which we reference to calculate price_change.
if context.base_price is None:
context.base_price = price
price_change = (price - context.base_price) / context.base_price
# Save values for later inspection
record(price=price,
cash=context.portfolio.cash,
price_change=price_change,
short_mavg=short_mavg,
long_mavg=long_mavg)
# Since we are using limit orders, some orders may not execute immediately
# we wait until all orders are executed before considering more trades.
orders = get_open_orders(context.asset)
if len(orders) > 0:
return
# Exit if we cannot trade
if not data.can_trade(context.asset):
return
# We check what's our position on our portfolio and trade accordingly
pos_amount = context.portfolio.positions[context.asset].amount
# Trading logic
if short_mavg > long_mavg and pos_amount == 0:
# we buy 100% of our portfolio for this asset
order_target_percent(context.asset, 1)
elif short_mavg < long_mavg and pos_amount > 0:
# we sell all our positions for this asset
order_target_percent(context.asset, 0)
def analyze(context, perf):
# Get the base_currency that was passed as a parameter to the simulation
base_currency = context.exchanges.values()[0].base_currency.upper()
# First chart: Plot portfolio value using base_currency
ax1 = plt.subplot(411)
perf.loc[:, ['portfolio_value']].plot(ax=ax1)
ax1.legend_.remove()
ax1.set_ylabel('Portfolio Value\n({})'.format(base_currency))
start, end = ax1.get_ylim()
ax1.yaxis.set_ticks(np.arange(start, end, (end-start)/5))
# Second chart: Plot asset price, moving averages and buys/sells
ax2 = plt.subplot(412, sharex=ax1)
perf.loc[:, ['price','short_mavg','long_mavg']].plot(ax=ax2, label='Price')
ax2.legend_.remove()
ax2.set_ylabel('{asset}\n({base})'.format(
asset = context.asset.symbol,
base = base_currency
))
start, end = ax2.get_ylim()
ax2.yaxis.set_ticks(np.arange(start, end, (end-start)/5))
transaction_df = extract_transactions(perf)
if not transaction_df.empty:
buy_df = transaction_df[transaction_df['amount'] > 0]
sell_df = transaction_df[transaction_df['amount'] < 0]
ax2.scatter(
buy_df.index.to_pydatetime(),
perf.loc[buy_df.index, 'price'],
marker='^',
s=100,
c='green',
label=''
)
ax2.scatter(
sell_df.index.to_pydatetime(),
perf.loc[sell_df.index, 'price'],
marker='v',
s=100,
c='red',
label=''
)
# Third chart: Compare percentage change between our portfolio
# and the price of the asset
ax3 = plt.subplot(413, sharex=ax1)
perf.loc[:, ['algorithm_period_return', 'price_change']].plot(ax=ax3)
ax3.legend_.remove()
ax3.set_ylabel('Percent Change')
start, end = ax3.get_ylim()
ax3.yaxis.set_ticks(np.arange(start, end, (end-start)/5))
# Fourth chart: Plot our cash
ax4 = plt.subplot(414, sharex=ax1)
perf.cash.plot(ax=ax4)
ax4.set_ylabel('Cash\n({})'.format(base_currency))
start, end = ax4.get_ylim()
ax4.yaxis.set_ticks(np.arange(0, end, end/5))
plt.show()
if __name__ == '__main__':
run_algorithm(
capital_base=1000,
data_frequency='minute',
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='bitfinex',
algo_namespace=NAMESPACE,
base_currency='usd',
start=pd.to_datetime('2017-9-22', utc=True),
end=pd.to_datetime('2017-9-23', utc=True),
)
.. image:: https://s3.amazonaws.com/enigmaco-docs/github.io/tutorial_dual_moving_average.png .. image:: https://s3.amazonaws.com/enigmaco-docs/github.io/tutorial_dual_moving_average.png
@@ -447,8 +118,6 @@ This strategy is covered in detail in the last part of
Mean Reversion Algorithm Mean Reversion Algorithm
~~~~~~~~~~~~~~~~~~~~~~~~ ~~~~~~~~~~~~~~~~~~~~~~~~
Source code: `examples/mean_reversion_simple.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/mean_reversion_simple.py>`_
This algorithm is based on a simple momentum strategy. When the cryptoasset goes This algorithm is based on a simple momentum strategy. When the cryptoasset goes
up quickly, we're going to buy; when it goes down quickly, we're going to sell. up quickly, we're going to buy; when it goes down quickly, we're going to sell.
Hopefully, we'll ride the waves. Hopefully, we'll ride the waves.
@@ -469,284 +138,10 @@ lines 218-245, so in order to run the algorithm we just type:
python mean_reversion_simple.py python mean_reversion_simple.py
.. code-block:: python Source code: `examples/mean_reversion_simple.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/mean_reversion_simple.py>`_
import os .. literalinclude:: ../../catalyst/examples/mean_reversion_simple.py
import tempfile :language: python
import time
import numpy as np
import pandas as pd
import talib
from logbook import Logger
from catalyst import run_algorithm
from catalyst.api import symbol, record, order_target_percent, get_open_orders
from catalyst.exchange.stats_utils import extract_transactions
# We give a name to the algorithm which Catalyst will use to persist its state.
# In this example, Catalyst will create the `.catalyst/data/live_algos`
# directory. If we stop and start the algorithm, Catalyst will resume its
# state using the files included in the folder.
from catalyst.utils.paths import ensure_directory
NAMESPACE = 'mean_reversion_simple'
log = Logger(NAMESPACE)
# To run an algorithm in Catalyst, you need two functions: initialize and
# handle_data.
def initialize(context):
# This initialize function sets any data or variables that you'll use in
# your algorithm. For instance, you'll want to define the trading pair (or
# trading pairs) you want to backtest. You'll also want to define any
# parameters or values you're going to use.
# In our example, we're looking at Neo in USD.
context.neo_eth = symbol('neo_usd')
context.base_price = None
context.current_day = None
context.RSI_OVERSOLD = 30
context.RSI_OVERBOUGHT = 80
context.CANDLE_SIZE = '15T'
context.start_time = time.time()
def handle_data(context, data):
# This handle_data function is where the real work is done. Our data is
# minute-level tick data, and each minute is called a frame. This function
# runs on each frame of the data.
# We flag the first period of each day.
# Since cryptocurrencies trade 24/7 the `before_trading_starts` handle
# would only execute once. This method works with minute and daily
# frequencies.
today = data.current_dt.floor('1D')
if today != context.current_day:
context.traded_today = False
context.current_day = today
# We're computing the volume-weighted-average-price of the security
# defined above, in the context.neo_eth variable. For this example, we're
# using three bars on the 15 min bars.
# The frequency attribute determine the bar size. We use this convention
# for the frequency alias:
# http://pandas.pydata.org/pandas-docs/stable/timeseries.html#offset-aliases
prices = data.history(
context.neo_eth,
fields='close',
bar_count=50,
frequency=context.CANDLE_SIZE
)
# Ta-lib calculates various technical indicator based on price and
# volume arrays.
# In this example, we are comp
rsi = talib.RSI(prices.values, timeperiod=14)
# We need a variable for the current price of the security to compare to
# the average. Since we are requesting two fields, data.current()
# returns a DataFrame with
current = data.current(context.neo_eth, fields=['close', 'volume'])
price = current['close']
# If base_price is not set, we use the current value. This is the
# price at the first bar which we reference to calculate price_change.
if context.base_price is None:
context.base_price = price
price_change = (price - context.base_price) / context.base_price
cash = context.portfolio.cash
# Now that we've collected all current data for this frame, we use
# the record() method to save it. This data will be available as
# a parameter of the analyze() function for further analysis.
record(
price=price,
volume=current['volume'],
price_change=price_change,
rsi=rsi[-1],
cash=cash
)
# We are trying to avoid over-trading by limiting our trades to
# one per day.
if context.traded_today:
return
# Since we are using limit orders, some orders may not execute immediately
# we wait until all orders are executed before considering more trades.
orders = get_open_orders(context.neo_eth)
if len(orders) > 0:
return
# Exit if we cannot trade
if not data.can_trade(context.neo_eth):
return
# Another powerful built-in feature of the Catalyst backtester is the
# portfolio object. The portfolio object tracks your positions, cash,
# cost basis of specific holdings, and more. In this line, we calculate
# how long or short our position is at this minute.
pos_amount = context.portfolio.positions[context.neo_eth].amount
if rsi[-1] <= context.RSI_OVERSOLD and pos_amount == 0:
log.info(
'{}: buying - price: {}, rsi: {}'.format(
data.current_dt, price, rsi[-1]
)
)
# Set a style for limit orders,
limit_price = price * 1.005
order_target_percent(
context.neo_eth, 1, limit_price=limit_price
)
context.traded_today = True
elif rsi[-1] >= context.RSI_OVERBOUGHT and pos_amount > 0:
log.info(
'{}: selling - price: {}, rsi: {}'.format(
data.current_dt, price, rsi[-1]
)
)
limit_price = price * 0.995
order_target_percent(
context.neo_eth, 0, limit_price=limit_price
)
context.traded_today = True
def analyze(context=None, perf=None):
end = time.time()
log.info('elapsed time: {}'.format(end - context.start_time))
import matplotlib.pyplot as plt
# The base currency of the algo exchange
base_currency = context.exchanges.values()[0].base_currency.upper()
# Plot the portfolio value over time.
ax1 = plt.subplot(611)
perf.loc[:, 'portfolio_value'].plot(ax=ax1)
ax1.set_ylabel('Portfolio\nValue\n({})'.format(base_currency))
# Plot the price increase or decrease over time.
ax2 = plt.subplot(612, sharex=ax1)
perf.loc[:, 'price'].plot(ax=ax2, label='Price')
ax2.set_ylabel('{asset}\n({base})'.format(
asset=context.neo_eth.symbol, base=base_currency
))
transaction_df = extract_transactions(perf)
if not transaction_df.empty:
buy_df = transaction_df[transaction_df['amount'] > 0]
sell_df = transaction_df[transaction_df['amount'] < 0]
ax2.scatter(
buy_df.index.to_pydatetime(),
perf.loc[buy_df.index.floor('1 min'), 'price'],
marker='^',
s=100,
c='green',
label=''
)
ax2.scatter(
sell_df.index.to_pydatetime(),
perf.loc[sell_df.index.floor('1 min'), 'price'],
marker='v',
s=100,
c='red',
label=''
)
ax4 = plt.subplot(613, sharex=ax1)
perf.loc[:, 'cash'].plot(
ax=ax4, label='Base Currency ({})'.format(base_currency)
)
ax4.set_ylabel('Cash\n({})'.format(base_currency))
perf['algorithm'] = perf.loc[:, 'algorithm_period_return']
ax5 = plt.subplot(614, sharex=ax1)
perf.loc[:, ['algorithm', 'price_change']].plot(ax=ax5)
ax5.set_ylabel('Percent\nChange')
ax6 = plt.subplot(615, sharex=ax1)
perf.loc[:, 'rsi'].plot(ax=ax6, label='RSI')
ax6.set_ylabel('RSI')
ax6.axhline(context.RSI_OVERBOUGHT, color='darkgoldenrod')
ax6.axhline(context.RSI_OVERSOLD, color='darkgoldenrod')
if not transaction_df.empty:
ax6.scatter(
buy_df.index.to_pydatetime(),
perf.loc[buy_df.index.floor('1 min'), 'rsi'],
marker='^',
s=100,
c='green',
label=''
)
ax6.scatter(
sell_df.index.to_pydatetime(),
perf.loc[sell_df.index.floor('1 min'), 'rsi'],
marker='v',
s=100,
c='red',
label=''
)
plt.legend(loc=3)
start, end = ax6.get_ylim()
ax6.yaxis.set_ticks(np.arange(0, end, end/5))
# Show the plot.
plt.gcf().set_size_inches(18, 8)
plt.show()
pass
if __name__ == '__main__':
# The execution mode: backtest or live
MODE = 'backtest'
if MODE == 'backtest':
folder = os.path.join(
tempfile.gettempdir(), 'catalyst', NAMESPACE
)
ensure_directory(folder)
timestr = time.strftime('%Y%m%d-%H%M%S')
out = os.path.join(folder, '{}.p'.format(timestr))
# catalyst run -f catalyst/examples/mean_reversion_simple.py -x bitfinex -s 2017-10-1 -e 2017-11-10 -c usdt -n mean-reversion --data-frequency minute --capital-base 10000
run_algorithm(
capital_base=10000,
data_frequency='minute',
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='bitfinex',
algo_namespace=NAMESPACE,
base_currency='usd',
start=pd.to_datetime('2017-10-01', utc=True),
end=pd.to_datetime('2017-11-10', utc=True),
output=out
)
log.info('saved perf stats: {}'.format(out))
elif MODE == 'live':
run_algorithm(
capital_base=0.5,
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='bittrex',
live=True,
algo_namespace=NAMESPACE,
base_currency='usd',
live_graph=False
)
.. image:: https://s3.amazonaws.com/enigmaco-docs/github.io/example_mean_reversion_simple.png .. image:: https://s3.amazonaws.com/enigmaco-docs/github.io/example_mean_reversion_simple.png
@@ -763,8 +158,6 @@ strategy.
Simple Universe Simple Universe
~~~~~~~~~~~~~~~ ~~~~~~~~~~~~~~~
Source code: `examples/simple_universe.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/simple_universe.py>`_
This example aims to provide an easy way for users to learn how to This example aims to provide an easy way for users to learn how to
collect data from any given exchange and select a subset of the available collect data from any given exchange and select a subset of the available
currency pairs for trading. You simply need to specify the exchange and currency pairs for trading. You simply need to specify the exchange and
@@ -791,142 +184,10 @@ of the file:
catalyst ingest-exchange -x bitfinex -f minute catalyst ingest-exchange -x bitfinex -f minute
.. code-block:: bash Source code: `examples/simple_universe.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/simple_universe.py>`_
python simple_universe.py
Credits: This code was originally submitted by `Abner Ayala-Acevedo
<https://github.com/abnera>`_. Thank you!
.. code-block:: python
from datetime import timedelta
import numpy as np
import pandas as pd
from catalyst import run_algorithm
from catalyst.exchange.utils.exchange_utils import get_exchange_symbols
from catalyst.api import (symbols, )
def initialize(context):
context.i = -1 # minute counter
context.exchange = context.exchanges.values()[0].name.lower()
context.base_currency = context.exchanges.values()[0].base_currency.lower()
def handle_data(context, data):
context.i += 1
lookback_days = 7 # 7 days
# current date & time in each iteration formatted into a string
now = data.current_dt
date, time = now.strftime('%Y-%m-%d %H:%M:%S').split(' ')
lookback_date = now - timedelta(days=lookback_days)
# keep only the date as a string, discard the time
lookback_date = lookback_date.strftime('%Y-%m-%d %H:%M:%S').split(' ')[0]
one_day_in_minutes = 1440 # 60 * 24 assumes data_frequency='minute'
# update universe everyday at midnight
if not context.i % one_day_in_minutes:
context.universe = universe(context, lookback_date, date)
# get data every 30 minutes
minutes = 30
# get lookback_days of history data: that is 'lookback' number of bins
lookback = one_day_in_minutes / minutes * lookback_days
if not context.i % minutes and context.universe:
# we iterate for every pair in the current universe
for coin in context.coins:
pair = str(coin.symbol)
# Get 30 minute interval OHLCV data. This is the standard data
# required for candlestick or indicators/signals. Return Pandas
# DataFrames. 30T means 30-minute re-sampling of one minute data.
# Adjust it to your desired time interval as needed.
opened = fill(data.history(coin, 'open',
bar_count=lookback, frequency='30T')).values
high = fill(data.history(coin, 'high',
bar_count=lookback, frequency='30T')).values
low = fill(data.history(coin, 'low',
bar_count=lookback, frequency='30T')).values
close = fill(data.history(coin, 'price',
bar_count=lookback, frequency='30T')).values
volume = fill(data.history(coin, 'volume',
bar_count=lookback, frequency='30T')).values
# close[-1] is the last value in the set, which is the equivalent
# to current price (as in the most recent value)
# displays the minute price for each pair every 30 minutes
print('{now}: {pair} -\tO:{o},\tH:{h},\tL:{c},\tC{c},\tV:{v}'.format(
now=now,
pair=pair,
o=opened[-1],
h=high[-1],
l=low[-1],
c=close[-1],
v=volume[-1],
))
# -------------------------------------------------------------
# --------------- Insert Your Strategy Here -------------------
# -------------------------------------------------------------
def analyze(context=None, results=None):
pass
# Get the universe for a given exchange and a given base_currency market
# Example: Poloniex BTC Market
def universe(context, lookback_date, current_date):
# get all the pairs for the given exchange
json_symbols = get_exchange_symbols(context.exchange)
# convert into a DataFrame for easier processing
df = pd.DataFrame.from_dict(json_symbols).transpose().astype(str)
df['base_currency'] = df.apply(lambda row: row.symbol.split('_')[1],axis=1)
df['market_currency'] = df.apply(lambda row: row.symbol.split('_')[0],axis=1)
# Filter all the pairs to get only the ones for a given base_currency
df = df[df['base_currency'] == context.base_currency]
# Filter all the pairs to ensure that pair existed in the current date range
df = df[df.start_date < lookback_date]
df = df[df.end_daily >= current_date]
context.coins = symbols(*df.symbol) # convert all the pairs to symbols
return df.symbol.tolist()
# Replace all NA, NAN or infinite values with its nearest value
def fill(series):
if isinstance(series, pd.Series):
return series.replace([np.inf, -np.inf], np.nan).ffill().bfill()
elif isinstance(series, np.ndarray):
return pd.Series(series).replace(
[np.inf, -np.inf], np.nan
).ffill().bfill().values
else:
return series
if __name__ == '__main__':
start_date = pd.to_datetime('2017-11-10', utc=True)
end_date = pd.to_datetime('2017-11-13', utc=True)
performance = run_algorithm(start=start_date, end=end_date,
capital_base=100.0, # amount of base_currency
initialize=initialize,
handle_data=handle_data,
analyze=analyze,
exchange_name='bitfinex',
data_frequency='minute',
base_currency='btc',
live=False,
live_graph=False,
algo_namespace='simple_universe')
.. literalinclude:: ../../catalyst/examples/simple_universe.py
:language: python
.. _portfolio_optimization: .. _portfolio_optimization:
@@ -940,135 +201,10 @@ use 180 days of historical data and rebalance every 30 days. This code was used
in writting the following article: in writting the following article:
`Markowitz Portfolio Optimization for Cryptocurrencies <https://blog.enigma.co/markowitz-portfolio-optimization-for-cryptocurrencies-in-catalyst-b23c38652556>`_. `Markowitz Portfolio Optimization for Cryptocurrencies <https://blog.enigma.co/markowitz-portfolio-optimization-for-cryptocurrencies-in-catalyst-b23c38652556>`_.
.. code-block:: python Source code: `examples/simple_universe.py <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/portfolio_optimization.py>`_
''' .. literalinclude:: ../../catalyst/examples/portfolio_optimization.py
You can run this code using the Python interpreter: :language: python
$ python portfolio_optimization.py
'''
from __future__ import division
import os
import pytz
import numpy as np
import pandas as pd
from scipy.optimize import minimize
import matplotlib.pyplot as plt
from datetime import datetime
from catalyst.api import record, symbol, symbols, order_target_percent
from catalyst.utils.run_algo import run_algorithm
np.set_printoptions(threshold='nan', suppress=True)
def initialize(context):
# Portfolio assets list
context.assets = symbols('btc_usdt', 'eth_usdt', 'ltc_usdt', 'dash_usdt',
'xmr_usdt')
context.nassets = len(context.assets)
# Set the time window that will be used to compute expected return
# and asset correlations
context.window = 180
# Set the number of days between each portfolio rebalancing
context.rebalance_period = 30
context.i = 0
def handle_data(context, data):
# Only rebalance at the beggining of the algorithm execution and
# every multiple of the rebalance period
if context.i == 0 or context.i%context.rebalance_period == 0:
n = context.window
prices = data.history(context.assets, fields='price',
bar_count=n+1, frequency='1d')
pr = np.asmatrix(prices)
t_prices = prices.iloc[1:n+1]
t_val = t_prices.values
tminus_prices = prices.iloc[0:n]
tminus_val = tminus_prices.values
# Compute daily returns (r)
r = np.asmatrix(t_val/tminus_val-1)
# Compute the expected returns of each asset with the average
# daily return for the selected time window
m = np.asmatrix(np.mean(r, axis=0))
# ###
stds = np.std(r, axis=0)
# Compute excess returns matrix (xr)
xr = r - m
# Matrix algebra to get variance-covariance matrix
cov_m = np.dot(np.transpose(xr),xr)/n
# Compute asset correlation matrix (informative only)
corr_m = cov_m/np.dot(np.transpose(stds),stds)
# Define portfolio optimization parameters
n_portfolios = 50000
results_array = np.zeros((3+context.nassets,n_portfolios))
for p in xrange(n_portfolios):
weights = np.random.random(context.nassets)
weights /= np.sum(weights)
w = np.asmatrix(weights)
p_r = np.sum(np.dot(w,np.transpose(m)))*365
p_std = np.sqrt(np.dot(np.dot(w,cov_m),np.transpose(w)))*np.sqrt(365)
#store results in results array
results_array[0,p] = p_r
results_array[1,p] = p_std
#store Sharpe Ratio (return / volatility) - risk free rate element
#excluded for simplicity
results_array[2,p] = results_array[0,p] / results_array[1,p]
i = 0
for iw in weights:
results_array[3+i,p] = weights[i]
i += 1
#convert results array to Pandas DataFrame
results_frame = pd.DataFrame(np.transpose(results_array),
columns=['r','stdev','sharpe']+context.assets)
#locate position of portfolio with highest Sharpe Ratio
max_sharpe_port = results_frame.iloc[results_frame['sharpe'].idxmax()]
#locate positon of portfolio with minimum standard deviation
min_vol_port = results_frame.iloc[results_frame['stdev'].idxmin()]
#order optimal weights for each asset
for asset in context.assets:
if data.can_trade(asset):
order_target_percent(asset, max_sharpe_port[asset])
#create scatter plot coloured by Sharpe Ratio
plt.scatter(results_frame.stdev,results_frame.r,c=results_frame.sharpe,cmap='RdYlGn')
plt.xlabel('Volatility')
plt.ylabel('Returns')
plt.colorbar()
#plot red star to highlight position of portfolio with highest Sharpe Ratio
plt.scatter(max_sharpe_port[1],max_sharpe_port[0],marker='o',color='b',s=200)
#plot green star to highlight position of minimum variance portfolio
plt.show()
print(max_sharpe_port)
record(pr=pr,r=r, m=m, stds=stds ,max_sharpe_port=max_sharpe_port, corr_m=corr_m)
context.i += 1
def analyze(context=None, results=None):
# Form DataFrame with selected data
data = results[['pr','r','m','stds','max_sharpe_port','corr_m','portfolio_value']]
# Save results in CSV file
filename = os.path.splitext(os.path.basename(__file__))[0]
data.to_csv(filename + '.csv')
# Bitcoin data is available from 2015-3-2. Dates vary for other tokens.
start = datetime(2017, 1, 1, 0, 0, 0, 0, pytz.utc)
end = datetime(2017, 8, 16, 0, 0, 0, 0, pytz.utc)
results = run_algorithm(initialize=initialize,
handle_data=handle_data,
analyze=analyze,
start=start,
end=end,
exchange_name='poloniex',
capital_base=100000, )
.. image:: https://cdn-images-1.medium.com/max/1600/0*EjjiKZHlYF3sn7yQ. .. image:: https://cdn-images-1.medium.com/max/1600/0*EjjiKZHlYF3sn7yQ.
:align: center :align: center
+9 -1
View File
@@ -89,7 +89,7 @@ Once either Conda or MiniConda has been set up you can install Catalyst:
.. code-block:: bash .. code-block:: bash
conda env create -f python2.7-environment.yml conda env create -f python2.7-environment.yml
4. Activate the environment (which you need to do every time you start a new 4. Activate the environment (which you need to do every time you start a new
session to run Catalyst): session to run Catalyst):
@@ -133,6 +133,14 @@ with the following steps:
2. Create the environment: 2. Create the environment:
for python 2.7:
.. code-block:: bash
conda create --name catalyst python=2.7 scipy zlib
or for python 3.6:
.. code-block:: bash .. code-block:: bash
conda create --name catalyst python=2.7 scipy zlib conda create --name catalyst python=2.7 scipy zlib
+9
View File
@@ -184,5 +184,14 @@ Here is the breakdown of the new arguments:
essentially sleep and when the predefined time comes, it would start executing. essentially sleep and when the predefined time comes, it would start executing.
The `catalyst live` command offers additional parameters.
You can learn more by running the following from the command line:
.. code-block:: bash
catalyst live --help
Here is a complete algorithm for reference: Here is a complete algorithm for reference:
`Buy Low and Sell High <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/buy_low_sell_high_live.py>`_ `Buy Low and Sell High <https://github.com/enigmampc/catalyst/blob/master/catalyst/examples/buy_low_sell_high_live.py>`_
+6 -2
View File
@@ -1,9 +1,11 @@
name: catalyst name: catalyst
channels: channels:
- defaults - defaults
- conda-forge
dependencies: dependencies:
- certifi=2016.2.28=py27_0 - certifi=2016.2.28=py27_0
- mkl=2017.0.3 - mkl=2017.0.3
- matplotlib=2.1.2=py36_0
- numpy=1.13.1=py27_0 - numpy=1.13.1=py27_0
- openssl=1.0.2l - openssl=1.0.2l
- pip=9.0.1=py27_1 - pip=9.0.1=py27_1
@@ -20,8 +22,10 @@ dependencies:
- bcolz==0.12.1 - bcolz==0.12.1
- bottleneck==1.2.1 - bottleneck==1.2.1
- chardet==3.0.4 - chardet==3.0.4
- ccxt==1.10.1094 - ccxt==1.11.22
- web3==4.0.0b7 # The Enigma Data Marketplace requires Python3 because it depends on
# web3, which requires Python3, as building its dependencies breaks in Python2
# - web3==4.0.0b7
- requests-toolbelt==0.8.0 - requests-toolbelt==0.8.0
- click==6.7 - click==6.7
- contextlib2==0.5.5 - contextlib2==0.5.5
+17 -22
View File
@@ -1,29 +1,24 @@
name: catalyst name: catalyst
channels: channels:
- defaults - defaults
- conda-forge
dependencies: dependencies:
- ca-certificates=2017.08.26=ha1e5d58_0 - ca-certificates=2017.08.26
- certifi=2018.1.18=py36_0 - certifi=2018.1.18
- intel-openmp=2018.0.0=h8158457_8 - intel-openmp=2018.0.0
- libcxx=4.0.1=h579ed51_0 - mkl=2018.0.1
- libcxxabi=4.0.1=hebd6815_0 - numpy=1.14.0
- libedit=3.1=hb4e282d_0 - openssl=1.0.2n
- libffi=3.2.1=h475c297_4 - matplotlib=2.1.2=py36_0
- libgfortran=3.0.1=h93005f0_2 - pip=9.0.1
- mkl=2018.0.1=hfbd8650_4 - python=3.6.4
- ncurses=6.0=hd04f020_2 - scipy=1.0.0
- numpy=1.14.0=py36h8a80b8c_1
- openssl=1.0.2n=hdbc3d79_0
- pip=9.0.1=py36h1555ced_4
- python=3.6.4=hc167b69_1
- readline=7.0=hc1231fa_4
- scipy=1.0.0=py36h1de22e9_0
- setuptools=38.4.0=py36_0 - setuptools=38.4.0=py36_0
- sqlite=3.22.0=h3efe00b_0 - sqlite=3.22.0
- tk=8.6.7=h35a86e2_3 - tk=8.6.7
- wheel=0.30.0=py36h5eb2c71_1 - wheel=0.30.0
- xz=5.2.3=h0278029_2 - xz=5.2.3
- zlib=1.2.11=hf3cbc9b_2 - zlib=1.2.11
- pip: - pip:
- aiodns==1.1.1 - aiodns==1.1.1
- aiohttp==3.0.1 - aiohttp==3.0.1
@@ -36,7 +31,7 @@ dependencies:
- botocore==1.8.41 - botocore==1.8.41
- bottleneck==1.2.1 - bottleneck==1.2.1
- cchardet==2.1.1 - cchardet==2.1.1
- ccxt==1.10.1102 - ccxt==1.11.22
- chardet==3.0.4 - chardet==3.0.4
- click==6.7 - click==6.7
- contextlib2==0.5.5 - contextlib2==0.5.5
+2 -2
View File
@@ -81,8 +81,8 @@ empyrical==0.2.1
tables==3.3.0 tables==3.3.0
#Catalyst dependencies #Catalyst dependencies
ccxt==1.10.1094 ccxt==1.11.22
boto3==1.4.8 boto3==1.4.8
redo==1.6 redo==1.6
web3==4.0.0b7 web3==4.0.0b11; python_version > '3.4'
requests-toolbelt==0.8.0 requests-toolbelt==0.8.0
-2
View File
@@ -1,2 +0,0 @@
web3==4.0.0b7
requests-toolbelt==0.8.0
+35 -17
View File
@@ -11,10 +11,11 @@ from catalyst.exchange.exchange_bundle import ExchangeBundle, \
BUNDLE_NAME_TEMPLATE BUNDLE_NAME_TEMPLATE
from catalyst.exchange.utils.bundle_utils import get_bcolz_chunk, \ from catalyst.exchange.utils.bundle_utils import get_bcolz_chunk, \
get_df_from_arrays get_df_from_arrays
from exchange.utils.datetime_utils import get_start_dt from catalyst.exchange.utils.datetime_utils import get_start_dt
from catalyst.exchange.utils.exchange_utils import get_exchange_folder from catalyst.exchange.utils.exchange_utils import get_exchange_folder
from catalyst.exchange.utils.factory import get_exchange from catalyst.exchange.utils.factory import get_exchange
from catalyst.exchange.utils.stats_utils import df_to_string from catalyst.exchange.utils.stats_utils import df_to_string, \
set_print_settings
from catalyst.utils.paths import ensure_directory from catalyst.utils.paths import ensure_directory
log = getLogger('test_exchange_bundle') log = getLogger('test_exchange_bundle')
@@ -42,16 +43,16 @@ class TestExchangeBundle:
def test_ingest_minute(self): def test_ingest_minute(self):
data_frequency = 'minute' data_frequency = 'minute'
exchange_name = 'poloniex' exchange_name = 'binance'
exchange = get_exchange(exchange_name) exchange = get_exchange(exchange_name)
exchange_bundle = ExchangeBundle(exchange) exchange_bundle = ExchangeBundle(exchange_name)
assets = [ assets = [
exchange.get_asset('eth_btc') exchange.get_asset('bch_eth')
] ]
start = pd.to_datetime('2016-03-01', utc=True) start = pd.to_datetime('2018-03-01', utc=True)
end = pd.to_datetime('2017-11-1', utc=True) end = pd.to_datetime('2018-03-8', utc=True)
log.info('ingesting exchange bundle {}'.format(exchange_name)) log.info('ingesting exchange bundle {}'.format(exchange_name))
exchange_bundle.ingest( exchange_bundle.ingest(
@@ -61,7 +62,8 @@ class TestExchangeBundle:
exclude_symbols=None, exclude_symbols=None,
start=start, start=start,
end=end, end=end,
show_progress=True show_progress=False,
show_breakdown=False
) )
reader = exchange_bundle.get_reader(data_frequency) reader = exchange_bundle.get_reader(data_frequency)
@@ -72,9 +74,15 @@ class TestExchangeBundle:
start_dt=start, start_dt=start,
end_dt=end end_dt=end
) )
print('found {} rows for {} ingestion\n{}'.format( periods = exchange_bundle.get_calendar_periods_range(
len(arrays[0]), asset.symbol, arrays[0]) start, end, data_frequency
) )
dx = get_df_from_arrays(arrays[0], periods)
set_print_settings()
print('found {} rows for last ingestion:\n{}\n{}'.format(
len(dx), dx.head(10), dx.tail(10)
))
pass pass
def test_ingest_minute_all(self): def test_ingest_minute_all(self):
@@ -101,7 +109,7 @@ class TestExchangeBundle:
# data_frequency = 'daily' # data_frequency = 'daily'
# include_symbols = 'neo_btc,bch_btc,eth_btc' # include_symbols = 'neo_btc,bch_btc,eth_btc'
exchange_name = 'bitfinex' exchange_name = 'binance'
data_frequency = 'minute' data_frequency = 'minute'
exchange = get_exchange(exchange_name) exchange = get_exchange(exchange_name)
@@ -222,9 +230,14 @@ class TestExchangeBundle:
start_dt=start, start_dt=start,
end_dt=end end_dt=end
) )
print('found {} rows for {} ingestion\n{}'.format( periods = exchange_bundle.get_calendar_periods_range(
len(arrays[0]), asset.symbol, arrays[0]) start, end, data_frequency
) )
dx = get_df_from_arrays(arrays, periods)
print('found {} rows for last ingestion'.format(
len(dx)
))
pass pass
def test_daily_data_to_minute_table(self): def test_daily_data_to_minute_table(self):
@@ -290,17 +303,22 @@ class TestExchangeBundle:
for asset in assets: for asset in assets:
sid = asset.sid sid = asset.sid
daily_values = reader.load_raw_arrays( arrays = reader.load_raw_arrays(
fields=['open', 'high', 'low', 'close', 'volume'], fields=['open', 'high', 'low', 'close', 'volume'],
start_dt=start, start_dt=start,
end_dt=end, end_dt=end,
sids=[sid], sids=[sid],
) )
print('found {} rows for last ingestion'.format( periods = exchange_bundle.get_calendar_periods_range(
len(daily_values[0])) start, end, data_frequency
) )
pass
dx = get_df_from_arrays(arrays, periods)
print('found {} rows for last ingestion'.format(
len(dx)
))
pass
def test_minute_bundle(self): def test_minute_bundle(self):
# exchange_name = 'poloniex' # exchange_name = 'poloniex'
+37 -4
View File
@@ -5,7 +5,8 @@ from catalyst.exchange.utils.stats_utils import set_print_settings
from .base import BaseExchangeTestCase from .base import BaseExchangeTestCase
from catalyst.exchange.ccxt.ccxt_exchange import CCXT from catalyst.exchange.ccxt.ccxt_exchange import CCXT
from catalyst.exchange.exchange_execution import ExchangeLimitOrder from catalyst.exchange.exchange_execution import ExchangeLimitOrder
from catalyst.exchange.utils.exchange_utils import get_exchange_auth from catalyst.exchange.utils.exchange_utils import get_exchange_auth, \
get_trades_df, candles_from_trades
from catalyst.finance.order import Order from catalyst.finance.order import Order
log = Logger('test_ccxt') log = Logger('test_ccxt')
@@ -14,12 +15,13 @@ log = Logger('test_ccxt')
class TestCCXT(BaseExchangeTestCase): class TestCCXT(BaseExchangeTestCase):
@classmethod @classmethod
def setup(self): def setup(self):
exchange_name = 'bittrex' exchange_name = 'binance'
auth = get_exchange_auth(exchange_name) auth = get_exchange_auth(exchange_name)
self.exchange = CCXT( self.exchange = CCXT(
exchange_name=exchange_name, exchange_name=exchange_name,
key=auth['key'], key=auth['key'],
secret=auth['secret'], secret=auth['secret'],
password=None,
base_currency='usdt', base_currency='usdt',
) )
self.exchange.init() self.exchange.init()
@@ -58,9 +60,9 @@ class TestCCXT(BaseExchangeTestCase):
log.info('retrieving candles') log.info('retrieving candles')
candles = self.exchange.get_candles( candles = self.exchange.get_candles(
freq='1T', freq='1T',
assets=[self.exchange.get_asset('eth_btc')], assets=[self.exchange.get_asset('eng_eth')],
bar_count=200, bar_count=200,
# start_dt=pd.to_datetime('2017-09-01', utc=True), start_dt=pd.to_datetime('2017-09-01', utc=True),
) )
for asset in candles: for asset in candles:
@@ -90,6 +92,37 @@ class TestCCXT(BaseExchangeTestCase):
assert trades assert trades
pass pass
def test_validate_volume(self):
asset = self.exchange.get_asset('eng_eth')
candles = self.exchange.get_candles(
freq='1T',
assets=[asset],
bar_count=10,
)
df = pd.DataFrame(candles[asset])
df.set_index('last_traded', drop=True, inplace=True)
df.drop_duplicates()
df.sort_index(inplace=True, ascending=False)
assert candles
start_dt = df.index[-1]
trades = self.exchange.get_trades(
asset, start_dt=start_dt, my_trades=False
)
assert trades
trades_df = get_trades_df(trades)
df2 = candles_from_trades(trades_df, '1T')
set_print_settings()
log.info(
'comparing candles / resampled trades:\n{}\n{}'.format(
df, df2
)
)
pass
def test_get_executed_order(self): def test_get_executed_order(self):
log.info('retrieving executed order') log.info('retrieving executed order')
asset = self.exchange.get_asset('eng_eth') asset = self.exchange.get_asset('eng_eth')
+1 -1
View File
@@ -9,7 +9,7 @@ from catalyst.exchange.exchange_data_portal import (
) )
from catalyst.exchange.utils.exchange_utils import get_common_assets from catalyst.exchange.utils.exchange_utils import get_common_assets
from catalyst.exchange.utils.factory import get_exchanges from catalyst.exchange.utils.factory import get_exchanges
from test_utils import rnd_history_date_days, rnd_bar_count from .test_utils import rnd_history_date_days, rnd_bar_count
log = Logger('test_bitfinex') log = Logger('test_bitfinex')
+175
View File
@@ -0,0 +1,175 @@
from catalyst.exchange.utils.exchange_utils import transform_candles_to_df, \
forward_fill_df_if_needed, get_candles_df
from catalyst.testing.fixtures import WithLogger, ZiplineTestCase
from datetime import timedelta
from pandas import Timestamp, DataFrame, concat
import numpy as np
class TestExchangeUtils(WithLogger, ZiplineTestCase):
@classmethod
def get_specific_field_from_df(cls, df, field, asset):
new_df = DataFrame(df[field])
new_df.columns = [asset]
new_df.index.name = None
return new_df
@classmethod
def verify_forward_fill_df_if_needed(cls, candles, periods, expected_df):
observed_df = forward_fill_df_if_needed(
transform_candles_to_df(candles),
periods)
assert (expected_df.equals(observed_df))
@classmethod
def verify_get_candles_df(cls, assets, candles, end_fixed_dt,
expected_df, check_next_candle=False):
# run on all the fields
for field in ['volume', 'open', 'close', 'high', 'low']:
field_dt = cls.get_specific_field_from_df(expected_df,
field,
assets[0])
# run on several timestamps
for delta in range(5):
end_dt = end_fixed_dt + timedelta(minutes=delta)
assert (field_dt.equals(get_candles_df({assets[0]: candles},
field, '5T', 3,
end_dt=end_dt)))
field_dt_a1 = cls.get_specific_field_from_df(expected_df,
field,
assets[0])
field_dt_a2 = cls.get_specific_field_from_df(expected_df,
field,
assets[1])
observed_df = get_candles_df({assets[0]: candles,
assets[1]: candles},
field, '5T', 3,
end_dt=end_dt)
assert (observed_df.equals(concat([field_dt_a1, field_dt_a2],
axis=1)))
if check_next_candle:
# one candle forward
end_dt = end_fixed_dt + timedelta(minutes=6)
observed_df = get_candles_df({assets[0]: candles,
assets[1]: candles},
field, '5T', 3,
end_dt=end_dt)
assert (not observed_df.equals(concat([field_dt_a1,
field_dt_a2],
axis=1)))
assert (concat([field_dt_a1, field_dt_a2],
axis=1)[1:].equals(observed_df[:-1]))
def test_get_candles_df(self):
assets = ['btc_usdt', 'eth_usdt']
# test forward fill in the end
candles = [{'high': 595, 'volume': 10, 'low': 594,
'close': 595, 'open': 594,
'last_traded': Timestamp('2018-03-01 09:45:00+0000',
tz='UTC')
},
{'high': 594, 'volume': 108, 'low': 592,
'close': 593, 'open': 592,
'last_traded': Timestamp('2018-03-01 09:50:00+0000',
tz='UTC')
}]
expected = [{'high': 595.0, 'volume': 10.0, 'low': 594.0,
'close': 595.0, 'open': 594.0,
'last_traded': Timestamp('2018-03-01 09:45:00+0000',
tz='UTC')
},
{'high': 594.0, 'volume': 108.0, 'low': 592.0,
'close': 593.0, 'open': 592.0,
'last_traded': Timestamp('2018-03-01 09:50:00+0000',
tz='UTC')
},
{'high': 593.0, 'volume': 0.0, 'low': 593.0,
'close': 593.0, 'open': 593.0,
'last_traded': Timestamp('2018-03-01 09:55:00+0000',
tz='UTC')
}]
periods = [Timestamp('2018-03-01 09:45:00+0000', tz='UTC'),
Timestamp('2018-03-01 09:50:00+0000', tz='UTC'),
Timestamp('2018-03-01 09:55:00+0000', tz='UTC')]
expected_df = transform_candles_to_df(expected)
self.verify_forward_fill_df_if_needed(candles, periods,
expected_df)
self.verify_get_candles_df(assets, candles, periods[2],
expected_df, True)
# test forward fill in the middle
candles = [{'high': 595, 'volume': 10, 'low': 594,
'close': 595, 'open': 594,
'last_traded': Timestamp('2018-03-01 09:45:00+0000',
tz='UTC')
},
{'high': 594, 'volume': 108, 'low': 592,
'close': 593, 'open': 592,
'last_traded': Timestamp('2018-03-01 09:55:00+0000',
tz='UTC')
}]
expected = [{'high': 595.0, 'volume': 10.0, 'low': 594.0,
'close': 595.0, 'open': 594.0,
'last_traded': Timestamp('2018-03-01 09:45:00+0000',
tz='UTC')
},
{'high': 595.0, 'volume': 0.0, 'low': 595.0,
'close': 595.0, 'open': 595.0,
'last_traded': Timestamp('2018-03-01 09:50:00+0000',
tz='UTC')
},
{'high': 594.0, 'volume': 108.0, 'low': 592.0,
'close': 593.0, 'open': 592.0,
'last_traded': Timestamp('2018-03-01 09:55:00+0000',
tz='UTC')
}]
expected_df = transform_candles_to_df(expected)
self.verify_forward_fill_df_if_needed(candles, periods, expected_df)
self.verify_get_candles_df(assets, candles, periods[2], expected_df)
# test "forward fill" at the beginning
candles = [{'high': 595, 'volume': 10, 'low': 594,
'close': 595, 'open': 594,
'last_traded': Timestamp('2018-03-01 09:50:00+0000',
tz='UTC')
},
{'high': 594, 'volume': 108, 'low': 592,
'close': 593, 'open': 592,
'last_traded': Timestamp('2018-03-01 09:55:00+0000',
tz='UTC')
}]
expected = [{'high': np.NaN, 'volume': 0.0, 'low': np.NaN,
'close': np.NaN, 'open': np.NaN,
'last_traded': Timestamp('2018-03-01 09:45:00+0000',
tz='UTC')
},
{'high': 595, 'volume': 10, 'low': 594,
'close': 595, 'open': 594,
'last_traded': Timestamp('2018-03-01 09:50:00+0000',
tz='UTC')
},
{'high': 594, 'volume': 108, 'low': 592,
'close': 593, 'open': 592,
'last_traded': Timestamp('2018-03-01 09:55:00+0000',
tz='UTC')
}]
expected_df = transform_candles_to_df(expected)
self.verify_forward_fill_df_if_needed(candles, periods, expected_df)
# Not the same due to dropna - commenting out for now
# self.verify_get_candles_df(assets, candles, periods[2], expected_df)
+16 -12
View File
@@ -107,14 +107,14 @@ class TestSuiteBundle:
print('saved {} test results: {}'.format(end_dt, folder)) print('saved {} test results: {}'.format(end_dt, folder))
assert_frame_equal( assert_frame_equal(
right=data['bundle'], right=data['bundle'][:-1],
left=data['exchange'], left=data['exchange'][:-1],
check_less_precise=1, check_less_precise=1,
) )
try: try:
assert_frame_equal( assert_frame_equal(
right=data['bundle'], right=data['bundle'][:-1],
left=data['exchange'], left=data['exchange'][:-1],
check_less_precise=min([a.decimals for a in assets]), check_less_precise=min([a.decimals for a in assets]),
) )
except Exception as e: except Exception as e:
@@ -197,24 +197,28 @@ class TestSuiteBundle:
# population=exchange_population, # population=exchange_population,
# features=[bundle], # features=[bundle],
# ) # Type: list[Exchange] # ) # Type: list[Exchange]
exchanges = [get_exchange('binance', skip_init=True)] # TODO: currently focusing on Binance, try other exchanges
exchanges = [get_exchange('poloniex', skip_init=True)]
data_portal = TestSuiteBundle.get_data_portal(exchanges) data_portal = TestSuiteBundle.get_data_portal(exchanges)
for exchange in exchanges: for exchange in exchanges:
exchange.init() exchange.init()
frequencies = exchange.get_candle_frequencies(data_frequency) frequencies = exchange.get_candle_frequencies(data_frequency)
freq = random.sample(frequencies, 1)[0] # freq = random.sample(frequencies, 1)[0]
freq = '5T'
rnd = random.SystemRandom() rnd = random.SystemRandom()
# field = rnd.choice(['open', 'high', 'low', 'close', 'volume']) # field = rnd.choice(['open', 'high', 'low', 'close', 'volume'])
field = rnd.choice(['volume']) field = rnd.choice(['close'])
bar_count = random.randint(3, 6) # bar_count = random.randint(3, 6)
bar_count = 5
assets = select_random_assets( # assets = select_random_assets(
exchange.assets, asset_population # exchange.assets, asset_population
) # )
end_dt = None assets = [exchange.get_asset('bch_eth')]
end_dt = pd.to_datetime('2018-03-01', utc=True)
for asset in assets: for asset in assets:
attribute = 'end_{}'.format(data_frequency) attribute = 'end_{}'.format(data_frequency)
asset_end_dt = getattr(asset, attribute) asset_end_dt = getattr(asset, attribute)
+2 -3
View File
@@ -1,6 +1,5 @@
from catalyst.marketplace.marketplace import Marketplace from catalyst.marketplace.marketplace import Marketplace
from catalyst.testing.fixtures import WithLogger, ZiplineTestCase from catalyst.testing.fixtures import WithLogger, ZiplineTestCase
import pandas as pd
class TestMarketplace(WithLogger, ZiplineTestCase): class TestMarketplace(WithLogger, ZiplineTestCase):
@@ -16,12 +15,12 @@ class TestMarketplace(WithLogger, ZiplineTestCase):
def test_subscribe(self): def test_subscribe(self):
marketplace = Marketplace() marketplace = Marketplace()
marketplace.subscribe('marketcap2222') marketplace.subscribe('marketcap')
pass pass
def test_ingest(self): def test_ingest(self):
marketplace = Marketplace() marketplace = Marketplace()
ds_def = marketplace.ingest('github') ds_def = marketplace.ingest('marketcap')
pass pass
def test_publish(self): def test_publish(self):