mirror of
https://github.com/wassname/catalyst.git
synced 2026-08-04 12:45:06 +08:00
Five minute bar handle_data and pipline data
This commit is contained in:
+15
-10
@@ -746,14 +746,15 @@ class TradingAlgorithm(object):
|
||||
for perf in self.get_generator():
|
||||
perfs.append(perf)
|
||||
|
||||
|
||||
# convert perf dict to pandas dataframe
|
||||
daily_stats = self._create_daily_stats(perfs)
|
||||
stats = self._create_daily_stats(perfs)
|
||||
|
||||
self.analyze(daily_stats)
|
||||
self.analyze(stats)
|
||||
finally:
|
||||
self.data_portal = None
|
||||
|
||||
return daily_stats
|
||||
return stats
|
||||
|
||||
def _write_and_map_id_index_to_sids(self, identifiers, as_of_date):
|
||||
# Build new Assets for identifiers that can't be resolved as
|
||||
@@ -1144,14 +1145,12 @@ class TradingAlgorithm(object):
|
||||
|
||||
date_rule = date_rule or date_rules.every_day()
|
||||
if freq is 'daily':
|
||||
# ignore time rule in daily mode
|
||||
# Ignore any time rules in daily mode.
|
||||
# every_minute in daily mode does nothing.
|
||||
time_rule = time_rules.every_minute()
|
||||
else:
|
||||
# use provided time rule or default to every minute or 5 minutes
|
||||
# based on desired data frequency.
|
||||
time_rule = time_rule or (time_rules.every_5_minutes()
|
||||
if freq is '5-minute' else
|
||||
time_rules.every_minute())
|
||||
# use provided time rule or default to every minute
|
||||
time_rule = time_rule or time_rules.every_minute()
|
||||
|
||||
# Check the type of the algorithm's schedule before pulling calendar
|
||||
# Note that the ExchangeTradingSchedule is currently the only
|
||||
@@ -1175,7 +1174,13 @@ class TradingAlgorithm(object):
|
||||
)
|
||||
|
||||
self.add_event(
|
||||
make_eventrule(date_rule, time_rule, cal, half_days),
|
||||
make_eventrule(
|
||||
date_rule,
|
||||
time_rule,
|
||||
cal,
|
||||
half_days=half_days,
|
||||
data_frequency=self.data_frequency,
|
||||
),
|
||||
func,
|
||||
)
|
||||
|
||||
|
||||
@@ -60,7 +60,7 @@ OPEN_FIVE_MINUTES_PER_DAY = 288
|
||||
|
||||
DEFAULT_EXPECTEDLEN_CRYPTO = OPEN_FIVE_MINUTES_PER_DAY * 366 * 15
|
||||
|
||||
OHLC_RATIO = 1000000
|
||||
OHLC_RATIO = 1000
|
||||
|
||||
OHLC = frozenset(['open', 'high', 'low', 'close'])
|
||||
OHLCV = frozenset(['open', 'high', 'low', 'close', 'volume'])
|
||||
|
||||
@@ -189,14 +189,14 @@ class PerformanceTracker(object):
|
||||
|
||||
@property
|
||||
def progress(self):
|
||||
if self.emission_rate == 'minute':
|
||||
if self.emission_rate in set(('minute', '5-minute')):
|
||||
# Fake a value
|
||||
return 1.0
|
||||
elif self.emission_rate == 'daily':
|
||||
return self.session_count / self.total_session_count
|
||||
|
||||
def set_date(self, date):
|
||||
if self.emission_rate == 'minute':
|
||||
if self.emission_rate in set(('minute', '5-minute')):
|
||||
self.saved_dt = date
|
||||
self.todays_performance.period_close = self.saved_dt
|
||||
|
||||
@@ -370,7 +370,9 @@ class PerformanceTracker(object):
|
||||
bench_since_open,
|
||||
account.leverage)
|
||||
|
||||
assert self.emission_rate in set(('minute', '5-minute'))
|
||||
minute_packet = self.to_dict(emission_type='minute')
|
||||
|
||||
return minute_packet
|
||||
|
||||
def handle_market_close(self, dt, data_portal):
|
||||
|
||||
@@ -47,6 +47,8 @@ __all__ = [
|
||||
'NDaysBeforeLastTradingDayOfMonth',
|
||||
'StatefulRule',
|
||||
'OncePerDay',
|
||||
'OncePerFiveMinutes',
|
||||
'OncePerMinute',
|
||||
|
||||
# Factory API
|
||||
'date_rules',
|
||||
@@ -552,15 +554,18 @@ class StatefulRule(EventRule):
|
||||
"""
|
||||
self.should_trigger = callable_
|
||||
|
||||
|
||||
class OncePerDay(StatefulRule):
|
||||
class OncePerInterval(StatefulRule):
|
||||
def __init__(self, rule=None):
|
||||
self.triggered = False
|
||||
|
||||
self.date = None
|
||||
self.next_date = None
|
||||
|
||||
super(OncePerDay, self).__init__(rule)
|
||||
super(OncePerInterval, self).__init__(rule)
|
||||
|
||||
@lazyval
|
||||
def interval(self):
|
||||
raise NotImplementedError
|
||||
|
||||
def should_trigger(self, dt):
|
||||
if self.date is None or dt >= self.next_date:
|
||||
@@ -570,11 +575,28 @@ class OncePerDay(StatefulRule):
|
||||
|
||||
# record the timestamp for the next day, so that we can use it
|
||||
# to know if we've moved to the next day
|
||||
self.next_date = dt + pd.Timedelta(1, unit="d")
|
||||
self.next_date = dt + self.interval
|
||||
|
||||
if not self.triggered and self.rule.should_trigger(dt):
|
||||
self.triggered = True
|
||||
return True
|
||||
|
||||
|
||||
|
||||
class OncePerDay(OncePerInterval):
|
||||
@lazyval
|
||||
def interval(self):
|
||||
return pd.Timedelta(1, unit='d')
|
||||
|
||||
class OncePerFiveMinutes(OncePerInterval):
|
||||
@lazyval
|
||||
def interval(self):
|
||||
return pd.Timedelta(5, unit='m')
|
||||
|
||||
class OncePerMinute(OncePerInterval):
|
||||
@lazyval
|
||||
def interval(self):
|
||||
return pd.Timedelta(1, unit='m')
|
||||
|
||||
|
||||
# Factory API
|
||||
@@ -612,7 +634,11 @@ class calendars(object):
|
||||
US_FUTURES = sentinel('US_FUTURES')
|
||||
|
||||
|
||||
def make_eventrule(date_rule, time_rule, cal, half_days=True):
|
||||
def make_eventrule(date_rule,
|
||||
time_rule,
|
||||
cal,
|
||||
half_days=True,
|
||||
data_frequency=None):
|
||||
"""
|
||||
Constructs an event rule from the factory api.
|
||||
"""
|
||||
@@ -628,4 +654,15 @@ def make_eventrule(date_rule, time_rule, cal, half_days=True):
|
||||
nhd_rule.cal = cal
|
||||
inner_rule = date_rule & time_rule & nhd_rule
|
||||
|
||||
return OncePerDay(rule=inner_rule)
|
||||
if data_frequency == 'daily':
|
||||
return OncePerDay(rule=inner_rule)
|
||||
elif data_frequency == '5-minute':
|
||||
return OncePerFiveMinutes(rule=inner_rule)
|
||||
elif data_frequency == 'minute':
|
||||
return OncePerMinute(rule=inner_rule)
|
||||
else:
|
||||
raise ValueError(
|
||||
'Cannot make event rule for data frequency: {}'.format(
|
||||
data_frequency,
|
||||
)
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user