mirror of
https://github.com/wassname/catalyst.git
synced 2026-07-21 12:30:16 +08:00
187 lines
6.2 KiB
Python
187 lines
6.2 KiB
Python
import json, time, csv
|
|
from datetime import datetime
|
|
import pandas as pd
|
|
import os
|
|
import time
|
|
import requests
|
|
import logbook
|
|
|
|
import catalyst.data.bundles.core as bundles
|
|
|
|
DT_START = time.mktime(datetime(2010, 01, 01, 0, 0).timetuple())
|
|
# DT_START = time.mktime(datetime(2017, 06, 13, 0, 0).timetuple()) # TODO: remove temp
|
|
CSV_OUT_FOLDER = '/var/tmp/catalyst/data/poloniex/'
|
|
CONN_RETRIES = 2
|
|
|
|
logbook.StderrHandler().push_application()
|
|
log = logbook.Logger(__name__)
|
|
|
|
class PoloniexDataGenerator(object):
|
|
"""
|
|
OHLCV data feed generator for crypto data. Based on Poloniex market data
|
|
"""
|
|
|
|
_api_path = 'https://poloniex.com/public?'
|
|
currency_pairs = []
|
|
|
|
def __init__(self):
|
|
if not os.path.exists(CSV_OUT_FOLDER):
|
|
try:
|
|
os.makedirs(CSV_OUT_FOLDER)
|
|
except Exception as e:
|
|
log.error('Failed to create data folder: %s' % CSV_OUT_FOLDER)
|
|
log.exception(e)
|
|
|
|
def get_currency_pairs(self):
|
|
url = self._api_path + 'command=returnTicker'
|
|
|
|
try:
|
|
response = requests.get(url)
|
|
except Exception as e:
|
|
log.error('Failed to retrieve list of currency pairs')
|
|
log.exception(e)
|
|
return None
|
|
|
|
data = response.json()
|
|
self.currency_pairs = []
|
|
for ticker in data:
|
|
self.currency_pairs.append(ticker)
|
|
self.currency_pairs.sort()
|
|
|
|
log.debug('Currency pairs retrieved successfully: %d' % (len(self.currency_pairs)))
|
|
|
|
def _get_start_date(self, csv_fn):
|
|
''' Function returns latest appended date, if the file has been previously written
|
|
the last line is an empty one, so we have to read the second to last line
|
|
'''
|
|
try:
|
|
with open(csv_fn, 'ab+') as f:
|
|
f.seek(0, os.SEEK_END) # First check file is not zero size
|
|
if(f.tell() > 2):
|
|
f.seek(-2, os.SEEK_END) # Jump to the second last byte.
|
|
while f.read(1) != b"\n": # Until EOL is found...
|
|
f.seek(-2, os.SEEK_CUR) # ...jump back the read byte plus one more.
|
|
lastrow = f.readline()
|
|
return int(lastrow.split(',')[0]) + 300
|
|
|
|
except Exception as e:
|
|
log.error('Error opening file: %s' % csv_fn)
|
|
log.exception(e)
|
|
|
|
return DT_START
|
|
|
|
def get_data(self, currencyPair, start, end=9999999999, period=300):
|
|
url = self._api_path + 'command=returnChartData¤cyPair=' + currencyPair + '&start=' + str(start) + '&end=' + str(end) + '&period=' + str(period)
|
|
|
|
try:
|
|
response = requests.get(url)
|
|
except Exception as e:
|
|
log.error('Failed to retrieve candlestick chart data for %s' % currencyPair)
|
|
log.exception(e)
|
|
return None
|
|
|
|
return response.json()
|
|
|
|
'''
|
|
Pulls latest data for a single pair
|
|
'''
|
|
def append_data_single_pair(self, currencyPair, repeat=0):
|
|
log.debug('Getting data for %s' % currencyPair)
|
|
csv_fn = CSV_OUT_FOLDER + 'crypto_prices-' + currencyPair + '.csv'
|
|
start = self._get_start_date(csv_fn)
|
|
if (time.time() > start): # Only fetch data if more than 5min have passed since last fetch
|
|
data = self.get_data(currencyPair, start)
|
|
if data is not None:
|
|
try:
|
|
with open(csv_fn, 'ab') as csvfile:
|
|
csvwriter = csv.writer(csvfile)
|
|
for item in data:
|
|
if item['date'] == 0:
|
|
continue
|
|
csvwriter.writerow([item['date'], item['open'], item['high'], item['low'], item['close'], item['volume']])
|
|
except Exception as e:
|
|
log.error('Error opening %s' % csv_fn)
|
|
log.exception(e)
|
|
elif (repeat < CONN_RETRIES):
|
|
log.debug('Retrying: attemt %d' % (repeat+1) )
|
|
self.append_data_single_pair(currencyPair, repeat + 1)
|
|
|
|
'''
|
|
Pulls latest data for all currency pairs
|
|
'''
|
|
def append_data(self):
|
|
for currencyPair in self.currency_pairs:
|
|
self.append_data_single_pair(currencyPair)
|
|
time.sleep(0.17) # Rate limit is 6 calls per second, sleep 1sec/6 to be safe
|
|
|
|
'''
|
|
Returns a data frame for all pairs, or for the requests currency pair.
|
|
Makes sure data is up to date
|
|
'''
|
|
def to_dataframe(self, start, end, currencyPair=None):
|
|
csv_fn = CSV_OUT_FOLDER + 'crypto_prices-' + currencyPair + '.csv'
|
|
last_date = self._get_start_date(csv_fn)
|
|
if last_date + 300 < end or not os.path.exists(csv_fn):
|
|
# get latest data
|
|
self.append_data_single_pair(currencyPair)
|
|
|
|
# CSV holds the latest snapshot
|
|
df = pd.read_csv(csv_fn, names=['date', 'open', 'high', 'low', 'close', 'volume'])
|
|
df['date']=pd.to_datetime(df['date'],unit='s')
|
|
df.set_index('date', inplace=True)
|
|
|
|
#return df.loc[(df.index > start) & (df.index <= end)]
|
|
return df[datetime.fromtimestamp(start):datetime.fromtimestamp(end-1)]
|
|
|
|
if __name__ == '__main__':
|
|
pdg = PoloniexDataGenerator()
|
|
pdg.get_currency_pairs()
|
|
pdg.append_data()
|
|
|
|
|
|
|
|
# from zipline.utils.calendars import get_calendar
|
|
# from zipline.data.us_equity_pricing import (
|
|
# BcolzDailyBarWriter,
|
|
# BcolzDailyBarReader,
|
|
# )
|
|
|
|
# open_calendar = get_calendar('OPEN')
|
|
|
|
# start_session = pd.Timestamp('2012-12-31', tz='UTC')
|
|
# end_session = pd.Timestamp('2015-01-01', tz='UTC')
|
|
|
|
# file_path = 'test.bcolz'
|
|
|
|
# writer = BcolzDailyBarWriter(
|
|
# file_path,
|
|
# open_calendar,
|
|
# start_session,
|
|
# end_session
|
|
# )
|
|
|
|
# index = open_calendar.schedule.index
|
|
# index = index[
|
|
# (index.date >= start_session.date()) &
|
|
# (index.date <= end_session.date())
|
|
# ]
|
|
|
|
# data = pd.DataFrame(
|
|
# 0,
|
|
# index=index,
|
|
# columns=['open', 'high', 'low', 'close', 'volume'],
|
|
# )
|
|
|
|
# writer.write(
|
|
# [(0, data)],
|
|
# assets=[0],
|
|
# show_progress=True
|
|
# )
|
|
|
|
# print 'len(index):', len(index)
|
|
|
|
# reader = BcolzDailyBarReader(file_path)
|
|
|
|
# print 'first_rows:', reader._first_rows
|
|
# print 'last_rows:', reader._last_rows
|