mirror of
https://github.com/wassname/catalyst.git
synced 2026-09-10 11:50:32 +08:00
ENH: Send transactions and orders as standalone events.
- Add transaction and order types - Move TransactionSimulator from trading.py to tradesimulation.py (only used by other members of the tradesimulation module) - Make Transaction an independent event, like dividend - Add Blotter class. - Flatten the transaction events to be independent of trade bar events - Make orders into events that reach performance (need to add handling) - Issue IDs to orders and tracking each transaction's order id. - Make volume share slippage fill orders independently, rather than aggregating them into a single transaction. - Perf tracker holds orders, serializes them with transactions. - Order state defined and maintained by order class. - Minutely emission of orders based on last_modified date.
This commit is contained in:
+87
-51
@@ -24,11 +24,13 @@ from operator import attrgetter
|
||||
|
||||
import zipline.utils.factory as factory
|
||||
import zipline.finance.performance as perf
|
||||
from zipline.finance.slippage import Transaction
|
||||
from zipline.finance.slippage import Transaction, create_transaction
|
||||
|
||||
from zipline.gens.composites import date_sorted_sources
|
||||
from zipline.finance.trading import SimulationParameters
|
||||
from zipline.gens.tradesimulation import Order
|
||||
import zipline.finance.trading as trading
|
||||
from zipline.protocol import DATASOURCE_TYPE
|
||||
from zipline.utils.factory import create_random_simulation_parameters
|
||||
|
||||
onesec = datetime.timedelta(seconds=1)
|
||||
@@ -36,6 +38,10 @@ oneday = datetime.timedelta(days=1)
|
||||
tradingday = datetime.timedelta(hours=6, minutes=30)
|
||||
|
||||
|
||||
def create_txn(sid, price, amount, dt):
|
||||
return create_transaction(sid, amount, price, dt, "fakeuid")
|
||||
|
||||
|
||||
class TestDividendPerformance(unittest.TestCase):
|
||||
|
||||
def setUp(self):
|
||||
@@ -78,8 +84,8 @@ class TestDividendPerformance(unittest.TestCase):
|
||||
events[2].dt
|
||||
)
|
||||
|
||||
txn = factory.create_txn(1, 10.0, 100, events[0].dt)
|
||||
events[0].TRANSACTION = txn
|
||||
txn = create_txn(1, 10.0, 100, events[0].dt)
|
||||
events.insert(0, txn)
|
||||
events.insert(1, dividend)
|
||||
perf_tracker = perf.PerformanceTracker(self.sim_params)
|
||||
transformed_events = list(perf_tracker.transform(
|
||||
@@ -136,8 +142,8 @@ class TestDividendPerformance(unittest.TestCase):
|
||||
)
|
||||
|
||||
events.insert(1, dividend)
|
||||
txn = factory.create_txn(1, 10.0, 100, events[3].dt)
|
||||
events[3].TRANSACTION = txn
|
||||
txn = create_txn(1, 10.0, 100, events[3].dt)
|
||||
events.insert(4, txn)
|
||||
perf_tracker = perf.PerformanceTracker(self.sim_params)
|
||||
transformed_events = list(perf_tracker.transform(
|
||||
((event.dt, [event]) for event in events))
|
||||
@@ -183,11 +189,11 @@ class TestDividendPerformance(unittest.TestCase):
|
||||
events[3].dt
|
||||
)
|
||||
|
||||
buy_txn = factory.create_txn(1, 10.0, 100, events[0].dt)
|
||||
events[0].TRANSACTION = buy_txn
|
||||
sell_txn = factory.create_txn(1, 10.0, -100, events[2].dt)
|
||||
events[2].TRANSACTION = sell_txn
|
||||
events.insert(1, dividend)
|
||||
buy_txn = create_txn(1, 10.0, 100, events[0].dt)
|
||||
events.insert(1, buy_txn)
|
||||
sell_txn = create_txn(1, 10.0, -100, events[3].dt)
|
||||
events.insert(4, sell_txn)
|
||||
events.insert(0, dividend)
|
||||
perf_tracker = perf.PerformanceTracker(self.sim_params)
|
||||
transformed_events = list(perf_tracker.transform(
|
||||
((event.dt, [event]) for event in events))
|
||||
@@ -233,10 +239,10 @@ class TestDividendPerformance(unittest.TestCase):
|
||||
events[5].dt
|
||||
)
|
||||
|
||||
buy_txn = factory.create_txn(1, 10.0, 100, events[1].dt)
|
||||
events[1].TRANSACTION = buy_txn
|
||||
sell_txn = factory.create_txn(1, 10.0, -100, events[2].dt)
|
||||
events[2].TRANSACTION = sell_txn
|
||||
buy_txn = create_txn(1, 10.0, 100, events[1].dt)
|
||||
events.insert(2, buy_txn)
|
||||
sell_txn = create_txn(1, 10.0, -100, events[3].dt)
|
||||
events.insert(4, sell_txn)
|
||||
events.insert(1, dividend)
|
||||
perf_tracker = perf.PerformanceTracker(self.sim_params)
|
||||
transformed_events = list(perf_tracker.transform(
|
||||
@@ -280,11 +286,11 @@ class TestDividendPerformance(unittest.TestCase):
|
||||
10.00,
|
||||
events[0].dt,
|
||||
events[1].dt,
|
||||
events[-1].dt + 10*oneday
|
||||
events[-1].dt + 10 * oneday
|
||||
)
|
||||
|
||||
buy_txn = factory.create_txn(1, 10.0, 100, events[1].dt)
|
||||
events[1].TRANSACTION = buy_txn
|
||||
buy_txn = create_txn(1, 10.0, 100, events[1].dt)
|
||||
events.insert(2, buy_txn)
|
||||
events.insert(1, dividend)
|
||||
perf_tracker = perf.PerformanceTracker(self.sim_params)
|
||||
transformed_events = list(perf_tracker.transform(
|
||||
@@ -334,8 +340,8 @@ class TestDividendPerformance(unittest.TestCase):
|
||||
events[2].dt
|
||||
)
|
||||
|
||||
txn = factory.create_txn(1, 10.0, -100, self.dt+oneday)
|
||||
events[0].TRANSACTION = txn
|
||||
txn = create_txn(1, 10.0, -100, self.dt + oneday)
|
||||
events.insert(1, txn)
|
||||
events.insert(0, dividend)
|
||||
perf_tracker = perf.PerformanceTracker(self.sim_params)
|
||||
transformed_events = list(perf_tracker.transform(
|
||||
@@ -431,7 +437,7 @@ class TestPositionPerformance(unittest.TestCase):
|
||||
self.sim_params
|
||||
)
|
||||
|
||||
txn = factory.create_txn(1, 10.0, 100, self.dt + onesec)
|
||||
txn = create_txn(1, 10.0, 100, self.dt + onesec)
|
||||
pp = perf.PerformancePeriod(1000.0)
|
||||
|
||||
pp.execute_transaction(txn)
|
||||
@@ -502,7 +508,7 @@ single short-sale transaction"""
|
||||
|
||||
trades_1 = trades[:-2]
|
||||
|
||||
txn = factory.create_txn(1, 10.0, -100, self.dt + onesec)
|
||||
txn = create_txn(1, 10.0, -100, self.dt + onesec)
|
||||
pp = perf.PerformancePeriod(1000.0)
|
||||
|
||||
pp.execute_transaction(txn)
|
||||
@@ -690,14 +696,14 @@ trade after cover"""
|
||||
self.sim_params
|
||||
)
|
||||
|
||||
short_txn = factory.create_txn(
|
||||
short_txn = create_txn(
|
||||
1,
|
||||
10.0,
|
||||
-100,
|
||||
self.dt + onesec
|
||||
)
|
||||
|
||||
cover_txn = factory.create_txn(1, 7.0, 100, self.dt + onesec * 6)
|
||||
cover_txn = create_txn(1, 7.0, 100, self.dt + onesec * 6)
|
||||
pp = perf.PerformancePeriod(1000.0)
|
||||
|
||||
pp.execute_transaction(short_txn)
|
||||
@@ -805,7 +811,7 @@ shares in position"
|
||||
400
|
||||
)
|
||||
|
||||
saleTxn = factory.create_txn(
|
||||
saleTxn = create_txn(
|
||||
1,
|
||||
10.0,
|
||||
-100,
|
||||
@@ -871,6 +877,11 @@ shares in position"
|
||||
|
||||
class TestPerformanceTracker(unittest.TestCase):
|
||||
|
||||
def setUp(self):
|
||||
|
||||
self.sim_params, self.dt, self.end_dt = \
|
||||
create_random_simulation_parameters()
|
||||
|
||||
NumDaysToDelete = collections.namedtuple(
|
||||
'NumDaysToDelete', ('start', 'middle', 'end'))
|
||||
|
||||
@@ -984,11 +995,15 @@ class TestPerformanceTracker(unittest.TestCase):
|
||||
|
||||
events = date_sorted_sources(trade_history, trade_history2)
|
||||
|
||||
events = [self.event_with_txn(event, trade_history[0].dt)
|
||||
for event in events]
|
||||
events = [event for event in
|
||||
self.trades_with_txns(events, trade_history[0].dt)]
|
||||
|
||||
# Extract events with transactions to use for verification.
|
||||
events_with_txns = [event for event in events if event.TRANSACTION]
|
||||
txns = [event for event in
|
||||
events if event.type == DATASOURCE_TYPE.TRANSACTION]
|
||||
|
||||
orders = [event for event in
|
||||
events if event.type == DATASOURCE_TYPE.ORDER]
|
||||
|
||||
perf_messages = \
|
||||
[msg for date, snapshot in
|
||||
@@ -1002,10 +1017,11 @@ class TestPerformanceTracker(unittest.TestCase):
|
||||
perf_messages.extend(end_perf_messages)
|
||||
|
||||
#we skip two trades, to test case of None transaction
|
||||
self.assertEqual(perf_tracker.txn_count, len(events_with_txns))
|
||||
self.assertEqual(perf_tracker.txn_count, len(txns))
|
||||
self.assertEqual(perf_tracker.txn_count, len(orders))
|
||||
|
||||
cumulative_pos = perf_tracker.cumulative_performance.positions[sid]
|
||||
expected_size = len(events_with_txns) / 2 * -25
|
||||
expected_size = len(txns) / 2 * -25
|
||||
self.assertEqual(cumulative_pos.amount, expected_size)
|
||||
|
||||
self.assertEqual(perf_tracker.last_close,
|
||||
@@ -1014,22 +1030,30 @@ class TestPerformanceTracker(unittest.TestCase):
|
||||
self.assertEqual(len(perf_messages),
|
||||
sim_params.days_in_period)
|
||||
|
||||
def event_with_txn(self, event, no_txn_dt):
|
||||
#create a transaction for all but
|
||||
#first trade in each sid, to simulate None transaction
|
||||
if event.dt != no_txn_dt:
|
||||
txn = Transaction(**{
|
||||
'sid': event.sid,
|
||||
'amount': -25,
|
||||
'dt': event.dt,
|
||||
'price': 10.0,
|
||||
'commission': 0.50
|
||||
})
|
||||
else:
|
||||
txn = None
|
||||
event['TRANSACTION'] = txn
|
||||
def trades_with_txns(self, events, no_txn_dt):
|
||||
for event in events:
|
||||
|
||||
return event
|
||||
#create a transaction for all but
|
||||
#first trade in each sid, to simulate None transaction
|
||||
if event.dt != no_txn_dt:
|
||||
order = Order(**{
|
||||
'sid': event.sid,
|
||||
'amount': -25,
|
||||
'dt': event.dt
|
||||
})
|
||||
yield order
|
||||
yield event
|
||||
txn = Transaction(**{
|
||||
'sid': event.sid,
|
||||
'amount': -25,
|
||||
'dt': event.dt,
|
||||
'price': 10.0,
|
||||
'commission': 0.50,
|
||||
'order_id': order.id
|
||||
})
|
||||
yield txn
|
||||
else:
|
||||
yield event
|
||||
|
||||
@trading.use_environment(trading.TradingEnvironment())
|
||||
def test_minute_tracker(self):
|
||||
@@ -1047,13 +1071,17 @@ class TestPerformanceTracker(unittest.TestCase):
|
||||
tracker = perf.PerformanceTracker(sim_params)
|
||||
|
||||
foo_event_1 = factory.create_trade('foo', 10.0, 20, start_dt)
|
||||
order_event_1 = Order(**{
|
||||
'sid': foo_event_1.sid,
|
||||
'amount': -25,
|
||||
'dt': foo_event_1.dt
|
||||
})
|
||||
bar_event_1 = factory.create_trade('bar', 100.0, 200, start_dt)
|
||||
txn = Transaction(sid=foo_event_1.sid,
|
||||
amount=-25,
|
||||
dt=foo_event_1.dt,
|
||||
price=10.0,
|
||||
commission=0.50)
|
||||
foo_event_1.TRANSACTION = txn
|
||||
txn_event_1 = Transaction(sid=foo_event_1.sid,
|
||||
amount=-25,
|
||||
dt=foo_event_1.dt,
|
||||
price=10.0,
|
||||
commission=0.50)
|
||||
|
||||
foo_event_2 = factory.create_trade(
|
||||
'foo', 11.0, 20, start_dt + datetime.timedelta(minutes=1))
|
||||
@@ -1062,13 +1090,15 @@ class TestPerformanceTracker(unittest.TestCase):
|
||||
|
||||
events = [
|
||||
foo_event_1,
|
||||
order_event_1,
|
||||
txn_event_1,
|
||||
bar_event_1,
|
||||
foo_event_2,
|
||||
bar_event_2
|
||||
]
|
||||
|
||||
import operator
|
||||
messages = {date: snapshot[0].perf_messages[0] for date, snapshot in
|
||||
messages = {date: snapshot[-1].perf_messages[0] for date, snapshot in
|
||||
tracker.transform(
|
||||
itertools.groupby(
|
||||
events,
|
||||
@@ -1085,6 +1115,12 @@ class TestPerformanceTracker(unittest.TestCase):
|
||||
self.assertEquals(0, len(msg_2['intraday_perf']['transactions']),
|
||||
"The second message should have no transactions.")
|
||||
|
||||
self.assertEquals(1, len(msg_1['intraday_perf']['orders']),
|
||||
"The first message should contain one orders.")
|
||||
# Check that orders aren't emitted for previous events.
|
||||
self.assertEquals(0, len(msg_2['intraday_perf']['orders']),
|
||||
"The second message should have no orders.")
|
||||
|
||||
# Ensure that period_close moves through time.
|
||||
# Also, ensure that the period_closes are the expected dts.
|
||||
self.assertEquals(foo_event_1.dt,
|
||||
|
||||
Reference in New Issue
Block a user