+4
-4
@@ -365,24 +365,24 @@ class ChanKLC():
|
|||||||
last_top.fx_confirmed = True
|
last_top.fx_confirmed = True
|
||||||
if last_top.fx_box:
|
if last_top.fx_box:
|
||||||
last_top.fx_box.end_time = self.end_time
|
last_top.fx_box.end_time = self.end_time
|
||||||
print(self.end_time, "fx_confirmed top")
|
#print(self.end_time, "fx_confirmed top")
|
||||||
else:
|
else:
|
||||||
high = last_top.high
|
high = last_top.high
|
||||||
low = self.low
|
low = self.low
|
||||||
last_top.fx_box = Chan_FX_Box.Chan_FX_Box(last_top.pre.end_time, self.end_time, high, low)
|
last_top.fx_box = Chan_FX_Box.Chan_FX_Box(last_top.pre.end_time, self.end_time, high, low)
|
||||||
print(self.end_time, "fx_confirmed new box top")
|
#print(self.end_time, "fx_confirmed new box top")
|
||||||
elif self.in_fx == False and last_bottom.fx_confirmed == False:
|
elif self.in_fx == False and last_bottom.fx_confirmed == False:
|
||||||
pre = last_bottom.pre
|
pre = last_bottom.pre
|
||||||
if pre.high < self.close:
|
if pre.high < self.close:
|
||||||
last_bottom.fx_confirmed = True
|
last_bottom.fx_confirmed = True
|
||||||
if last_bottom.fx_box:
|
if last_bottom.fx_box:
|
||||||
last_bottom.fx_box.end_time = self.end_time
|
last_bottom.fx_box.end_time = self.end_time
|
||||||
print(self.end_time, "fx_confirmed bottom")
|
#print(self.end_time, "fx_confirmed bottom")
|
||||||
else:
|
else:
|
||||||
high = self.high
|
high = self.high
|
||||||
low = last_bottom.low
|
low = last_bottom.low
|
||||||
last_bottom.fx_box = Chan_FX_Box.Chan_FX_Box(last_bottom.pre.end_time, self.end_time, high, low)
|
last_bottom.fx_box = Chan_FX_Box.Chan_FX_Box(last_bottom.pre.end_time, self.end_time, high, low)
|
||||||
print(self.end_time, "fx_confirmed new box bottom")
|
#print(self.end_time, "fx_confirmed new box bottom")
|
||||||
def add_klu(self, klu):
|
def add_klu(self, klu):
|
||||||
self.klu_list.append(klu)
|
self.klu_list.append(klu)
|
||||||
def set_end_klu(self, klu):
|
def set_end_klu(self, klu):
|
||||||
|
|||||||
@@ -1576,9 +1576,10 @@ class TF_DF():
|
|||||||
"""
|
"""
|
||||||
根据缠论线段中枢定义计算中枢
|
根据缠论线段中枢定义计算中枢
|
||||||
从第4根线段开始(索引3),每3根线段为一组检查
|
从第4根线段开始(索引3),每3根线段为一组检查
|
||||||
后一个中枢比前一个高 -> 上涨中枢,以下跌开始、以下跌结束
|
上涨中枢:后中枢 zd > 前中枢 zg(不重叠上移)
|
||||||
后一个中枢比前一个低 -> 下跌中枢,以上涨开始、以上涨结束
|
下跌中枢:后中枢 zg < 前中枢 zd(不重叠下移)
|
||||||
中枢可以扩展到5根、7根...
|
盘整/扩张:后中枢与前中枢整体区间(GG/DD)有交集
|
||||||
|
中枢可按两段一组继续扩展到5根、7根...
|
||||||
"""
|
"""
|
||||||
zs_list = []
|
zs_list = []
|
||||||
if len(seg_list) < 3:
|
if len(seg_list) < 3:
|
||||||
@@ -1609,32 +1610,34 @@ class TF_DF():
|
|||||||
|
|
||||||
if zg <= zd:
|
if zg <= zd:
|
||||||
start_idx += 1
|
start_idx += 1
|
||||||
|
#print(seg1.start_bi.start_klc.end_time, "not valid", zg, zd)
|
||||||
continue
|
continue
|
||||||
|
|
||||||
# 判断中枢类型
|
# 判断中枢类型(按注释定义)
|
||||||
# 上涨中枢:以下跌开始、以下跌结束(下跌+上涨+下跌)
|
# 上涨中枢:后中枢 zd > 前中枢 zg(不重叠上移)
|
||||||
# 下跌中枢:以上涨开始、以上涨结束(上涨+下跌+上涨)
|
# 下跌中枢:后中枢 zg < 前中枢 zd(不重叠下移)
|
||||||
|
# 盘整/扩张:后中枢与前中枢区间有交集
|
||||||
if last_zs is None:
|
if last_zs is None:
|
||||||
# 第一个中枢
|
# 第一个中枢仅按线段形态判定方向
|
||||||
if seg1.dir == Chan_SEG_DIR.DOWN:
|
if seg1.dir == Chan_SEG_DIR.DOWN:
|
||||||
# 下跌开始 -> 上涨中枢
|
# 下跌+上涨+下跌,对应上涨中枢
|
||||||
zs_dir = Chan_ZS_DIR.DOWN
|
zs_dir = Chan_ZS_DIR.UP
|
||||||
# 验证模式:下跌+上涨+下跌
|
|
||||||
valid = (seg2.dir == Chan_SEG_DIR.UP and seg3.dir == Chan_SEG_DIR.DOWN)
|
valid = (seg2.dir == Chan_SEG_DIR.UP and seg3.dir == Chan_SEG_DIR.DOWN)
|
||||||
else:
|
else:
|
||||||
# 上涨开始 -> 下跌中枢
|
# 上涨+下跌+上涨,对应下跌中枢
|
||||||
zs_dir = Chan_ZS_DIR.UP
|
zs_dir = Chan_ZS_DIR.DOWN
|
||||||
# 验证模式:上涨+下跌+上涨
|
|
||||||
valid = (seg2.dir == Chan_SEG_DIR.DOWN and seg3.dir == Chan_SEG_DIR.UP)
|
valid = (seg2.dir == Chan_SEG_DIR.DOWN and seg3.dir == Chan_SEG_DIR.UP)
|
||||||
else:
|
else:
|
||||||
# 根据与前一个中枢的高低比较判断
|
is_up_zs = zd > last_zs.zg
|
||||||
if zg > last_zs.zg:
|
is_down_zs = zg < last_zs.zd
|
||||||
# 上涨中枢:以下跌开始、以下跌结束
|
|
||||||
zs_dir = Chan_ZS_DIR.DOWN
|
if is_up_zs:
|
||||||
valid = (seg1.dir == Chan_SEG_DIR.DOWN and seg2.dir == Chan_SEG_DIR.UP and seg3.dir == Chan_SEG_DIR.DOWN)
|
# 不重叠上移
|
||||||
else:
|
|
||||||
# 下跌中枢:以上涨开始、以上涨结束
|
|
||||||
zs_dir = Chan_ZS_DIR.UP
|
zs_dir = Chan_ZS_DIR.UP
|
||||||
|
valid = (seg1.dir == Chan_SEG_DIR.DOWN and seg2.dir == Chan_SEG_DIR.UP and seg3.dir == Chan_SEG_DIR.DOWN)
|
||||||
|
elif is_down_zs:
|
||||||
|
# 不重叠下移
|
||||||
|
zs_dir = Chan_ZS_DIR.DOWN
|
||||||
valid = (seg1.dir == Chan_SEG_DIR.UP and seg2.dir == Chan_SEG_DIR.DOWN and seg3.dir == Chan_SEG_DIR.UP)
|
valid = (seg1.dir == Chan_SEG_DIR.UP and seg2.dir == Chan_SEG_DIR.DOWN and seg3.dir == Chan_SEG_DIR.UP)
|
||||||
|
|
||||||
# 验证是否有效
|
# 验证是否有效
|
||||||
@@ -1642,98 +1645,6 @@ class TF_DF():
|
|||||||
start_idx += 1
|
start_idx += 1
|
||||||
continue
|
continue
|
||||||
|
|
||||||
# 检查是否与前一个中枢重叠
|
|
||||||
if last_zs:
|
|
||||||
# 判断是否有重叠
|
|
||||||
overlap = (zg >= last_zs.zd and zd <= last_zs.zg)
|
|
||||||
|
|
||||||
if overlap:
|
|
||||||
# 有重叠,扩展中枢到5根、7根...(缠论:合并为同一中枢)
|
|
||||||
# 本组先纳入当前 3 根,再向后逐根尝试;遇到与 [zd,zg] 不重叠(离开中枢)则停止扩展
|
|
||||||
added_segs = [seg_list[start_idx], seg_list[start_idx + 1], seg_list[start_idx + 2]]
|
|
||||||
cur_idx = start_idx + 3
|
|
||||||
|
|
||||||
while cur_idx < len(seg_list):
|
|
||||||
next_seg = seg_list[cur_idx]
|
|
||||||
if not next_seg.is_sure:
|
|
||||||
break
|
|
||||||
|
|
||||||
# 扩展条件:新线段与中枢区间 [zd, zg] 有重叠即并入;不重叠则停止,离开中枢的线段不包含
|
|
||||||
# 用起止笔的极值算线段区间,避免 seg.high/seg.low 在个别线段上未同步导致的误判
|
|
||||||
seg_high = max(next_seg.start_bi.high, next_seg.end_bi.high) if next_seg.end_bi else next_seg.start_bi.high
|
|
||||||
seg_low = min(next_seg.start_bi.low, next_seg.end_bi.low) if next_seg.end_bi else next_seg.start_bi.low
|
|
||||||
overlap_with_zs = (seg_high >= last_zs.zd and seg_low <= last_zs.zg)
|
|
||||||
if not overlap_with_zs:
|
|
||||||
break
|
|
||||||
added_segs.append(next_seg)
|
|
||||||
cur_idx += 1
|
|
||||||
|
|
||||||
# 扩展中枢 = 原中枢线段 + 本组并入的线段(缠论合并)
|
|
||||||
segs_for_zs = list(last_zs.seg_list) + list(added_segs)
|
|
||||||
|
|
||||||
# 中枢开始与结束线段方向一致:上涨中枢结束于 DOWN,下跌中枢结束于 UP
|
|
||||||
required_end_seg_dir = Chan_SEG_DIR.DOWN if last_zs.dir == Chan_ZS_DIR.DOWN else Chan_SEG_DIR.UP
|
|
||||||
while len(segs_for_zs) >= 3 and segs_for_zs[-1].dir != required_end_seg_dir:
|
|
||||||
segs_for_zs.pop()
|
|
||||||
|
|
||||||
# 扩展时只更新 gg、dd 和 seg_list;zg、zd 由前 3 根线段确定,不随扩展改变
|
|
||||||
seg_highs = [s.high for s in segs_for_zs]
|
|
||||||
seg_lows = [s.low for s in segs_for_zs]
|
|
||||||
last_zs.set_gg(max(seg_highs))
|
|
||||||
last_zs.set_dd(min(seg_lows))
|
|
||||||
last_zs.seg_list = segs_for_zs
|
|
||||||
|
|
||||||
# 更新结束时间(以裁剪后的最后一段为准)
|
|
||||||
last_seg = segs_for_zs[-1]
|
|
||||||
if last_seg.end_bi:
|
|
||||||
last_zs.set_end_klc(last_seg.end_bi.end_klc, last_seg.sure_time, 0, last_seg)
|
|
||||||
last_zs.set_end_seg(last_seg)
|
|
||||||
|
|
||||||
# 下一组从「裁剪后保留的最后一段」的下一条线段开始,被裁掉的线段会参与下一组 3 根,避免漏到下一中枢
|
|
||||||
last_kept_idx = cur_idx - 1
|
|
||||||
for i in range(len(seg_list)):
|
|
||||||
if seg_list[i] == last_seg:
|
|
||||||
last_kept_idx = i
|
|
||||||
break
|
|
||||||
start_idx = last_kept_idx + 1
|
|
||||||
continue
|
|
||||||
else:
|
|
||||||
# 本组3段整体与前中枢不重叠,仍逐根检查:若某线段与 [zd,zg] 重叠(如离开后回抽回到前中枢)则并入扩展,不直接新中枢
|
|
||||||
added_after_leave = []
|
|
||||||
for k in range(start_idx, min(start_idx + 3, len(seg_list))):
|
|
||||||
s = seg_list[k]
|
|
||||||
if not s.is_sure:
|
|
||||||
break
|
|
||||||
sh = max(s.start_bi.high, s.end_bi.high) if s.end_bi else s.start_bi.high
|
|
||||||
sl = min(s.start_bi.low, s.end_bi.low) if s.end_bi else s.start_bi.low
|
|
||||||
if sh >= last_zs.zd and sl <= last_zs.zg:
|
|
||||||
added_after_leave.append(s)
|
|
||||||
else:
|
|
||||||
break
|
|
||||||
if added_after_leave:
|
|
||||||
segs_for_zs = list(last_zs.seg_list) + list(added_after_leave)
|
|
||||||
required_end_seg_dir = Chan_SEG_DIR.DOWN if last_zs.dir == Chan_ZS_DIR.DOWN else Chan_SEG_DIR.UP
|
|
||||||
while len(segs_for_zs) >= 3 and segs_for_zs[-1].dir != required_end_seg_dir:
|
|
||||||
segs_for_zs.pop()
|
|
||||||
seg_highs = [s.high for s in segs_for_zs]
|
|
||||||
seg_lows = [s.low for s in segs_for_zs]
|
|
||||||
last_zs.set_gg(max(seg_highs))
|
|
||||||
last_zs.set_dd(min(seg_lows))
|
|
||||||
last_zs.seg_list = segs_for_zs
|
|
||||||
last_seg = segs_for_zs[-1]
|
|
||||||
if last_seg.end_bi:
|
|
||||||
last_zs.set_end_klc(last_seg.end_bi.end_klc, last_seg.sure_time, 0, last_seg)
|
|
||||||
last_zs.set_end_seg(last_seg)
|
|
||||||
start_idx = start_idx + len(added_after_leave)
|
|
||||||
continue
|
|
||||||
# 没有任何线段与前中枢重叠,确认前中枢结束并创建新中枢
|
|
||||||
if last_zs.seg_list and len(last_zs.seg_list) > 0:
|
|
||||||
prev_zs_last_seg = last_zs.seg_list[-1]
|
|
||||||
if prev_zs_last_seg.end_bi:
|
|
||||||
last_zs.set_end_klc(prev_zs_last_seg.end_bi.end_klc, prev_zs_last_seg.sure_time, 0, prev_zs_last_seg)
|
|
||||||
last_zs.set_end_seg(prev_zs_last_seg)
|
|
||||||
last_zs.is_sure = True
|
|
||||||
|
|
||||||
# 创建新中枢
|
# 创建新中枢
|
||||||
gg = max(seg1.high, seg2.high, seg3.high)
|
gg = max(seg1.high, seg2.high, seg3.high)
|
||||||
dd = min(seg1.low, seg2.low, seg3.low)
|
dd = min(seg1.low, seg2.low, seg3.low)
|
||||||
@@ -1743,11 +1654,41 @@ class TF_DF():
|
|||||||
zs.set_zd(zd)
|
zs.set_zd(zd)
|
||||||
zs.set_gg(gg)
|
zs.set_gg(gg)
|
||||||
zs.set_dd(dd)
|
zs.set_dd(dd)
|
||||||
zs.set_end_klc(seg3.end_bi.end_klc, seg3.sure_time, 0, seg3)
|
|
||||||
zs.set_end_seg(seg3)
|
|
||||||
zs.is_sure = False
|
zs.is_sure = False
|
||||||
zs.seg_list = [seg1, seg2, seg3]
|
zs.seg_list = [seg1, seg2, seg3]
|
||||||
|
# 若第二线段与 [zd,zg] 重叠(如离开后回抽回到前中枢)则并入扩展
|
||||||
|
added_after_leave = []
|
||||||
|
leave_index = start_idx + 4
|
||||||
|
while leave_index < len(seg_list):
|
||||||
|
s = seg_list[leave_index]
|
||||||
|
if not s.is_sure:
|
||||||
|
break
|
||||||
|
sh = max(s.start_bi.high, s.end_bi.high) if s.end_bi else s.start_bi.high
|
||||||
|
sl = min(s.start_bi.low, s.end_bi.low) if s.end_bi else s.start_bi.low
|
||||||
|
if sh >= zs.zd and sl <= zs.zg:
|
||||||
|
added_after_leave.append(s.pre)
|
||||||
|
added_after_leave.append(s)
|
||||||
|
else:
|
||||||
|
break
|
||||||
|
leave_index += 2
|
||||||
|
if added_after_leave:
|
||||||
|
#print(len(added_after_leave))
|
||||||
|
segs_for_zs = list(zs.seg_list) + list(added_after_leave)
|
||||||
|
seg_highs = [s.high for s in segs_for_zs]
|
||||||
|
seg_lows = [s.low for s in segs_for_zs]
|
||||||
|
zs.set_gg(max(seg_highs))
|
||||||
|
zs.set_dd(min(seg_lows))
|
||||||
|
zs.seg_list = segs_for_zs
|
||||||
|
seg = segs_for_zs[-1]
|
||||||
|
if seg.end_bi:
|
||||||
|
zs.set_end_klc(seg.end_bi.end_klc, seg.sure_time, 0, seg)
|
||||||
|
zs.set_end_seg(seg)
|
||||||
|
zs.is_sure = True
|
||||||
|
start_idx = start_idx + len(added_after_leave)
|
||||||
|
else:
|
||||||
|
zs.set_end_klc(seg3.end_bi.end_klc, seg3.sure_time, 0, seg3)
|
||||||
|
zs.set_end_seg(seg3)
|
||||||
|
zs.is_sure = True
|
||||||
if last_zs:
|
if last_zs:
|
||||||
last_zs.set_next(zs)
|
last_zs.set_next(zs)
|
||||||
zs.set_pre(last_zs)
|
zs.set_pre(last_zs)
|
||||||
@@ -1756,8 +1697,9 @@ class TF_DF():
|
|||||||
last_zs = zs
|
last_zs = zs
|
||||||
|
|
||||||
# 移动到下一组
|
# 移动到下一组
|
||||||
start_idx += 3
|
start_idx += 4
|
||||||
|
if last_zs:
|
||||||
|
last_zs.is_sure = seg_list[-1].is_sure
|
||||||
# 处理最后一个未确认的中枢 - 不自动扩展,保持未完成状态
|
# 处理最后一个未确认的中枢 - 不自动扩展,保持未完成状态
|
||||||
if last_zs and not last_zs.is_sure:
|
if last_zs and not last_zs.is_sure:
|
||||||
# 获取中枢最后一个线段的索引
|
# 获取中枢最后一个线段的索引
|
||||||
|
|||||||
+53
-2
@@ -1,3 +1,8 @@
|
|||||||
|
"""
|
||||||
|
Chan 数据服务:用 ccxt 从交易所拉取 K 线,内存缓存 + CSV 落盘;
|
||||||
|
后台线程定期增量刷新,断线时记录 resume_since 以免漏 K;
|
||||||
|
配置中的基础周期(如 1m/1h)可合成 DERIVED_TIMEFRAME_PLAN 中的衍生周期。
|
||||||
|
"""
|
||||||
import asyncio
|
import asyncio
|
||||||
import csv
|
import csv
|
||||||
import json
|
import json
|
||||||
@@ -20,14 +25,17 @@ from technical.util import resample_to_interval
|
|||||||
# docker compose logs --tail=200
|
# docker compose logs --tail=200
|
||||||
# docker compose down && docker compose build --no-cache && docker compose up -d
|
# docker compose down && docker compose build --no-cache && docker compose up -d
|
||||||
|
|
||||||
|
# 基础周期枚举顺序(用于衍生周期展示顺序);仅允许集合内周期作为交易所直接拉取的 tf
|
||||||
TIMEFRAME_ORDER = ["1m", "1h", "1d", "1w"]
|
TIMEFRAME_ORDER = ["1m", "1h", "1d", "1w"]
|
||||||
ALLOWED_TIMEFRAMES = set(TIMEFRAME_ORDER)
|
ALLOWED_TIMEFRAMES = set(TIMEFRAME_ORDER)
|
||||||
|
# 各基础周期一根 K 线的毫秒长度(用于历史分页与断线回退)
|
||||||
TIMEFRAME_TO_MS: Dict[str, int] = {
|
TIMEFRAME_TO_MS: Dict[str, int] = {
|
||||||
"1m": 60_000,
|
"1m": 60_000,
|
||||||
"1h": 3_600_000,
|
"1h": 3_600_000,
|
||||||
"1d": 86_400_000,
|
"1d": 86_400_000,
|
||||||
"1w": 604_800_000,
|
"1w": 604_800_000,
|
||||||
}
|
}
|
||||||
|
# 每个基础周期可派生出的合成周期列表(由该基础周期 K 线 resample 得到)
|
||||||
DERIVED_TIMEFRAME_PLAN: Dict[str, List[str]] = {
|
DERIVED_TIMEFRAME_PLAN: Dict[str, List[str]] = {
|
||||||
"1m": ["2m", "3m", "4m", "5m", "10m", "15m", "20m", "25m", "30m", "45m"],
|
"1m": ["2m", "3m", "4m", "5m", "10m", "15m", "20m", "25m", "30m", "45m"],
|
||||||
"1h": ["2h", "3h", "4h", "5h", "6h", "7h", "8h", "9h", "10", "11h", "12h", "16h", "20h"],
|
"1h": ["2h", "3h", "4h", "5h", "6h", "7h", "8h", "9h", "10", "11h", "12h", "16h", "20h"],
|
||||||
@@ -37,8 +45,8 @@ DERIVED_TIMEFRAME_PLAN: Dict[str, List[str]] = {
|
|||||||
CSV_FIELDNAMES = ["timestamp", "datetime", "open", "high", "low", "close", "volume"]
|
CSV_FIELDNAMES = ["timestamp", "datetime", "open", "high", "low", "close", "volume"]
|
||||||
DEFAULT_LIMIT = 500
|
DEFAULT_LIMIT = 500
|
||||||
RECENT_CANDLE_LIMIT = 10
|
RECENT_CANDLE_LIMIT = 10
|
||||||
RECENT_FETCH_INTERVAL = 5
|
RECENT_FETCH_INTERVAL = 5 # 后台刷新循环休眠秒数
|
||||||
PERSIST_INTERVAL = 600
|
PERSIST_INTERVAL = 600 # 全量落盘周期(秒)
|
||||||
|
|
||||||
|
|
||||||
logger = logging.getLogger("data_provider")
|
logger = logging.getLogger("data_provider")
|
||||||
@@ -49,11 +57,13 @@ logging.basicConfig(
|
|||||||
|
|
||||||
|
|
||||||
def to_utc_iso(timestamp_ms: int) -> str:
|
def to_utc_iso(timestamp_ms: int) -> str:
|
||||||
|
"""将毫秒时间戳格式化为 UTC ISO 字符串(末尾 Z)。"""
|
||||||
dt = datetime.fromtimestamp(timestamp_ms / 1000, tz=timezone.utc)
|
dt = datetime.fromtimestamp(timestamp_ms / 1000, tz=timezone.utc)
|
||||||
return dt.isoformat().replace("+00:00", "Z")
|
return dt.isoformat().replace("+00:00", "Z")
|
||||||
|
|
||||||
|
|
||||||
def parse_timestamp(value: Optional[object]) -> Optional[int]:
|
def parse_timestamp(value: Optional[object]) -> Optional[int]:
|
||||||
|
"""解析查询参数中的时间为 UTC 毫秒时间戳;支持数字或 ISO 字符串。"""
|
||||||
if value is None:
|
if value is None:
|
||||||
return None
|
return None
|
||||||
if isinstance(value, (int, float)):
|
if isinstance(value, (int, float)):
|
||||||
@@ -79,6 +89,7 @@ def parse_timestamp(value: Optional[object]) -> Optional[int]:
|
|||||||
|
|
||||||
|
|
||||||
def candle_to_dict(candle: Iterable[float]) -> Dict[str, float]:
|
def candle_to_dict(candle: Iterable[float]) -> Dict[str, float]:
|
||||||
|
"""ccxt OHLCV 单根 [ts, o, h, l, c, v] 转为内部字典结构。"""
|
||||||
ts = int(candle[0])
|
ts = int(candle[0])
|
||||||
return {
|
return {
|
||||||
"timestamp": ts,
|
"timestamp": ts,
|
||||||
@@ -92,6 +103,7 @@ def candle_to_dict(candle: Iterable[float]) -> Dict[str, float]:
|
|||||||
|
|
||||||
|
|
||||||
def timeframe_to_minutes(tf: str) -> Optional[int]:
|
def timeframe_to_minutes(tf: str) -> Optional[int]:
|
||||||
|
"""将如 15m、2h 转为「分钟数」,供 resample 与衍生周期计算。"""
|
||||||
if not tf:
|
if not tf:
|
||||||
return None
|
return None
|
||||||
unit = tf[-1]
|
unit = tf[-1]
|
||||||
@@ -111,6 +123,8 @@ def timeframe_to_minutes(tf: str) -> Optional[int]:
|
|||||||
|
|
||||||
|
|
||||||
class DataProvider:
|
class DataProvider:
|
||||||
|
"""封装交易所连接、本地 CSV、内存缓存、断线恢复与衍生周期聚合。"""
|
||||||
|
|
||||||
def __init__(self, config_path: Path) -> None:
|
def __init__(self, config_path: Path) -> None:
|
||||||
self.config_path = config_path
|
self.config_path = config_path
|
||||||
self.config = self._load_config()
|
self.config = self._load_config()
|
||||||
@@ -126,10 +140,12 @@ class DataProvider:
|
|||||||
self.data: Dict[str, Dict[str, List[Dict[str, float]]]] = {
|
self.data: Dict[str, Dict[str, List[Dict[str, float]]]] = {
|
||||||
symbol: {tf: [] for tf in self.timeframes} for symbol in self.symbols
|
symbol: {tf: [] for tf in self.timeframes} for symbol in self.symbols
|
||||||
}
|
}
|
||||||
|
# 衍生周期 -> 用于合成的交易所基础周期(每个衍生只对应一个 base)
|
||||||
self.derived_map: Dict[str, str] = {}
|
self.derived_map: Dict[str, str] = {}
|
||||||
for base_tf in self.timeframes:
|
for base_tf in self.timeframes:
|
||||||
for derived_tf in DERIVED_TIMEFRAME_PLAN.get(base_tf, []):
|
for derived_tf in DERIVED_TIMEFRAME_PLAN.get(base_tf, []):
|
||||||
self.derived_map.setdefault(derived_tf, base_tf)
|
self.derived_map.setdefault(derived_tf, base_tf)
|
||||||
|
# 衍生周期展示顺序:按 TIMEFRAME_ORDER 中的基础周期依次展开
|
||||||
derived_order: List[str] = []
|
derived_order: List[str] = []
|
||||||
for base_tf in TIMEFRAME_ORDER:
|
for base_tf in TIMEFRAME_ORDER:
|
||||||
if base_tf not in self.timeframes:
|
if base_tf not in self.timeframes:
|
||||||
@@ -151,12 +167,14 @@ class DataProvider:
|
|||||||
self._load_resume_since()
|
self._load_resume_since()
|
||||||
|
|
||||||
def _load_config(self) -> Dict[str, object]:
|
def _load_config(self) -> Dict[str, object]:
|
||||||
|
"""读取 JSON 配置文件。"""
|
||||||
if not self.config_path.exists():
|
if not self.config_path.exists():
|
||||||
raise FileNotFoundError(f"未找到配置文件: {self.config_path}")
|
raise FileNotFoundError(f"未找到配置文件: {self.config_path}")
|
||||||
with self.config_path.open("r", encoding="utf-8") as fp:
|
with self.config_path.open("r", encoding="utf-8") as fp:
|
||||||
return json.load(fp)
|
return json.load(fp)
|
||||||
|
|
||||||
def _load_symbols(self, config: Dict[str, object]) -> List[str]:
|
def _load_symbols(self, config: Dict[str, object]) -> List[str]:
|
||||||
|
"""从 symbols 列表、逗号分隔字符串或单字段 symbol 解析交易对,去重保序。"""
|
||||||
raw_symbols: List[str] = []
|
raw_symbols: List[str] = []
|
||||||
symbols_value = config.get("symbols")
|
symbols_value = config.get("symbols")
|
||||||
if isinstance(symbols_value, list):
|
if isinstance(symbols_value, list):
|
||||||
@@ -175,6 +193,7 @@ class DataProvider:
|
|||||||
return unique
|
return unique
|
||||||
|
|
||||||
def _validate_timeframes(self, configured: Optional[Iterable[str]]) -> List[str]:
|
def _validate_timeframes(self, configured: Optional[Iterable[str]]) -> List[str]:
|
||||||
|
"""校验周期在允许集合内;未配置则默认 TIMEFRAME_ORDER 全部;顺序优先按 TIMEFRAME_ORDER。"""
|
||||||
if not configured:
|
if not configured:
|
||||||
return list(TIMEFRAME_ORDER)
|
return list(TIMEFRAME_ORDER)
|
||||||
invalid = [tf for tf in configured if tf not in ALLOWED_TIMEFRAMES]
|
invalid = [tf for tf in configured if tf not in ALLOWED_TIMEFRAMES]
|
||||||
@@ -193,6 +212,7 @@ class DataProvider:
|
|||||||
return unique
|
return unique
|
||||||
|
|
||||||
def _init_exchange(self):
|
def _init_exchange(self):
|
||||||
|
"""实例化 ccxt 交易所,币安期货默认 defaultType=future,并 load_markets。"""
|
||||||
if not hasattr(ccxt, self.exchange_name):
|
if not hasattr(ccxt, self.exchange_name):
|
||||||
raise ValueError(f"不支持的交易所: {self.exchange_name}")
|
raise ValueError(f"不支持的交易所: {self.exchange_name}")
|
||||||
exchange_class = getattr(ccxt, self.exchange_name)
|
exchange_class = getattr(ccxt, self.exchange_name)
|
||||||
@@ -204,10 +224,12 @@ class DataProvider:
|
|||||||
return exchange
|
return exchange
|
||||||
|
|
||||||
def _data_file_path(self, symbol: str, timeframe: str) -> Path:
|
def _data_file_path(self, symbol: str, timeframe: str) -> Path:
|
||||||
|
"""单交易对单周期的 CSV 路径:data_dir/tf/exchange_symbol_tf.csv。"""
|
||||||
symbol_safe = symbol.replace("/", "_").replace(":", "_")
|
symbol_safe = symbol.replace("/", "_").replace(":", "_")
|
||||||
return self.data_dir / timeframe / f"{self.exchange.id}_{symbol_safe}_{timeframe}.csv"
|
return self.data_dir / timeframe / f"{self.exchange.id}_{symbol_safe}_{timeframe}.csv"
|
||||||
|
|
||||||
def _load_local(self, symbol: str, timeframe: str) -> List[Dict[str, float]]:
|
def _load_local(self, symbol: str, timeframe: str) -> List[Dict[str, float]]:
|
||||||
|
"""启动时从磁盘加载已有 K 线,损坏行跳过,按时间排序。"""
|
||||||
path = self._data_file_path(symbol, timeframe)
|
path = self._data_file_path(symbol, timeframe)
|
||||||
if not path.exists():
|
if not path.exists():
|
||||||
return []
|
return []
|
||||||
@@ -239,6 +261,7 @@ class DataProvider:
|
|||||||
base: List[Dict[str, float]],
|
base: List[Dict[str, float]],
|
||||||
new_candles: Iterable[Iterable[float]],
|
new_candles: Iterable[Iterable[float]],
|
||||||
) -> List[Dict[str, float]]:
|
) -> List[Dict[str, float]]:
|
||||||
|
"""按 timestamp 去重合并,新数据覆盖同时间戳旧数据。"""
|
||||||
merged = {entry["timestamp"]: entry for entry in base}
|
merged = {entry["timestamp"]: entry for entry in base}
|
||||||
for candle in new_candles:
|
for candle in new_candles:
|
||||||
entry = candle_to_dict(candle)
|
entry = candle_to_dict(candle)
|
||||||
@@ -248,6 +271,7 @@ class DataProvider:
|
|||||||
return ordered
|
return ordered
|
||||||
|
|
||||||
def _write_to_disk(self, symbol: str, timeframe: str, data: List[Dict[str, float]]) -> None:
|
def _write_to_disk(self, symbol: str, timeframe: str, data: List[Dict[str, float]]) -> None:
|
||||||
|
"""先写临时文件再 replace,避免写入中断导致 CSV 损坏。"""
|
||||||
path = self._data_file_path(symbol, timeframe)
|
path = self._data_file_path(symbol, timeframe)
|
||||||
path.parent.mkdir(parents=True, exist_ok=True)
|
path.parent.mkdir(parents=True, exist_ok=True)
|
||||||
tmp_path = path.with_suffix(path.suffix + ".tmp")
|
tmp_path = path.with_suffix(path.suffix + ".tmp")
|
||||||
@@ -266,6 +290,7 @@ class DataProvider:
|
|||||||
logger.info("交易对 %s 时间周期 %s 已写入磁盘 (%s 根K线)", symbol, timeframe, len(data))
|
logger.info("交易对 %s 时间周期 %s 已写入磁盘 (%s 根K线)", symbol, timeframe, len(data))
|
||||||
|
|
||||||
def _fetch_history(self, symbol: str, timeframe: str, since_ms: int) -> List[List[float]]:
|
def _fetch_history(self, symbol: str, timeframe: str, since_ms: int) -> List[List[float]]:
|
||||||
|
"""从 since_ms 分页拉取直到接近当前时间;遇限频则 sleep 重试。"""
|
||||||
results: List[List[float]] = []
|
results: List[List[float]] = []
|
||||||
limit = 1500
|
limit = 1500
|
||||||
now_ms = self.exchange.milliseconds()
|
now_ms = self.exchange.milliseconds()
|
||||||
@@ -302,6 +327,7 @@ class DataProvider:
|
|||||||
return results
|
return results
|
||||||
|
|
||||||
def initialize(self) -> None:
|
def initialize(self) -> None:
|
||||||
|
"""阻塞式启动:加载本地、从倒数第二根或配置起点补历史、写盘并 set _ready。"""
|
||||||
logger.info("开始初始化数据提供商")
|
logger.info("开始初始化数据提供商")
|
||||||
for symbol in self.symbols:
|
for symbol in self.symbols:
|
||||||
for timeframe in self.timeframes:
|
for timeframe in self.timeframes:
|
||||||
@@ -310,6 +336,7 @@ class DataProvider:
|
|||||||
last_ts = existing[-1]["timestamp"] if existing else None
|
last_ts = existing[-1]["timestamp"] if existing else None
|
||||||
if last_ts is not None:
|
if last_ts is not None:
|
||||||
if len(existing) >= 2:
|
if len(existing) >= 2:
|
||||||
|
# 从倒数第二根起拉,避免最后一根未收盘重复/缺口
|
||||||
fetch_since = existing[-2]["timestamp"]
|
fetch_since = existing[-2]["timestamp"]
|
||||||
else:
|
else:
|
||||||
fetch_since = max(0, last_ts - tf_ms)
|
fetch_since = max(0, last_ts - tf_ms)
|
||||||
@@ -334,9 +361,11 @@ class DataProvider:
|
|||||||
logger.info("数据初始化完成")
|
logger.info("数据初始化完成")
|
||||||
|
|
||||||
def resample_df(self, df: pd.DataFrame, interval: int) -> pd.DataFrame:
|
def resample_df(self, df: pd.DataFrame, interval: int) -> pd.DataFrame:
|
||||||
|
"""将基础周期 DataFrame 聚合为 interval 分钟周期(freqtrade technical.util)。"""
|
||||||
return resample_to_interval(df, interval)
|
return resample_to_interval(df, interval)
|
||||||
|
|
||||||
def _save_resume_since(self) -> None:
|
def _save_resume_since(self) -> None:
|
||||||
|
"""将断线恢复点持久化到 resume_since.json(原子替换)。"""
|
||||||
path = self._resume_file
|
path = self._resume_file
|
||||||
path.parent.mkdir(parents=True, exist_ok=True)
|
path.parent.mkdir(parents=True, exist_ok=True)
|
||||||
tmp_path = path.with_suffix(path.suffix + ".tmp")
|
tmp_path = path.with_suffix(path.suffix + ".tmp")
|
||||||
@@ -358,6 +387,7 @@ class DataProvider:
|
|||||||
logger.debug("恢复点已保存到磁盘: %s", path)
|
logger.debug("恢复点已保存到磁盘: %s", path)
|
||||||
|
|
||||||
def _load_resume_since(self) -> None:
|
def _load_resume_since(self) -> None:
|
||||||
|
"""启动时加载恢复点;与内存合并时取更早的 since,避免漏拉。"""
|
||||||
path = self._resume_file
|
path = self._resume_file
|
||||||
if not path.exists():
|
if not path.exists():
|
||||||
return
|
return
|
||||||
@@ -395,10 +425,12 @@ class DataProvider:
|
|||||||
logger.info("已加载恢复点: %s", path)
|
logger.info("已加载恢复点: %s", path)
|
||||||
|
|
||||||
def _get_resume_since(self, symbol: str, timeframe: str) -> Optional[int]:
|
def _get_resume_since(self, symbol: str, timeframe: str) -> Optional[int]:
|
||||||
|
"""若曾断线,返回应从哪一毫秒起补拉该 symbol/tf。"""
|
||||||
with self._lock:
|
with self._lock:
|
||||||
return self._resume_since.get(symbol, {}).get(timeframe)
|
return self._resume_since.get(symbol, {}).get(timeframe)
|
||||||
|
|
||||||
def _set_resume_since(self, symbol: str, timeframe: str, since_ms: int) -> None:
|
def _set_resume_since(self, symbol: str, timeframe: str, since_ms: int) -> None:
|
||||||
|
"""断线时写入恢复点(取更早的 since 以免漏数据),并持久化到磁盘。"""
|
||||||
with self._lock:
|
with self._lock:
|
||||||
per_symbol = self._resume_since.setdefault(symbol, {})
|
per_symbol = self._resume_since.setdefault(symbol, {})
|
||||||
prev = per_symbol.get(timeframe)
|
prev = per_symbol.get(timeframe)
|
||||||
@@ -416,6 +448,7 @@ class DataProvider:
|
|||||||
self._save_resume_since()
|
self._save_resume_since()
|
||||||
|
|
||||||
def _clear_resume_since(self, symbol: str, timeframe: str) -> None:
|
def _clear_resume_since(self, symbol: str, timeframe: str) -> None:
|
||||||
|
"""补数成功后清除该 symbol/tf 的恢复点。"""
|
||||||
with self._lock:
|
with self._lock:
|
||||||
if symbol in self._resume_since and timeframe in self._resume_since[symbol]:
|
if symbol in self._resume_since and timeframe in self._resume_since[symbol]:
|
||||||
del self._resume_since[symbol][timeframe]
|
del self._resume_since[symbol][timeframe]
|
||||||
@@ -426,6 +459,7 @@ class DataProvider:
|
|||||||
self._save_resume_since()
|
self._save_resume_since()
|
||||||
|
|
||||||
def start_background_workers(self) -> None:
|
def start_background_workers(self) -> None:
|
||||||
|
"""启动增量刷新线程与周期性落盘线程。"""
|
||||||
if self._fetch_thread and self._fetch_thread.is_alive():
|
if self._fetch_thread and self._fetch_thread.is_alive():
|
||||||
return
|
return
|
||||||
self._stop_event.clear()
|
self._stop_event.clear()
|
||||||
@@ -436,6 +470,7 @@ class DataProvider:
|
|||||||
logger.info("后台线程已启动")
|
logger.info("后台线程已启动")
|
||||||
|
|
||||||
def stop(self) -> None:
|
def stop(self) -> None:
|
||||||
|
"""停止后台线程(应用关闭时 lifespan finally 调用)。"""
|
||||||
self._stop_event.set()
|
self._stop_event.set()
|
||||||
if self._fetch_thread:
|
if self._fetch_thread:
|
||||||
self._fetch_thread.join(timeout=5)
|
self._fetch_thread.join(timeout=5)
|
||||||
@@ -444,6 +479,7 @@ class DataProvider:
|
|||||||
logger.info("数据提供商已停止")
|
logger.info("数据提供商已停止")
|
||||||
|
|
||||||
def _refresh_loop(self) -> None:
|
def _refresh_loop(self) -> None:
|
||||||
|
"""轮询各 symbol/tf:有恢复点则先补历史,否则 fetch 最近 RECENT_CANDLE_LIMIT 根。"""
|
||||||
while not self._stop_event.is_set():
|
while not self._stop_event.is_set():
|
||||||
for symbol in self.symbols:
|
for symbol in self.symbols:
|
||||||
for timeframe in self.timeframes:
|
for timeframe in self.timeframes:
|
||||||
@@ -496,10 +532,12 @@ class DataProvider:
|
|||||||
break
|
break
|
||||||
|
|
||||||
def _persist_loop(self) -> None:
|
def _persist_loop(self) -> None:
|
||||||
|
"""每隔 PERSIST_INTERVAL 秒把内存快照写 CSV 并保存恢复点。"""
|
||||||
while not self._stop_event.wait(PERSIST_INTERVAL):
|
while not self._stop_event.wait(PERSIST_INTERVAL):
|
||||||
self._persist_all()
|
self._persist_all()
|
||||||
|
|
||||||
def _persist_all(self) -> None:
|
def _persist_all(self) -> None:
|
||||||
|
"""在锁内复制 data 后落盘,避免长时间持锁。"""
|
||||||
if not self._ready.is_set():
|
if not self._ready.is_set():
|
||||||
return
|
return
|
||||||
with self._lock:
|
with self._lock:
|
||||||
@@ -533,6 +571,7 @@ class DataProvider:
|
|||||||
end_ms: Optional[int],
|
end_ms: Optional[int],
|
||||||
limit: Optional[int],
|
limit: Optional[int],
|
||||||
) -> List[Dict[str, float]]:
|
) -> List[Dict[str, float]]:
|
||||||
|
"""从内存读取已缓存的基础周期 K 线并按时间/limit 裁剪。"""
|
||||||
with self._lock:
|
with self._lock:
|
||||||
candles = list(self.data.get(symbol, {}).get(timeframe, []))
|
candles = list(self.data.get(symbol, {}).get(timeframe, []))
|
||||||
if start_ms is not None:
|
if start_ms is not None:
|
||||||
@@ -551,6 +590,7 @@ class DataProvider:
|
|||||||
end_time: Optional[object] = None,
|
end_time: Optional[object] = None,
|
||||||
limit: Optional[int] = None,
|
limit: Optional[int] = None,
|
||||||
) -> List[Dict[str, float]]:
|
) -> List[Dict[str, float]]:
|
||||||
|
"""对外查询:基础周期直接返回;衍生周期从 derived_map 取 base,resample 后对齐时间戳再裁剪。"""
|
||||||
if symbol not in self.symbols:
|
if symbol not in self.symbols:
|
||||||
raise HTTPException(status_code=404, detail=f"symbol {symbol} 不可用")
|
raise HTTPException(status_code=404, detail=f"symbol {symbol} 不可用")
|
||||||
self.wait_ready()
|
self.wait_ready()
|
||||||
@@ -565,6 +605,7 @@ class DataProvider:
|
|||||||
if target_minutes is None:
|
if target_minutes is None:
|
||||||
raise HTTPException(status_code=400, detail=f"不支持的时间周期: {timeframe}")
|
raise HTTPException(status_code=400, detail=f"不支持的时间周期: {timeframe}")
|
||||||
target_ms = target_minutes * 60_000
|
target_ms = target_minutes * 60_000
|
||||||
|
# 起点前移一根目标周期长度,保证首根合成 K 边界完整
|
||||||
adjusted_start = None if start_ms is None else max(0, start_ms - target_ms)
|
adjusted_start = None if start_ms is None else max(0, start_ms - target_ms)
|
||||||
base_candles = self._get_base_klines(symbol, base_tf, adjusted_start, end_ms, None)
|
base_candles = self._get_base_klines(symbol, base_tf, adjusted_start, end_ms, None)
|
||||||
if not base_candles:
|
if not base_candles:
|
||||||
@@ -574,9 +615,11 @@ class DataProvider:
|
|||||||
return []
|
return []
|
||||||
df = df.drop_duplicates(subset=["timestamp"], keep="last").sort_values("timestamp")
|
df = df.drop_duplicates(subset=["timestamp"], keep="last").sort_values("timestamp")
|
||||||
df["date"] = pd.to_datetime(df["timestamp"], unit="ms", utc=True)
|
df["date"] = pd.to_datetime(df["timestamp"], unit="ms", utc=True)
|
||||||
|
# resample_to_interval 按「分钟」目标周期聚合 OHLCV
|
||||||
resampled = self.resample_df(df, target_minutes)
|
resampled = self.resample_df(df, target_minutes)
|
||||||
if resampled is None or resampled.empty:
|
if resampled is None or resampled.empty:
|
||||||
return []
|
return []
|
||||||
|
# 统一得到毫秒 timestamp 列(resample 可能返回 date 或 DatetimeIndex)
|
||||||
if "timestamp" in resampled.columns:
|
if "timestamp" in resampled.columns:
|
||||||
resampled_df = resampled.copy()
|
resampled_df = resampled.copy()
|
||||||
else:
|
else:
|
||||||
@@ -621,6 +664,8 @@ class DataProvider:
|
|||||||
|
|
||||||
|
|
||||||
def create_app(provider: DataProvider) -> FastAPI:
|
def create_app(provider: DataProvider) -> FastAPI:
|
||||||
|
"""构造 FastAPI 应用:lifespan 内同步 initialize 并启动后台拉数。"""
|
||||||
|
|
||||||
@asynccontextmanager
|
@asynccontextmanager
|
||||||
async def lifespan(app: FastAPI):
|
async def lifespan(app: FastAPI):
|
||||||
loop = asyncio.get_running_loop()
|
loop = asyncio.get_running_loop()
|
||||||
@@ -643,6 +688,7 @@ def create_app(provider: DataProvider) -> FastAPI:
|
|||||||
|
|
||||||
@app.get("/health")
|
@app.get("/health")
|
||||||
async def health() -> Dict[str, object]:
|
async def health() -> Dict[str, object]:
|
||||||
|
"""存活检查:交易所、交易对、基础/衍生周期、是否已完成冷启动。"""
|
||||||
return {
|
return {
|
||||||
"status": "ok",
|
"status": "ok",
|
||||||
"exchange": provider.exchange_name,
|
"exchange": provider.exchange_name,
|
||||||
@@ -655,6 +701,7 @@ def create_app(provider: DataProvider) -> FastAPI:
|
|||||||
|
|
||||||
@app.get("/timeframes")
|
@app.get("/timeframes")
|
||||||
async def list_timeframes() -> Dict[str, List[str]]:
|
async def list_timeframes() -> Dict[str, List[str]]:
|
||||||
|
"""返回配置的基础周期与可合成的衍生周期列表。"""
|
||||||
provider.wait_ready()
|
provider.wait_ready()
|
||||||
return {
|
return {
|
||||||
"base_timeframes": provider.timeframes,
|
"base_timeframes": provider.timeframes,
|
||||||
@@ -670,11 +717,13 @@ def create_app(provider: DataProvider) -> FastAPI:
|
|||||||
end: Optional[int] = Query(None, description="结束时间戳(ms)"),
|
end: Optional[int] = Query(None, description="结束时间戳(ms)"),
|
||||||
limit: Optional[int] = Query(None, description="可选,限制返回数量"),
|
limit: Optional[int] = Query(None, description="可选,限制返回数量"),
|
||||||
):
|
):
|
||||||
|
"""按交易对与时间周期返回 OHLCV;tf 支持配置的基础周期及衍生合成周期。"""
|
||||||
data = provider.get_klines(symbol=symbol, timeframe=tf, start_time=start, end_time=end, limit=limit)
|
data = provider.get_klines(symbol=symbol, timeframe=tf, start_time=start, end_time=end, limit=limit)
|
||||||
return data
|
return data
|
||||||
|
|
||||||
@app.get("/")
|
@app.get("/")
|
||||||
async def root() -> Dict[str, object]:
|
async def root() -> Dict[str, object]:
|
||||||
|
"""根路径:服务名、交易所、交易对与可用周期(含 ready 标志)。"""
|
||||||
return {
|
return {
|
||||||
"service": "Data Provider",
|
"service": "Data Provider",
|
||||||
"exchange": provider.exchange_name,
|
"exchange": provider.exchange_name,
|
||||||
@@ -689,6 +738,7 @@ def create_app(provider: DataProvider) -> FastAPI:
|
|||||||
|
|
||||||
|
|
||||||
def build_app() -> FastAPI:
|
def build_app() -> FastAPI:
|
||||||
|
"""默认入口:从环境变量 CONFIG_PATH(或 config.json)加载配置并创建 FastAPI app。"""
|
||||||
config_path = Path(os.getenv("CONFIG_PATH", "config.json"))
|
config_path = Path(os.getenv("CONFIG_PATH", "config.json"))
|
||||||
provider = DataProvider(config_path)
|
provider = DataProvider(config_path)
|
||||||
return create_app(provider)
|
return create_app(provider)
|
||||||
@@ -698,6 +748,7 @@ app = build_app()
|
|||||||
|
|
||||||
|
|
||||||
def main() -> None:
|
def main() -> None:
|
||||||
|
"""直接运行本模块时启动 uvicorn(监听 UVICORN_HOST / UVICORN_PORT)。"""
|
||||||
host = os.getenv("UVICORN_HOST", "0.0.0.0")
|
host = os.getenv("UVICORN_HOST", "0.0.0.0")
|
||||||
port = int(os.getenv("UVICORN_PORT", "9009"))
|
port = int(os.getenv("UVICORN_PORT", "9009"))
|
||||||
|
|
||||||
|
|||||||
@@ -1,54 +1,83 @@
|
|||||||
|
### 线段生成(get_seg_list,按代码实现)
|
||||||
|
|
||||||
|
输入为已生成的笔列表 bi_list,输出为线段列表 seg_list。线段 ChanSEG 由若干笔 ChanBI 顺序组成。缠论里用「特征序列」刻画线段划分;本实现中 ChanSBI 即特征序列的元素(由同向笔经包含处理合并而成的一段区间,链成特征序列),在反向笔方向上维护 down_sbi_list / up_sbi_list 即维护该方向上的特征序列。当出现「非包含的新笔」、特征序列上可构成三元素结构时,对前一个 ChanSBI 做分型判断,决定是否结束当前线段并生成反向线段。循环末尾会维护 last_up_bi、last_down_bi 及 up_bi_list、down_bi_list;最后调用 cal_bi_zs(seg_list) 做线段与笔中枢的后续计算。函数末尾大段被注释掉的「线段破坏重划」逻辑当前未执行,以下不描述。
|
||||||
|
|
||||||
|
1. 首根线段的产生(seg_list 为空时)
|
||||||
|
|
||||||
|
仅当当前笔 bi 满足 bi.check_overlap() 为真时才允许生成第一段线段;否则该笔不会开启线段。check_overlap 要求:bi 之后至少还有两笔(next、next.next),且 next.next 已确认(is_sure)。在此前提下:
|
||||||
|
- 若 bi 为向上笔:须满足 bi.high > bi.next.low 且 bi.high < bi.next.next.high(当前向上笔的高点落在后两笔的重叠/交错区间内);
|
||||||
|
- 若 bi 为向下笔:须满足 bi.high > bi.next.high 且 bi.low > bi.next.next.low。
|
||||||
|
满足后:向上笔则新建方向为「上涨」的 ChanSEG,并以该笔初始化一条向上 ChanSBI;向下笔则新建「下跌」线段并初始化向下 ChanSBI。
|
||||||
|
|
||||||
|
2. 已有线段时的同向延伸(笔与线段方向相同)
|
||||||
|
|
||||||
|
- 当前线段为上涨、新笔仍为向上:若已有 last_up_sbi,则用 ChanSBI.check_bi_included 判断新笔是否被并入当前向上的特征序列末项;不并入则链上追加新的 ChanSBI。当前笔始终 last_seg.add_bi(bi) 并入该线段。
|
||||||
|
- 当前线段为下跌、新笔仍为向下:对称地维护 last_down_sbi 与向下特征序列,并 last_seg.add_bi(bi)。
|
||||||
|
|
||||||
|
3. 反向笔与特征序列(ChanSBI 链,为分型做准备)
|
||||||
|
|
||||||
|
- 当前线段为上涨、新笔为向下:在向下笔序列上维护 down_sbi_list(向下特征序列)。若链上已有不止一个向下 ChanSBI,且新向下笔不能 check_bi_included 入上一档 last_down_sbi,则先为当前笔新建 ChanSBI 并接到链上,再对「原来的」last_down_sbi 调用 check_fx()(此时该节点已有 pre 与 next,特征序列上相邻三元素可验分型)。
|
||||||
|
- 若向下链上仅有一个 ChanSBI:同样先判包含;若不包含则追加第二个节点;若尚未形成多节点,则用当前笔直接新建 ChanSBI 加入链。
|
||||||
|
- 当前线段为下跌、新笔为向上:在 up_sbi_list 上对称处理,对 last_up_sbi 调用 check_fx()。
|
||||||
|
|
||||||
|
4. ChanSBI.check_bi_included(特征序列元素的笔包含)
|
||||||
|
|
||||||
|
若当前特征序列项区间「左包含」新笔(实现上:当前 high、low 与新笔 high、low 满足代码中的包含关系,且与 pre 组合满足扩展条件),则将新笔并入当前 ChanSBI,并按方向更新极值(向下特征序列取更低价为 low,向上特征序列取更高价为 high),返回已包含;否则返回不包含,由上层新建下一档 ChanSBI。
|
||||||
|
|
||||||
|
5. ChanSBI.check_fx 与 has_fx_gap(特征序列上的分型与缺口)
|
||||||
|
|
||||||
|
当某 ChanSBI 同时存在 pre、next 且已 set_end_bi 时(特征序列上相邻三元素齐备):
|
||||||
|
- 顶分型 TOP:中间元素高点高于前、后特征序列元素的高点;若中间低点高于前一段的高点,则 has_fx_gap 为真(顶分型处出现缺口类关系)。
|
||||||
|
- 底分型 BOTTOM:中间元素低点低于前、后特征序列元素的低点;若中间高点低于前一段的低点,则 has_fx_gap 为真。
|
||||||
|
|
||||||
|
6. 出现分型后如何结束当前线段并开新线段(核心分支)
|
||||||
|
|
||||||
|
在上涨线段末端、向下笔链上形成顶分型 TOP 时(对 last_down_sbi.check_fx() 得到 TOP):
|
||||||
|
- 若此前处于 look_for_top 状态:对 seg_list 中倒数第二条线段调用 set_sure(bi),并清除 look_for_top。
|
||||||
|
- 若该顶分型 has_fx_gap:置 look_for_bottom;对当前线段 pre_set_end_bi(结束笔取到分型前一笔等,按代码索引);以 last_down_sbi.start_bi 为起点新建下跌线段;向上特征序列(up_sbi_list)重置为仅含 last_up_bi 的新链。
|
||||||
|
- 若无缺口:若当前处于 look_for_bottom,则做「回补」式调整(如 set_start_bi、倒数第二条 set_end_bi、重置 up_sbi 等)并 last_seg.add_bi(bi);否则对当前线段 set_end_bi,再新建下跌线段,并同样重置向上特征序列链。
|
||||||
|
|
||||||
|
在下跌线段末端、向上笔链上形成底分型 BOTTOM 时:与上对称,使用 look_for_bottom / look_for_top、last_up_sbi.has_fx_gap,并新开上涨线段、重置向下特征序列(down_sbi_list)。
|
||||||
|
|
||||||
|
上述分支中凡未单独说明的,仍会按路径将当前笔 add_bi 到相应线段。
|
||||||
|
|
||||||
|
7. 小结
|
||||||
|
|
||||||
|
线段方向由首段 check_overlap 与后续「反向笔方向上的特征序列分型」共同决定;同向笔持续并入当前线段;反向笔先经包含处理叠成特征序列(ChanSBI 链),在三元素结构成立时识别顶/底分型,再结合是否缺口(has_fx_gap)与 look_for_top / look_for_bottom 状态机切换线段,并标记确认(set_sure)。
|
||||||
|
|
||||||
|
|
||||||
### 缠论中枢识别核心逻辑(按代码实现)
|
### 缠论中枢识别核心逻辑(按代码实现)
|
||||||
|
|
||||||
1. **基础中枢生成规则**
|
依据缠论线段中枢定义计算中枢:`get_zs_list` 从第 4 根线段开始(线段列表索引 3),每 3 根线段为一组检查;上涨中枢为后中枢 zd > 前中枢 zg(不重叠上移),下跌中枢为后中枢 zg < 前中枢 zd(不重叠下移),盘整 / 扩张为后中枢与前中枢的整体区间(GG / DD)存在交集;中枢在初成三段之后,可按两段一组继续并入线段,扩展为 5 根、7 根……
|
||||||
- 从第4根线段起,每连续3段已确认的线段为一组;方向严格交替:上涨中枢 down→up→down,下跌中枢 up→down→up;
|
|
||||||
- 3段线段极值须有重叠(min(三高点) > max(三低点)),否则跳过本组;
|
1. 基础中枢生成规则
|
||||||
|
- 从第4根线段起(索引3),每连续3段已确认的线段为一组;
|
||||||
|
- 3段线段极值须有重叠(min(三高点) > max(三低点),即 zg > zd),否则跳过本组;
|
||||||
- 中枢核心区间(只读、扩展时不变):
|
- 中枢核心区间(只读、扩展时不变):
|
||||||
- zg = 初始3段高点的最小值(重叠区间上沿);
|
- zg = 初始3段高点的最小值(重叠区间上沿);
|
||||||
- zd = 初始3段低点的最大值(重叠区间下沿);
|
- zd = 初始3段低点的最大值(重叠区间下沿);
|
||||||
- 中枢极值(扩展时可更新):
|
- 中枢极值(扩展时可更新):
|
||||||
- GG = 参与中枢的所有线段的高点最大值;
|
- GG = 参与中枢的所有线段的高点最大值;
|
||||||
- DD = 参与中枢的所有线段的低点最小值。
|
- DD = 参与中枢的所有线段的低点最小值。
|
||||||
|
- 相对前一中枢的区间关系(用于区分走势类型):
|
||||||
|
- 上涨中枢:后中枢 zd > 前中枢 zg(与前一中枢 [zd,zg] 不重叠、整体上移,即「不重叠上移」);
|
||||||
|
- 下跌中枢:后中枢 zg < 前中枢 zd(与前一中枢不重叠、整体下移,即「不重叠下移」);
|
||||||
|
- 盘整 / 扩张:后中枢与前中枢在整体区间(GG / DD)上存在交集(与上述纯阶梯式上移、下移相区别)。
|
||||||
|
- 线段方向:在判定为上涨 / 下跌中枢时,须满足与中枢类型对应的严格交替——上涨中枢 down→up→down,下跌中枢 up→down→up(首段中枢仅按前三段形态定方向)。
|
||||||
|
|
||||||
2. **中枢扩展规则**
|
2. 中枢扩展规则
|
||||||
- 若下一组3段的区间 [zd,zg] 与当前中枢区间有重叠(zg≥当前zd 且 zd≤当前zg),则进行扩展,不新建中枢;
|
- 当形成三段中枢后,以中枢之后已完成的线段每两段为一组(与本中枢离开、回抽相关的成对线段)检查是否满足扩展条件;满足时可继续并入,使中枢覆盖 5 根、7 根……线段;
|
||||||
- 若本组3段整体与前中枢不重叠(例如已离开后新中枢区间在上/下方),仍逐根检查:只要某线段与当前中枢 [zd,zg] 有重叠(例如离开后的回抽段又回到前中枢内),则将该线段并入前中枢扩展,不因此直接新建中枢;从第一根不重叠的线段起再考虑新中枢;
|
- 若下一组两段的第二段与当前中枢区间有重叠(第二段高点≥当前zd且第二段低点≤当前zg),则进行扩展,不新建中枢;
|
||||||
- 扩展时:先将本组3段并入(或逐根并入见上),再向后逐根线段判断;线段与当前中枢 [zd,zg] 有重叠(线段区间与 [zd,zg] 相交)则并入,遇首根不重叠则停止扩展;
|
- 若下一组两段的第二段与当前中枢区间没有重叠,那么这两段都不并入当前中枢,当前中枢结束;从该两段组的第二段开始作为后续扫描起点,按中枢离开判断规则和新中枢生成规则继续;
|
||||||
- 线段区间用该线段的起止笔极值:seg_high = max(start_bi.high, end_bi.high),seg_low = min(start_bi.low, end_bi.low);
|
|
||||||
- 扩展后裁剪:中枢参与线段的结束段方向须与形成时一致——上涨中枢结束于 down 段,下跌中枢结束于 up 段;从末尾向前 pop 至满足为止;
|
|
||||||
- 扩展时仅更新 GG、DD 和 seg_list;zg、zd 永不修改;
|
- 扩展时仅更新 GG、DD 和 seg_list;zg、zd 永不修改;
|
||||||
- 下一组起点:从「裁剪后保留的最后一段」的下一条线段开始,被裁掉的线段参与下一组3根,避免漏入错误中枢。
|
|
||||||
|
|
||||||
3. **中枢离开判定规则**
|
3. 新中枢生成规则
|
||||||
- 向上离开:up 方向完成线段的 low > 原中枢 zg;强离开标注(非判定条件):up 线段 high > 原中枢 GG;
|
- 若下一组3段与当前中枢区间不重叠,且本组内逐根检查后没有任何线段与前中枢 [zd,zg] 重叠(即没有“离开后回抽回到前中枢”),并满足上涨中枢或下跌中枢的区间关系(后中枢 zd > 前中枢 zg,或后中枢 zg < 前中枢 zd)及对应线段方向;
|
||||||
- 向下离开:down 方向完成线段的 high < 原中枢 zd;强离开标注(非判定条件):down 线段 low < 原中枢 DD;
|
- 新中枢类型按区间关系判定:后中枢 zd > 前中枢 zg → 上涨中枢(down→up→down);后中枢 zg < 前中枢 zd → 下跌中枢(up→down→up);
|
||||||
- 有效离开(触发第三类买卖点):
|
|
||||||
- 向上有效离开:up 离开后,down 方向回抽线段的 low > 原中枢 zg 不进原中枢,触发第三类买点;
|
|
||||||
- 向下有效离开:down 离开后,up 方向回抽线段的 high < 原中枢 zd 不进原中枢,触发第三类卖点;
|
|
||||||
- 无效离开:回抽线段进入原中枢 [zd, zg] 区间,不触发第三类买卖点。
|
|
||||||
|
|
||||||
4. **新中枢生成规则**
|
|
||||||
- 若下一组3段与当前中枢区间不重叠,且本组内逐根检查后没有任何线段与前中枢 [zd,zg] 重叠(即没有“离开后回抽回到前中枢”),则当前中枢确认结束,本组3段生成新中枢;
|
|
||||||
- 新中枢类型与区间:后中枢 zg > 前中枢 zg → 上涨中枢(down→up→down);否则为下跌中枢(up→down→up);
|
|
||||||
- 新中枢的 zg/zd/GG/DD 仅基于自身初始3段计算;zg、zd 之后不再修改。
|
- 新中枢的 zg/zd/GG/DD 仅基于自身初始3段计算;zg、zd 之后不再修改。
|
||||||
|
|
||||||
|
|
||||||
5. **中枢趋势判定规则**
|
|
||||||
- 上涨中枢:后中枢 zg > 前中枢 zg(即后中枢 zd > 前中枢 zg 的等价表述);
|
|
||||||
- 下跌中枢:后中枢 zg < 前中枢 zg;
|
|
||||||
- 盘整/扩张:后中枢与前中枢区间重叠(后 zd < 前 zg 且 后 zg > 前 zd)。
|
|
||||||
|
|
||||||
6. **大级别中枢(扩张)**
|
|
||||||
- 多个笔/线段中枢若区间两两重叠,可合并为一个大级别中枢;
|
|
||||||
- 重叠定义:两中枢 [zd,zg] 有交集,即 (zs_i.zg >= zs_j.zd and zs_i.zd <= zs_j.zg);
|
|
||||||
- 大级别中枢的 zd/zg 取子中枢的并集(包住所有子中枢),用于显示更大级别震荡区间。
|
|
||||||
线段高低点的判断
|
线段高低点的判断
|
||||||
注意,这里必须提醒一句,就是这在以前也曾说过,就是,如果线段中,最高或最低点不是线段的端点,那么,在任何以线段为基础的分析中,例如把线段为基础构成最小级别的中枢等,都可以把该线段标准化为最高低点都在端点。因为, 在以线段为基础的分析中,都把线段当成一个没有内部 结构的基本部件,所以,只需要关心这线段的实际区间就可以,这样就可以只看其高低点。
|
注意,这里必须提醒一句,就是这在以前也曾说过,就是,如果线段中,最高或最低点不是线段的端点,那么,在任何以线段为基础的分析中,例如把线段为基础构成最小级别的中枢等,都可以把该线段标准化为最高低点都在端点。因为, 在以线段为基础的分析中,都把线段当成一个没有内部 结构的基本部件,所以,只需要关心这线段的实际区间就可以,这样就可以只看其高低点。
|
||||||
经过标准化处理后,所有向上线段都是以最低点开始最高点结束,向下线段都是以最高点开始最低点结束,这样,所以线段的连接,就形成一条延续不断、首尾相连的折线,这样,复杂的图形,就会十分地标准化,也为后面的中枢、走势类型等分析提供了最标准且基础的部件。
|
经过标准化处理后,所有向上线段都是以最低点开始最高点结束,向下线段都是以最高点开始最低点结束,这样,所以线段的连接,就形成一条延续不断、首尾相连的折线,这样,复杂的图形,就会十分地标准化,也为后面的中枢、走势类型等分析提供了最标准且基础的部件。
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
概率引擎
|
|
||||||
|
|
||||||
|
|
||||||
策略引擎
|
|
||||||
分批建仓
|
|
||||||
Reference in New Issue
Block a user