From b357a0656a0072dcb7d6081c38f2f21e9f4adb12 Mon Sep 17 00:00:00 2001 From: fredfortier Date: Wed, 18 Oct 2017 13:58:37 -0400 Subject: [PATCH] Added a modified bcolz writer / reader --- catalyst/exchange/exchange_bcolz.py | 65 +++++++++++++++++++++++++++++ tests/exchange/test_bundle.py | 40 ++++-------------- 2 files changed, 72 insertions(+), 33 deletions(-) create mode 100644 catalyst/exchange/exchange_bcolz.py diff --git a/catalyst/exchange/exchange_bcolz.py b/catalyst/exchange/exchange_bcolz.py new file mode 100644 index 00000000..e24a2d41 --- /dev/null +++ b/catalyst/exchange/exchange_bcolz.py @@ -0,0 +1,65 @@ +import numpy as np + +from catalyst import get_calendar +from catalyst.data.minute_bars import BcolzMinuteBarReader, \ + BcolzMinuteBarWriter + + +class BcolzExchangeBarWriter(BcolzMinuteBarWriter): + def __init__(self, *args, **kwargs): + self._data_frequency = kwargs.pop('data_frequency', None) + kwargs.pop('minutes_per_day', None) + kwargs.pop('default_ohlc_ratio', None) + kwargs.pop('calendar', None) + + minutes_per_day = 1440 if self._data_frequency == 'minute' else 1 + default_ohlc_ratio = 1000000 + calendar = get_calendar('OPEN') + + super(BcolzExchangeBarWriter, self) \ + .__init__(*args, **dict(kwargs, + minutes_per_day=minutes_per_day, + default_ohlc_ratio=default_ohlc_ratio, + calendar=calendar + )) + + +class BcolzExchangeBarReader(BcolzMinuteBarReader): + def __init__(self, *args, **kwargs): + self._data_frequency = kwargs.pop('data_frequency', None) + + super(BcolzExchangeBarReader, self).__init__(*args, **kwargs) + + def load_raw_arrays(self, fields, start_dt, end_dt, sids): + + if self._data_frequency == 'minute': + return super(BcolzExchangeBarReader, self) \ + .load_raw_arrays(fields, start_dt, end_dt, sids) + + else: + return self._load_daily_raw_arrays(fields, start_dt, end_dt, sids) + + def _load_daily_raw_arrays(self, fields, start_dt, end_dt, sids): + start_idx = self._find_position_of_minute(start_dt) + end_idx = self._find_position_of_minute(end_dt) + + num_days = (end_idx - start_idx + 1) + shape = num_days, len(sids) + + data = [] + for field in fields: + out = np.full(shape, np.nan) + + for i, sid in enumerate(sids): + carray = self._open_minute_file(field, sid) + a = carray[start_idx:end_idx + 1] + + where = a != 0 + + out[:len(where), i][where] = ( + a[where] * self._ohlc_ratio_inverse_for_sid(sid) + ) + + data.append(out) + + return data diff --git a/tests/exchange/test_bundle.py b/tests/exchange/test_bundle.py index 99f10c24..6a460695 100644 --- a/tests/exchange/test_bundle.py +++ b/tests/exchange/test_bundle.py @@ -7,6 +7,8 @@ from catalyst import get_calendar from catalyst.data.minute_bars import BcolzMinuteBarReader, \ BcolzMinuteBarWriter from catalyst.exchange.bundle_utils import get_bcolz_chunk, get_periods_range +from catalyst.exchange.exchange_bcolz import BcolzExchangeBarReader, \ + BcolzExchangeBarWriter from catalyst.exchange.exchange_bundle import ExchangeBundle, \ BUNDLE_NAME_TEMPLATE from catalyst.exchange.exchange_utils import get_exchange_folder @@ -182,14 +184,12 @@ class ExchangeBundleTestCase: # I tried setting the minutes_per_day to 1 will not create # unnecessary bars - writer = BcolzMinuteBarWriter( + writer = BcolzExchangeBarWriter( rootdir=path, - calendar=calendar, - minutes_per_day=1, + data_frequency=data_frequency, start_session=start, end_session=end, - write_metadata=True, - default_ohlc_ratio=exchange_bundle.default_ohlc_ratio + write_metadata=True ) # This will read the daily data in a bundle created by @@ -209,34 +209,8 @@ class ExchangeBundleTestCase: empty_rows_behavior='strip' ) - # Simplifying the data reader to play nice with 1 minute per day - class BcolzDayBarReader(BcolzMinuteBarReader): - def load_raw_arrays(self, fields, start_dt, end_dt, sids): - start_idx = self._find_position_of_minute(start_dt) - end_idx = self._find_position_of_minute(end_dt) - - num_days = (end_idx - start_idx + 1) - shape = num_days, len(sids) - - data = [] - for field in fields: - out = np.full(shape, np.nan) - - for i, sid in enumerate(sids): - carray = reader._open_minute_file(field, sid) - a = carray[start_idx:end_idx + 1] - - where = a != 0 - - out[:len(where), i][where] = ( - a[where] * self._ohlc_ratio_inverse_for_sid(sid) - ) - - data.append(out) - - return data - - reader = BcolzDayBarReader(path) + reader = BcolzExchangeBarReader(rootdir=path, + data_frequency=data_frequency) # Reading the two assets to ensure that no data was lost for asset in assets: