From 04d3633545ad2fa430c996a13a4ef78ee9041efc Mon Sep 17 00:00:00 2001 From: Meng Xiangzhuo Date: Wed, 30 Oct 2024 04:03:41 +0800 Subject: [PATCH] feat: implement fetch data from data.binance.vision --- freqtrade/exchange/binance_public_data.py | 179 +++++++++++++++- tests/exchange/test_binance_public_data.py | 192 ++++++++++++++++++ ...utures-um-klines-BTCUSDT-1h-2024-10-28.zip | Bin 0 -> 1533 bytes .../spot-klines-BTCUSDT-1h-2024-10-28.zip | Bin 0 -> 1578 bytes 4 files changed, 370 insertions(+), 1 deletion(-) create mode 100644 tests/exchange/test_binance_public_data.py create mode 100644 tests/testdata/binance/binance_public_data/futures-um-klines-BTCUSDT-1h-2024-10-28.zip create mode 100644 tests/testdata/binance/binance_public_data/spot-klines-BTCUSDT-1h-2024-10-28.zip diff --git a/freqtrade/exchange/binance_public_data.py b/freqtrade/exchange/binance_public_data.py index 23cb02989..adc1ea511 100644 --- a/freqtrade/exchange/binance_public_data.py +++ b/freqtrade/exchange/binance_public_data.py @@ -1,2 +1,179 @@ -async def fetch_ohlcv(*args, **kwargs): +""" +Fetch daily-archived OHLCV data from https://data.binance.vision/ +""" + +import asyncio +import datetime +import io +import itertools +import logging +import zipfile + +import aiohttp +import pandas as pd +from pandas import DataFrame + +from freqtrade.enums import CandleType +from freqtrade.util.datetime_helpers import dt_from_ts, dt_now + + +logger = logging.getLogger(__name__) + + +class BadHttpStatus(Exception): + """Not 200/404""" + pass + + +async def fetch_ohlcv( + candle_type: CandleType, pair: str, timeframe: str, since_ms: int, until_ms: int | None +) -> DataFrame: + """ + Fetch OHLCV data from https://data.binance.vision/ + :candle_type: Currently only spot and futures are supported + :param until_ms: `None` indicates the timestamp of the latest available data + :return: None if no data available in the time range + """ + if candle_type == CandleType.SPOT: + asset_type = "spot" + elif candle_type == CandleType.FUTURES: + asset_type = "futures/um" + else: + raise ValueError(f"Unsupported CandleType: {candle_type}") + symbol = symbol_ccxt_to_binance(pair) + start = dt_from_ts(since_ms) + end = dt_from_ts(until_ms) if until_ms else dt_now() + + # We use two days ago as the last available day because the daily archives are daily uploaded + # and have several hours delay + last_available_date = dt_now() - datetime.timedelta(days=2) + end = min(end, last_available_date) + if start >= end: + return DataFrame() + return await _fetch_ohlcv(asset_type, symbol, timeframe, start, end) + + +def symbol_ccxt_to_binance(symbol: str) -> str: + """ + Convert ccxt symbol notation to binance notation + e.g. BTC/USDT -> BTCUSDT, BTC/USDT:USDT -> BTCUSDT + """ + if ":" in symbol: + parts = symbol.split() + if len(parts) != 2: + raise ValueError(f"Cannot recognize symbol: {symbol}") + return parts[0].replace("/", "") + else: + return symbol.replace("/", "") + + +def concat(dfs) -> DataFrame: + if all(df is None for df in dfs): + return DataFrame() + else: + return pd.concat(dfs) + + +async def _fetch_ohlcv(asset_type, symbol, timeframe, start, end) -> DataFrame: + dfs: list[DataFrame | None] = [] + + connector = aiohttp.TCPConnector(limit=100) + async with aiohttp.ClientSession(connector=connector) as session: + coroutines = [ + get_daily_ohlcv(asset_type, symbol, timeframe, date, session) + for date in date_range(start, end) + ] + # the HTTP connections has been throttled by TCPConnector + for batch in itertools.batched(coroutines, 1000): + results = await asyncio.gather(*batch) + for result in results: + if isinstance(result, BaseException): + logger.warning(f"An exception raised: : {result}") + # Directly return the existing data, do not allow the gap + # between the data + return concat(dfs) + else: + dfs.append(result) + return concat(dfs) + + +def date_range(start: datetime.date, end: datetime.date): + date = start + while date <= end: + yield date + date += datetime.timedelta(days=1) + + +def format_date(date: datetime.date) -> str: + return date.strftime("%Y-%m-%d") + + +def zip_name(symbol: str, timeframe: str, date: datetime.date) -> str: + return f"{symbol}-{timeframe}-{format_date(date)}.zip" + + +async def get_daily_ohlcv( + asset_type: str, + symbol: str, + timeframe: str, + date: datetime.date, + session: aiohttp.ClientSession, + retry_count: int = 3, +) -> DataFrame | None: + """ + Get daily OHLCV from https://data.binance.vision + See https://github.com/binance/binance-public-data + """ + + # example urls: + # https://data.binance.vision/data/spot/daily/klines/BTCUSDT/1s/BTCUSDT-1s-2023-10-27.zip + # https://data.binance.vision/data/futures/um/daily/klines/BTCUSDT/1h/BTCUSDT-1h-2023-10-27.zip + url = ( + f"https://data.binance.vision/data/{asset_type}/daily/klines/{symbol}/{timeframe}/" + f"{zip_name(symbol, timeframe, date)}" + ) + + logger.debug(f"download data from binance: {url}") + + retry = 0 + while True: + if retry > 0: + sleep_secs = retry * 0.5 + logger.debug( + f"[{retry}/{retry_count}] retry to download {url} after {sleep_secs} seconds" + ) + await asyncio.sleep(sleep_secs) + try: + async with session.get(url) as resp: + if resp.status == 200: + content = await resp.read() + logger.debug(f"Successfully downloaded {url}") + with zipfile.ZipFile(io.BytesIO(content)) as zipf: + with zipf.open(zipf.namelist()[0]) as csvf: + # https://github.com/binance/binance-public-data/issues/283 + first_byte = csvf.read(1)[0] + if chr(first_byte).isdigit(): + header = None + else: + header = 0 + csvf.seek(0) + + df = pd.read_csv( + csvf, + usecols=[0, 1, 2, 3, 4, 5], + names=["date", "open", "high", "low", "close", "volume"], + header=header, + ) + df["date"] = pd.to_datetime(df["date"], unit="ms", utc=True) + return df + elif resp.status == 404: + logger.warning(f"No data available for {symbol} in {format_date(date)}") + return None + else: + raise BadHttpStatus(f"{resp.status} - {resp.reason}") + except Exception as e: + retry += 1 + if retry >= retry_count: + logger.warning(f"Failed to get data from {url}: {e}") + raise diff --git a/tests/exchange/test_binance_public_data.py b/tests/exchange/test_binance_public_data.py new file mode 100644 index 000000000..5e0a414f2 --- /dev/null +++ b/tests/exchange/test_binance_public_data.py @@ -0,0 +1,192 @@ +import datetime +import io +import re +import zipfile +from datetime import timedelta + +import aiohttp +import pandas as pd +import pytest + +from freqtrade.enums import CandleType +from freqtrade.exchange.binance_public_data import ( + BadHttpStatus, + fetch_ohlcv, + get_daily_ohlcv, + zip_name, +) +from freqtrade.util.datetime_helpers import dt_ts, dt_utc + + +# spot klines archive csv file format, the futures/um klines don't have the header line +# +# open_time,open,high,low,close,volume,close_time,quote_volume,count,taker_buy_volume,taker_buy_quote_volume,ignore # noqa: E501 +# 1698364800000,34161.6,34182.5,33977.4,34024.2,409953,1698368399999,1202.97118037,15095,192220,564.12041453,0 # noqa: E501 +# 1698368400000,34024.2,34060.1,33776.4,33848.4,740960,1698371999999,2183.75671155,23938,368266,1085.17080793,0 # noqa: E501 +# 1698372000000,33848.5,34150.0,33815.1,34094.2,390376,1698375599999,1147.73267094,13854,231446,680.60405822,0 # noqa: E501 + + +def make_daily_df(date, timeframe): + start = dt_utc(date.year, date.month, date.day) + end = start + timedelta(days=1) + date_col = pd.date_range(start, end, freq=timeframe.replace("m", "min"), inclusive="left") + cols = ( + "open_time,open,high,low,close,volume,close_time,quote_volume,count,taker_buy_volume," + "taker_buy_quote_volume,ignore" + ) + df = pd.DataFrame(columns=cols.split(","), dtype=float) + df["open_time"] = date_col.astype("int64") // 10**6 + df["open"] = df["high"] = df["low"] = df["close"] = df["volume"] = 1.0 + return df + + +def make_daily_zip(asset_type, symbol, timeframe, date) -> bytes: + df = make_daily_df(date, timeframe) + if asset_type == "spot": + header = True + elif asset_type == "futures/um": + header = None + else: + raise ValueError + csv = df.to_csv(index=False, header=header) + zip_buffer = io.BytesIO() + with zipfile.ZipFile(zip_buffer, "w") as zipf: + zipf.writestr(zip_name(symbol, timeframe, date), csv) + return zip_buffer.getvalue() + + +class MockResponse: + def __init__(self, content, status, reason=""): + self._content = content + self.status = status + self.reason = reason + + async def read(self): + return self._content + + async def __aexit__(self, exc_type, exc, tb): + pass + + async def __aenter__(self): + return self + + +def make_response_from_url(start_date, end_date): + def make_response(url): + pattern = ( + r"https://data.binance.vision/data/(?Pspot|futures/um)/daily/klines/" + r"(?P.*?)/(?P.*?)/(?P=symbol)-(?P=timeframe)-" + r"(?P\d{4}-\d{2}-\d{2}).zip" + ) + m = re.match(pattern, url) + if not m: + return MockResponse(content="", status=404) + + date = datetime.datetime.strptime(m["date"], "%Y-%m-%d").date() + if date < start_date or date > end_date: + return MockResponse(content="", status=404) + + zip_file = make_daily_zip(m["asset_type"], m["symbol"], m["timeframe"], date) + return MockResponse(content=zip_file, status=200) + + return make_response + + +@pytest.mark.parametrize( + "since,until,first_date,last_date", + [ + (dt_utc(2020, 1, 1), dt_utc(2020, 1, 2), dt_utc(2020, 1, 1), dt_utc(2020, 1, 2, 23)), + ( + dt_utc(2020, 1, 1), + dt_utc(2020, 1, 1, 23, 59, 59), + dt_utc(2020, 1, 1), + dt_utc(2020, 1, 1, 23), + ), + ( + dt_utc(2020, 1, 1), + dt_utc(2020, 1, 5), + dt_utc(2020, 1, 1), + dt_utc(2020, 1, 3, 23), + ), + ( + dt_utc(2019, 1, 1), + dt_utc(2020, 1, 5), + dt_utc(2020, 1, 1), + dt_utc(2020, 1, 3, 23), + ), + ( + dt_utc(2019, 1, 1), + dt_utc(2019, 1, 5), + None, + None, + ), + ( + dt_utc(2021, 1, 1), + dt_utc(2021, 1, 5), + None, + None, + ), + ( + dt_utc(2020, 1, 2), + None, + dt_utc(2020, 1, 2), + dt_utc(2020, 1, 3, 23), + ), + ], +) +async def test_fetch_ohlcv(mocker, since, until, first_date, last_date): + history_start = dt_utc(2020, 1, 1).date() + history_end = dt_utc(2020, 1, 3).date() + candle_type = CandleType.SPOT + pair = "BTC/USDT" + timeframe = "1h" + + since_ms = dt_ts(since) + until_ms = dt_ts(until) + + mocker.patch( + "aiohttp.ClientSession.get", side_effect=make_response_from_url(history_start, history_end) + ) + df = await fetch_ohlcv(candle_type, pair, timeframe, since_ms, until_ms) + + if df.empty: + assert first_date is None and last_date is None + else: + assert df["date"].iloc[0] == first_date + assert df["date"].iloc[-1] == last_date + + +async def test_get_daily_ohlcv(mocker, testdatadir): + symbol = "BTCUSDT" + timeframe = "1h" + date = dt_utc(2024, 10, 28).date() + first_date = dt_utc(2024, 10, 28) + last_date = dt_utc(2024, 10, 28, 23) + + async with aiohttp.ClientSession() as session: + path = testdatadir / "binance/binance_public_data/spot-klines-BTCUSDT-1h-2024-10-28.zip" + mocker.patch("aiohttp.ClientSession.get", return_value=MockResponse(path.read_bytes(), 200)) + df = await get_daily_ohlcv("spot", symbol, timeframe, date, session) + assert df["date"].iloc[0] == first_date + assert df["date"].iloc[-1] == last_date + + path = ( + testdatadir / "binance/binance_public_data/futures-um-klines-BTCUSDT-1h-2024-10-28.zip" + ) + mocker.patch("aiohttp.ClientSession.get", return_value=MockResponse(path.read_bytes(), 200)) + df = await get_daily_ohlcv("futures/um", symbol, timeframe, date, session) + assert df["date"].iloc[0] == first_date + assert df["date"].iloc[-1] == last_date + + mocker.patch("aiohttp.ClientSession.get", return_value=MockResponse(b"", 404)) + df = await get_daily_ohlcv("spot", symbol, timeframe, date, session) + assert df is None + + mocker.patch("aiohttp.ClientSession.get", return_value=MockResponse(b"", 500)) + mocker.patch("asyncio.sleep") + with pytest.raises(BadHttpStatus): + df = await get_daily_ohlcv("spot", symbol, timeframe, date, session) + + mocker.patch("aiohttp.ClientSession.get", return_value=MockResponse(b"nop", 200)) + with pytest.raises(zipfile.BadZipFile): + df = await get_daily_ohlcv("spot", symbol, timeframe, date, session) diff --git a/tests/testdata/binance/binance_public_data/futures-um-klines-BTCUSDT-1h-2024-10-28.zip b/tests/testdata/binance/binance_public_data/futures-um-klines-BTCUSDT-1h-2024-10-28.zip new file mode 100644 index 0000000000000000000000000000000000000000..5bda1b27101dedd7d639726b70ebb2229076c411 GIT binary patch literal 1533 zcmaLXdpOez7zgm*%q^3V+gx)?ib;!+<}%AAOD$Q&$lYeIKo8Vn;M39+690sqtuAML%RY2K!Q_X01%wq zS&`>@gJ8zuE#C-&bZELvn>GoexhtnDkj&zK3@^JCmVx~o(3wKy)0;3 z%de{s3ljH!*#5E1|Ofq$b%cF8^@2@@cbWB;K9r1sVn@x5rQ?L#DL>et^ceZdy zK08bc8}uO$*_HHo@O%cxMD{u4P7)=^b?F$16=dr3HS-0@zDTnUEDeV8JAV%n;poan zOgY$GflhU9nR+AD>AE0eBb#{6CCBf{sO9E#)<&p$b<8+Ag;sp@L9@zCk} z+N+@_#8og#eoyKhjVm>f0CTqY1iJ<_R6l2x@;Pyi@5=VTz>QaA5os|_J{Qt+pRO2P z!WwvsHB4~~I)`kTcJL15uSSqc^n{wSZTj?pR;eDr3vb>XXP z?AH+xKF1;ai^T`hL$KC#EFbcx6x&8s6c!pbKO|nDQ5WJ-tGmMBTblBb#cTd7w#QUF z1x->BC+G4Qb?&KL=0`m>#3ZM|T9`rtTBfK1BWan|XVJ1<23*&Y%fTvrD>3}bn!&3> zI1TJQs~YlE#U!;woMg&rWN1)2{tM2aB9i&sP|IAK4054C%iKvD!rc@UJK}glUv)OU zQ7j;5i53&SAT0hDW4glMZ}K3R75vqBKoYca1P zO*47ZPcLI7=rW*a21stckSC|u$rpsjcGMpXq*oN%E4I^wc(u4fvCG=1>;2tT_0f-H zkC7?@ua%bI(4R#RS6U>Q333rP;*SaNnfo0BL|8n?8k7qeIxJpE;Y_JNZAK9 zVfE%$ye`NoBn{|9awXg%dg1j9!3W{-Ri=`k7Bk7+Nz73sZVr}rW>2Ar(u+%vjDfC$ zjV@zgLpS>*{G8-|hUllsE>%J>D0tvnt+2Q(x-2+9yMSH5xv{NUhDh`Cg3g-;wFd&H zxeL+CMvN?D?_^}(GZ!s>r&%foCqP=&%KrL9JD{WFS!qJl-M5soktM5M{rupi2wtG3 zrDRL)ypMC-Oz)v#v}~)}SPW~TS@^IZq|IoJ&Er)EaO1v9K zw+dep0!__`v2+?_(ee(g*7$Lc)Efl!{J@?}79<~wN)xXxk1yxdqXJ27r5cI}MXM%q z)%BsCb*cECTf^i$=0&r&iW-|G{XfcxXuUNiM=T2|?R#@Uin;NcqsrtCvW?E2>d?FO zQdtv!cX#dH#DpsCQXR6>zHNQ(F93%H38;eqOJdlL&;S71NeTV;^*=5P{p0e#OvYh_ SgnvJRcE)|D<9Cn;fWHCTt-L1y literal 0 HcmV?d00001 diff --git a/tests/testdata/binance/binance_public_data/spot-klines-BTCUSDT-1h-2024-10-28.zip b/tests/testdata/binance/binance_public_data/spot-klines-BTCUSDT-1h-2024-10-28.zip new file mode 100644 index 0000000000000000000000000000000000000000..b94090741ef3e3fee9aeb72c96313a5a9b351777 GIT binary patch literal 1578 zcmZ{lYdjMQ0L3Q_MP(YgH|CXH*E|=7A|cXLLLS{^+LnY+Uh|m9qlecNn&y!f@>-7> zd9P7h%S5GG#S*sNY$eNLrCr_co%7-Re&=^S|BoM1c`XJQjMb9J!t&Payc|RxUO_JHJ`R&~U zc?jvu{1--xtLuEbaPRh0F6kSjnP3~K=Ik_Zehy?Ft*yA;WGhsrrmeK9z7)5`#iQbc zE!38YH@$JmEneu6Y2b)t4Z{0;ek3j8(sS*9bvVGS+=e8lTBX$Ef#eFobHvVXN5}Yh zL0JR&2$7CoJm6iRT4P;zniJ4361h8}8lKN3pho70(;g^21s5ZxaOpSzXA%^KL@WzB4Xf((y zCP#fGw8R8bO8C)u1s&{iQjUqp=^?okVl@zUK>>9?vyPE9qZ^f{J;YIv?hmH$q#jJr zN`8RnaBciOy*ji(vWmOIGmL!0!S1jQO-hq~J)O4VzcN}k^IB=S4O_T5JnE8$u>Ku7 zs~f;c`x?C8&W(m(z}p5&95aRs&OeQwUU@%V%z}I@FdVr+sN~Y%E?1&kxvQV&1_Wdp;6qRGrOw)kDBs&b(X@ ziLReWU>MK?t5G&(_T+iGq+~<=u_XGN!zS%~NH1 zOhN-S(4;ge+F7Gdf zeQO#Fp18q;KVsyfV+REEAslFHU9!Sl9v-fenNnCMbG~uiQf0AeE;{P%w=>pi=P#ve zQL3#LwYwv!;M|UTyxWYb?4|-LAl~ZWE;WMgIN~O~B*TZ=4^|@)%76x-kW;xPpCC-z zYuC8;(rvsY?%O7} zGvYOD1uSuG=k~n6`fN=@*3aJiD5`fCI$YL+vc<#twoo|@T6Q+Uw~YQ;0%@Pgwr2M? zMW4FLm{xtCb-$3_MbsJA0vmXJ0J4^;a)N6L$eEN{GvS2*8T1D2uREkEC(bOjrck()9YU)?Aj&^YKg^D2VerY z=<2Z>pF@1Kv^VifXOK>~NA$m|KNqEUeH1sqqAP1Gv`~XaT|rFf$8R@3-7?iDrEfTJ!b$HtX?XO$M_)Ke!tFdyn8fBzJ(jVG8BB3&W??B~=)!>! z?oZ#{&~~y-r*m3$2kQc#veJBC+fjh0?0LZbkaT#`gOpnZ#=p;-#Ns4Nx#|JM2I}Hv zRq3`zQ5HYas?aw>3Q6M$BC=SJ_wwmE3bq_Ez>l|?_J5ItJ1GoayNwb}dCZB? z9K1(tFL-;kXPFl|nhW3Qt5<0j8g7*E}Jq7&hHoWtLEOfkhG zODqL%?#8QWgnr2PS9v%uO_}}boKfV8S*GFkTaCT4*Vvr`9ai`umDYfje~