mirror of
https://github.com/wassname/catalyst.git
synced 2026-08-11 11:16:15 +08:00
ENH: Add single-column input/output capabilities to pipeline terms
This commit is contained in:
+134
-130
@@ -72,9 +72,14 @@ from zipline.pipeline.loaders.synthetic import (
|
||||
make_bar_data,
|
||||
expected_bar_values_2d,
|
||||
)
|
||||
from zipline.pipeline.term import NotSpecified
|
||||
from zipline.pipeline.sentinels import NotSpecified
|
||||
from zipline.testing import (
|
||||
AssetID,
|
||||
AssetIDPlusDay,
|
||||
check_arrays,
|
||||
make_alternating_boolean_array,
|
||||
make_cascading_boolean_array,
|
||||
OpenPrice,
|
||||
parameter_space,
|
||||
product_upper_triangle,
|
||||
)
|
||||
@@ -95,38 +100,6 @@ class RollingSumDifference(CustomFactor):
|
||||
out[:] = (open - close).sum(axis=0)
|
||||
|
||||
|
||||
class AssetID(CustomFactor):
|
||||
"""
|
||||
CustomFactor that returns the AssetID of each asset.
|
||||
|
||||
Useful for providing a Factor that produces a different value for each
|
||||
asset.
|
||||
"""
|
||||
window_length = 1
|
||||
# HACK: We currently decide whether to load or compute a Term based on the
|
||||
# length of its inputs. This means we have to provide a dummy input.
|
||||
inputs = [USEquityPricing.close]
|
||||
|
||||
def compute(self, today, assets, out, close):
|
||||
out[:] = assets
|
||||
|
||||
|
||||
class AssetIDPlusDay(CustomFactor):
|
||||
window_length = 1
|
||||
inputs = [USEquityPricing.close]
|
||||
|
||||
def compute(self, today, assets, out, close):
|
||||
out[:] = assets + today.day
|
||||
|
||||
|
||||
class OpenPrice(CustomFactor):
|
||||
window_length = 1
|
||||
inputs = [USEquityPricing.open]
|
||||
|
||||
def compute(self, today, assets, out, open):
|
||||
out[:] = open
|
||||
|
||||
|
||||
class MultipleOutputs(CustomFactor):
|
||||
window_length = 1
|
||||
inputs = [USEquityPricing.open, USEquityPricing.close]
|
||||
@@ -421,6 +394,8 @@ class ConstantInputTestCase(WithTradingEnvironment, ZiplineTestCase):
|
||||
assets = self.assets
|
||||
asset_ids = self.asset_ids
|
||||
constants = self.constants
|
||||
num_dates = len(dates)
|
||||
num_assets = len(assets)
|
||||
open = USEquityPricing.open
|
||||
close = USEquityPricing.close
|
||||
engine = SimplePipelineEngine(
|
||||
@@ -435,19 +410,13 @@ class ConstantInputTestCase(WithTradingEnvironment, ZiplineTestCase):
|
||||
return DataFrame(expected_values, index=dates, columns=assets)
|
||||
|
||||
cascading_mask = AssetIDPlusDay() < (asset_ids[-1] + dates[0].day)
|
||||
expected_cascading_mask_result = array(
|
||||
[[True, True, True, False],
|
||||
[True, True, False, False],
|
||||
[True, False, False, False]],
|
||||
dtype=bool,
|
||||
expected_cascading_mask_result = make_cascading_boolean_array(
|
||||
shape=(num_dates, num_assets),
|
||||
)
|
||||
|
||||
alternating_mask = (AssetIDPlusDay() % 2).eq(0)
|
||||
expected_alternating_mask_result = array(
|
||||
[[False, True, False, True],
|
||||
[True, False, True, False],
|
||||
[False, True, False, True]],
|
||||
dtype=bool,
|
||||
expected_alternating_mask_result = make_alternating_boolean_array(
|
||||
shape=(num_dates, num_assets), first_value=False,
|
||||
)
|
||||
|
||||
masks = cascading_mask, alternating_mask
|
||||
@@ -592,6 +561,8 @@ class ConstantInputTestCase(WithTradingEnvironment, ZiplineTestCase):
|
||||
assets = self.assets
|
||||
asset_ids = self.asset_ids
|
||||
constants = self.constants
|
||||
num_dates = len(dates)
|
||||
num_assets = len(assets)
|
||||
open = USEquityPricing.open
|
||||
close = USEquityPricing.close
|
||||
engine = SimplePipelineEngine(
|
||||
@@ -603,32 +574,17 @@ class ConstantInputTestCase(WithTradingEnvironment, ZiplineTestCase):
|
||||
return DataFrame(expected_values, index=dates, columns=assets)
|
||||
|
||||
cascading_mask = AssetIDPlusDay() < (asset_ids[-1] + dates[0].day)
|
||||
expected_cascading_mask_result = array(
|
||||
[[True, True, True, False],
|
||||
[True, True, False, False],
|
||||
[True, False, False, False],
|
||||
[False, False, False, False],
|
||||
[False, False, False, False]],
|
||||
dtype=bool,
|
||||
expected_cascading_mask_result = make_cascading_boolean_array(
|
||||
shape=(num_dates, num_assets),
|
||||
)
|
||||
|
||||
alternating_mask = (AssetIDPlusDay() % 2).eq(0)
|
||||
expected_alternating_mask_result = array(
|
||||
[[False, True, False, True],
|
||||
[True, False, True, False],
|
||||
[False, True, False, True],
|
||||
[True, False, True, False],
|
||||
[False, True, False, True]],
|
||||
dtype=bool,
|
||||
expected_alternating_mask_result = make_alternating_boolean_array(
|
||||
shape=(num_dates, num_assets), first_value=False,
|
||||
)
|
||||
|
||||
expected_no_mask_result = array(
|
||||
[[True, True, True, True],
|
||||
[True, True, True, True],
|
||||
[True, True, True, True],
|
||||
[True, True, True, True],
|
||||
[True, True, True, True]],
|
||||
dtype=bool,
|
||||
expected_no_mask_result = full(
|
||||
shape=(num_dates, num_assets), fill_value=True, dtype=bool,
|
||||
)
|
||||
|
||||
masks = cascading_mask, alternating_mask, NotSpecified
|
||||
@@ -1258,19 +1214,39 @@ class ParameterizedFactorTestCase(WithTradingEnvironment, ZiplineTestCase):
|
||||
`RollingSpearmanOfReturns`.
|
||||
"""
|
||||
my_asset_column = 0
|
||||
start_date_index = 6
|
||||
end_date_index = 10
|
||||
start_date_index = 14
|
||||
end_date_index = 18
|
||||
|
||||
assets = self.asset_finder.retrieve_all(self.sids)
|
||||
sids = self.sids
|
||||
dates = self.dates
|
||||
assets = self.asset_finder.retrieve_all(sids)
|
||||
my_asset = assets[my_asset_column]
|
||||
my_asset_filter = (AssetID() != (my_asset_column + 1))
|
||||
num_days = end_date_index - start_date_index + 1
|
||||
num_assets = len(assets)
|
||||
|
||||
# Our correlation factors require that their target asset is not
|
||||
# filtered out, so make sure that masking out our target asset does not
|
||||
# take effect. That is, a filter which filters out only our target
|
||||
# asset should produce the same result as if no mask was passed at all.
|
||||
for mask in (NotSpecified, my_asset_filter):
|
||||
cascading_mask = \
|
||||
AssetIDPlusDay() < (sids[-1] + dates[start_date_index].day)
|
||||
expected_cascading_mask_result = make_cascading_boolean_array(
|
||||
shape=(num_days, num_assets),
|
||||
)
|
||||
|
||||
alternating_mask = (AssetIDPlusDay() % 2).eq(0)
|
||||
expected_alternating_mask_result = make_alternating_boolean_array(
|
||||
shape=(num_days, num_assets),
|
||||
)
|
||||
|
||||
expected_no_mask_result = full(
|
||||
shape=(num_days, num_assets), fill_value=True, dtype=bool,
|
||||
)
|
||||
|
||||
masks = cascading_mask, alternating_mask, NotSpecified
|
||||
expected_mask_results = (
|
||||
expected_cascading_mask_result,
|
||||
expected_alternating_mask_result,
|
||||
expected_no_mask_result,
|
||||
)
|
||||
|
||||
for mask, expected_mask in zip(masks, expected_mask_results):
|
||||
pearson_factor = RollingPearsonOfReturns(
|
||||
target=my_asset,
|
||||
returns_length=returns_length,
|
||||
@@ -1284,18 +1260,23 @@ class ParameterizedFactorTestCase(WithTradingEnvironment, ZiplineTestCase):
|
||||
mask=mask,
|
||||
)
|
||||
|
||||
pipeline = Pipeline(
|
||||
columns={
|
||||
'pearson_factor': pearson_factor,
|
||||
'spearman_factor': spearman_factor,
|
||||
},
|
||||
)
|
||||
if mask is not NotSpecified:
|
||||
pipeline.add(mask, 'mask')
|
||||
|
||||
results = self.engine.run_pipeline(
|
||||
Pipeline(
|
||||
columns={
|
||||
'pearson_factor': pearson_factor,
|
||||
'spearman_factor': spearman_factor,
|
||||
},
|
||||
),
|
||||
self.dates[start_date_index],
|
||||
self.dates[end_date_index],
|
||||
pipeline, dates[start_date_index], dates[end_date_index],
|
||||
)
|
||||
pearson_results = results['pearson_factor'].unstack()
|
||||
spearman_results = results['spearman_factor'].unstack()
|
||||
if mask is not NotSpecified:
|
||||
mask_results = results['mask'].unstack()
|
||||
check_arrays(mask_results.values, expected_mask)
|
||||
|
||||
# Run a separate pipeline that calculates returns starting
|
||||
# (correlation_length - 1) days prior to our start date. This is
|
||||
@@ -1304,8 +1285,8 @@ class ParameterizedFactorTestCase(WithTradingEnvironment, ZiplineTestCase):
|
||||
returns = Returns(window_length=returns_length)
|
||||
results = self.engine.run_pipeline(
|
||||
Pipeline(columns={'returns': returns}),
|
||||
self.dates[start_date_index - (correlation_length - 1)],
|
||||
self.dates[end_date_index],
|
||||
dates[start_date_index - (correlation_length - 1)],
|
||||
dates[end_date_index],
|
||||
)
|
||||
returns_results = results['returns'].unstack()
|
||||
|
||||
@@ -1328,22 +1309,19 @@ class ParameterizedFactorTestCase(WithTradingEnvironment, ZiplineTestCase):
|
||||
my_asset_returns, other_asset_returns,
|
||||
)[0]
|
||||
|
||||
assert_frame_equal(
|
||||
pearson_results,
|
||||
DataFrame(
|
||||
expected_pearson_results,
|
||||
index=self.dates[start_date_index:end_date_index + 1],
|
||||
columns=assets,
|
||||
),
|
||||
expected_pearson_results = DataFrame(
|
||||
data=where(expected_mask, expected_pearson_results, nan),
|
||||
index=dates[start_date_index:end_date_index + 1],
|
||||
columns=assets,
|
||||
)
|
||||
assert_frame_equal(
|
||||
spearman_results,
|
||||
DataFrame(
|
||||
expected_spearman_results,
|
||||
index=self.dates[start_date_index:end_date_index + 1],
|
||||
columns=assets,
|
||||
),
|
||||
assert_frame_equal(pearson_results, expected_pearson_results)
|
||||
|
||||
expected_spearman_results = DataFrame(
|
||||
data=where(expected_mask, expected_spearman_results, nan),
|
||||
index=dates[start_date_index:end_date_index + 1],
|
||||
columns=assets,
|
||||
)
|
||||
assert_frame_equal(spearman_results, expected_spearman_results)
|
||||
|
||||
@parameter_space(returns_length=[2, 3], regression_length=[3, 4])
|
||||
def test_regression_of_returns_factor(self,
|
||||
@@ -1353,38 +1331,65 @@ class ParameterizedFactorTestCase(WithTradingEnvironment, ZiplineTestCase):
|
||||
Tests for the built-in factor `RollingLinearRegressionOfReturns`.
|
||||
"""
|
||||
my_asset_column = 0
|
||||
start_date_index = 6
|
||||
end_date_index = 10
|
||||
start_date_index = 14
|
||||
end_date_index = 18
|
||||
|
||||
assets = self.asset_finder.retrieve_all(self.sids)
|
||||
sids = self.sids
|
||||
dates = self.dates
|
||||
assets = self.asset_finder.retrieve_all(sids)
|
||||
my_asset = assets[my_asset_column]
|
||||
my_asset_filter = (AssetID() != (my_asset_column + 1))
|
||||
num_days = end_date_index - start_date_index + 1
|
||||
num_assets = len(assets)
|
||||
|
||||
cascading_mask = \
|
||||
AssetIDPlusDay() < (sids[-1] + dates[start_date_index].day)
|
||||
expected_cascading_mask_result = make_cascading_boolean_array(
|
||||
shape=(num_days, num_assets),
|
||||
)
|
||||
|
||||
alternating_mask = (AssetIDPlusDay() % 2).eq(0)
|
||||
expected_alternating_mask_result = make_alternating_boolean_array(
|
||||
shape=(num_days, num_assets),
|
||||
)
|
||||
|
||||
expected_no_mask_result = full(
|
||||
shape=(num_days, num_assets), fill_value=True, dtype=bool,
|
||||
)
|
||||
|
||||
masks = cascading_mask, alternating_mask, NotSpecified
|
||||
expected_mask_results = (
|
||||
expected_cascading_mask_result,
|
||||
expected_alternating_mask_result,
|
||||
expected_no_mask_result,
|
||||
)
|
||||
|
||||
# The order of these is meant to align with the output of `linregress`.
|
||||
outputs = ['beta', 'alpha', 'r_value', 'p_value', 'stderr']
|
||||
|
||||
# Our regression factor requires that its target asset is not filtered
|
||||
# out, so make sure that masking out our target asset does not take
|
||||
# effect. That is, a filter which filters out only our target asset
|
||||
# should produce the same result as if no mask was passed at all.
|
||||
for mask in (NotSpecified, my_asset_filter):
|
||||
for mask, expected_mask in zip(masks, expected_mask_results):
|
||||
regression_factor = RollingLinearRegressionOfReturns(
|
||||
target=my_asset,
|
||||
returns_length=returns_length,
|
||||
regression_length=regression_length,
|
||||
mask=mask,
|
||||
)
|
||||
results = self.engine.run_pipeline(
|
||||
Pipeline(
|
||||
columns={
|
||||
output: getattr(regression_factor, output)
|
||||
for output in outputs
|
||||
},
|
||||
),
|
||||
self.dates[start_date_index],
|
||||
self.dates[end_date_index],
|
||||
|
||||
pipeline = Pipeline(
|
||||
columns={
|
||||
output: getattr(regression_factor, output)
|
||||
for output in outputs
|
||||
},
|
||||
)
|
||||
if mask is not NotSpecified:
|
||||
pipeline.add(mask, 'mask')
|
||||
|
||||
results = self.engine.run_pipeline(
|
||||
pipeline, dates[start_date_index], dates[end_date_index],
|
||||
)
|
||||
if mask is not NotSpecified:
|
||||
mask_results = results['mask'].unstack()
|
||||
check_arrays(mask_results.values, expected_mask)
|
||||
|
||||
output_results = {}
|
||||
expected_output_results = {}
|
||||
for output in outputs:
|
||||
@@ -1393,15 +1398,15 @@ class ParameterizedFactorTestCase(WithTradingEnvironment, ZiplineTestCase):
|
||||
output_results[output], nan,
|
||||
)
|
||||
|
||||
# Run a separate pipeline that calculates returns starting 2 days
|
||||
# prior to our start date. This is because we need
|
||||
# (regression_length - 1) extra days of returns to compute our
|
||||
# expected regressions.
|
||||
# Run a separate pipeline that calculates returns starting
|
||||
# (regression_length - 1) days prior to our start date. This is
|
||||
# because we need (regression_length - 1) extra days of returns to
|
||||
# compute our expected regressions.
|
||||
returns = Returns(window_length=returns_length)
|
||||
results = self.engine.run_pipeline(
|
||||
Pipeline(columns={'returns': returns}),
|
||||
self.dates[start_date_index - (regression_length - 1)],
|
||||
self.dates[end_date_index],
|
||||
dates[start_date_index - (regression_length - 1)],
|
||||
dates[end_date_index],
|
||||
)
|
||||
returns_results = results['returns'].unstack()
|
||||
|
||||
@@ -1424,14 +1429,13 @@ class ParameterizedFactorTestCase(WithTradingEnvironment, ZiplineTestCase):
|
||||
expected_regression_results[i]
|
||||
|
||||
for output in outputs:
|
||||
assert_frame_equal(
|
||||
output_results[output],
|
||||
DataFrame(
|
||||
expected_output_results[output],
|
||||
index=self.dates[start_date_index:end_date_index + 1],
|
||||
columns=assets,
|
||||
),
|
||||
output_result = output_results[output]
|
||||
expected_output_result = DataFrame(
|
||||
where(expected_mask, expected_output_results[output], nan),
|
||||
index=dates[start_date_index:end_date_index + 1],
|
||||
columns=assets,
|
||||
)
|
||||
assert_frame_equal(output_result, expected_output_result)
|
||||
|
||||
def test_correlation_and_regression_with_bad_asset(self):
|
||||
"""
|
||||
@@ -1439,8 +1443,8 @@ class ParameterizedFactorTestCase(WithTradingEnvironment, ZiplineTestCase):
|
||||
`RollingLinearRegressionOfReturns` raise the proper exception when
|
||||
given a nonexistent target asset.
|
||||
"""
|
||||
start_date_index = 6
|
||||
end_date_index = 10
|
||||
start_date_index = 14
|
||||
end_date_index = 18
|
||||
my_asset = Equity(0)
|
||||
|
||||
# This filter is arbitrary; the important thing is that we test each
|
||||
|
||||
Reference in New Issue
Block a user