From 56a38accc867dffb92060a2d5852db724e7422e4 Mon Sep 17 00:00:00 2001 From: boris Date: Fri, 28 Aug 2026 17:02:35 +0800 Subject: [PATCH] =?UTF-8?q?=E4=B8=BA=E8=82=A1=E7=A5=A8=E5=BA=8F=E5=88=97?= =?UTF-8?q?=E5=A2=9E=E5=8A=A0=E6=9C=89=E7=95=8C=E4=BA=A4=E6=98=93=E6=97=A5?= =?UTF-8?q?=E4=BD=8D=E7=BD=AE=E7=B4=A2=E5=BC=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- crates/fidc-core/src/data.rs | 258 +++++++++++++++++++++++++++++++++-- 1 file changed, 245 insertions(+), 13 deletions(-) diff --git a/crates/fidc-core/src/data.rs b/crates/fidc-core/src/data.rs index 2995329..9ae5455 100644 --- a/crates/fidc-core/src/data.rs +++ b/crates/fidc-core/src/data.rs @@ -571,6 +571,16 @@ type DenseRowPositionIndex = BTreeMap>; const MISSING_ROW_POSITION: u32 = u32::MAX; const MAX_DENSE_ROW_INDEX_BYTES: usize = 256 * 1024 * 1024; +#[derive(Debug, Clone)] +struct SymbolSeriesEndPositions { + decision: Vec, + current: Vec, +} + +type SymbolSeriesEndPositionIndex = Vec>; + +const MAX_SERIES_END_POSITION_INDEX_BYTES: usize = 256 * 1024 * 1024; + #[derive(Debug, Clone)] struct AdjustedCloseSeries { dates: Vec, @@ -670,6 +680,14 @@ impl AdjustedCloseSeries { Err(0) => return [None; N], Err(index) => index, }; + self.moving_averages_at_end(end, lookbacks) + } + + fn moving_averages_at_end( + &self, + end: usize, + lookbacks: &[usize; N], + ) -> [Option; N] { std::array::from_fn(|index| self.moving_average_at_end(end, lookbacks[index])) } @@ -864,6 +882,15 @@ impl SymbolPriceSeries { return None; } let end = self.end_index(date)?; + self.moving_average_at_end(end, lookback, field) + } + + fn moving_average_at_end( + &self, + end: usize, + lookback: usize, + field: PriceField, + ) -> Option { if end < lookback { return None; } @@ -994,6 +1021,14 @@ impl SymbolPriceSeries { let Some(end) = self.rolling_end_index(date, include_now) else { return [None; N]; }; + self.volume_moving_averages_at_end(end, lookbacks) + } + + fn volume_moving_averages_at_end( + &self, + end: usize, + lookbacks: &[usize; N], + ) -> [Option; N] { std::array::from_fn(|index| { let lookback = lookbacks[index]; self.valid_volume_window(end, lookback).map(|(start, end)| { @@ -1292,6 +1327,7 @@ pub struct DataSet { adjusted_close_series_by_symbol: Arc>>, market_series_by_symbol_id: Arc>>>, adjusted_close_series_by_symbol_id: Arc>>>, + market_series_end_positions_by_symbol_id: Arc>, benchmark_series_cache: Arc, symbol_id_by_code: Arc>, eligible_universe_by_date: Arc>>>, @@ -1319,7 +1355,11 @@ impl DataSet { ) -> Self { let mut calendar_dates = self.calendar.days().to_vec(); calendar_dates.extend(dates); - self.calendar = Arc::new(TradingCalendar::new(calendar_dates)); + let calendar = Arc::new(TradingCalendar::new(calendar_dates)); + self.market_series_end_positions_by_symbol_id = Arc::new( + build_symbol_series_end_positions(&self.market_series_by_symbol_id, &calendar), + ); + self.calendar = calendar; self } @@ -1728,6 +1768,8 @@ impl DataSet { adjusted_close_series_by_symbol_id[symbol_id as usize] = Some(Arc::clone(series)); } } + let market_series_end_positions_by_symbol_id = + build_symbol_series_end_positions(&market_series_by_symbol_id, &calendar); let execution_quotes_by_date = build_execution_quote_index(execution_quotes); let mut execution_quote_dates = execution_quotes_by_date.keys().copied().collect::>(); execution_quote_dates.sort_unstable(); @@ -1759,6 +1801,9 @@ impl DataSet { adjusted_close_series_by_symbol: Arc::new(adjusted_close_series_by_symbol), market_series_by_symbol_id: Arc::new(market_series_by_symbol_id), adjusted_close_series_by_symbol_id: Arc::new(adjusted_close_series_by_symbol_id), + market_series_end_positions_by_symbol_id: Arc::new( + market_series_end_positions_by_symbol_id, + ), benchmark_series_cache: Arc::new(benchmark_series_cache), symbol_id_by_code: Arc::new(symbol_id_by_code), eligible_universe_by_date: Arc::new(OnceLock::new()), @@ -1852,6 +1897,27 @@ impl DataSet { .as_deref() } + fn market_series_end_index_by_symbol_id( + &self, + date: NaiveDate, + symbol_id: u32, + include_now: bool, + ) -> Option { + let calendar_index = self.calendar.index_of(date)?; + let positions = self + .market_series_end_positions_by_symbol_id + .as_ref() + .as_ref()? + .get(symbol_id as usize)? + .as_ref()?; + let end = if include_now { + positions.current.get(calendar_index) + } else { + positions.decision.get(calendar_index) + }?; + Some(*end as usize) + } + pub fn factor(&self, date: NaiveDate, symbol: &str) -> Option<&DailyFactorSnapshot> { let symbol_id = self.symbol_id(symbol)?; self.factor_by_symbol_id(date, symbol_id) @@ -1969,14 +2035,24 @@ impl DataSet { ) -> StandardRollingMeans { let close = if close_lookbacks.iter().any(|lookback| *lookback > 0) { self.adjusted_close_series_by_symbol_id(symbol_id) - .map(|series| series.moving_averages(date, close_lookbacks, include_now)) + .map(|series| { + self.market_series_end_index_by_symbol_id(date, symbol_id, include_now) + .map(|end| series.moving_averages_at_end(end, close_lookbacks)) + .unwrap_or_else(|| series.moving_averages(date, close_lookbacks, include_now)) + }) .unwrap_or([None; 7]) } else { [None; 7] }; let volume = if volume_lookbacks.iter().any(|lookback| *lookback > 0) { self.market_series_by_symbol_id(symbol_id) - .map(|series| series.volume_moving_averages(date, volume_lookbacks, include_now)) + .map(|series| { + self.market_series_end_index_by_symbol_id(date, symbol_id, include_now) + .map(|end| series.volume_moving_averages_at_end(end, volume_lookbacks)) + .unwrap_or_else(|| { + series.volume_moving_averages(date, volume_lookbacks, include_now) + }) + }) .unwrap_or([None; 5]) } else { [None; 5] @@ -3137,19 +3213,50 @@ impl DataSet { match field.as_ref() { "close" | "prev_close" | "stock_close" | "price" => self .adjusted_close_series_by_symbol_id(symbol_id) - .and_then(|series| series.decision_moving_average(date, lookback)), + .and_then(|series| { + self.market_series_end_index_by_symbol_id(date, symbol_id, false) + .map(|end| series.moving_average_at_end(end, lookback)) + .unwrap_or_else(|| series.decision_moving_average(date, lookback)) + }), "volume" | "stock_volume" => self .market_series_by_symbol_id(symbol_id) - .and_then(|series| series.decision_volume_moving_average(date, lookback)), + .and_then(|series| { + self.market_series_end_index_by_symbol_id(date, symbol_id, false) + .map(|end| { + series.valid_volume_window(end, lookback).map(|(start, end)| { + normalize_rolling_factor( + (series.valid_volume_sum_prefix[end] + - series.valid_volume_sum_prefix[start]) + / lookback as f64, + 12, + ) + }) + }) + .unwrap_or_else(|| series.decision_volume_moving_average(date, lookback)) + }), "day_open" | "dayopen" => self .market_series_by_symbol_id(symbol_id) - .and_then(|series| series.moving_average(date, lookback, PriceField::DayOpen)), + .and_then(|series| { + self.market_series_end_index_by_symbol_id(date, symbol_id, false) + .map(|end| series.moving_average_at_end(end, lookback, PriceField::DayOpen)) + .unwrap_or_else(|| { + series.moving_average(date, lookback, PriceField::DayOpen) + }) + }), "open" => self .market_series_by_symbol_id(symbol_id) - .and_then(|series| series.moving_average(date, lookback, PriceField::Open)), + .and_then(|series| { + self.market_series_end_index_by_symbol_id(date, symbol_id, false) + .map(|end| series.moving_average_at_end(end, lookback, PriceField::Open)) + .unwrap_or_else(|| series.moving_average(date, lookback, PriceField::Open)) + }), "last" | "last_price" => self .market_series_by_symbol_id(symbol_id) - .and_then(|series| series.moving_average(date, lookback, PriceField::Last)), + .and_then(|series| { + self.market_series_end_index_by_symbol_id(date, symbol_id, false) + .map(|end| series.moving_average_at_end(end, lookback, PriceField::Last)) + .unwrap_or_else(|| series.moving_average(date, lookback, PriceField::Last)) + }), other => self.factor_moving_average(date, symbol, other, lookback), } } @@ -3192,19 +3299,50 @@ impl DataSet { match field.as_ref() { "close" | "prev_close" | "stock_close" | "price" => self .adjusted_close_series_by_symbol_id(symbol_id) - .and_then(|series| series.current_moving_average(date, lookback)), + .and_then(|series| { + self.market_series_end_index_by_symbol_id(date, symbol_id, true) + .map(|end| series.moving_average_at_end(end, lookback)) + .unwrap_or_else(|| series.current_moving_average(date, lookback)) + }), "volume" | "stock_volume" => self .market_series_by_symbol_id(symbol_id) - .and_then(|series| series.current_volume_moving_average(date, lookback)), + .and_then(|series| { + self.market_series_end_index_by_symbol_id(date, symbol_id, true) + .map(|end| { + series.valid_volume_window(end, lookback).map(|(start, end)| { + normalize_rolling_factor( + (series.valid_volume_sum_prefix[end] + - series.valid_volume_sum_prefix[start]) + / lookback as f64, + 12, + ) + }) + }) + .unwrap_or_else(|| series.current_volume_moving_average(date, lookback)) + }), "day_open" | "dayopen" => self .market_series_by_symbol_id(symbol_id) - .and_then(|series| series.moving_average(date, lookback, PriceField::DayOpen)), + .and_then(|series| { + self.market_series_end_index_by_symbol_id(date, symbol_id, true) + .map(|end| series.moving_average_at_end(end, lookback, PriceField::DayOpen)) + .unwrap_or_else(|| { + series.moving_average(date, lookback, PriceField::DayOpen) + }) + }), "open" => self .market_series_by_symbol_id(symbol_id) - .and_then(|series| series.moving_average(date, lookback, PriceField::Open)), + .and_then(|series| { + self.market_series_end_index_by_symbol_id(date, symbol_id, true) + .map(|end| series.moving_average_at_end(end, lookback, PriceField::Open)) + .unwrap_or_else(|| series.moving_average(date, lookback, PriceField::Open)) + }), "last" | "last_price" => self .market_series_by_symbol_id(symbol_id) - .and_then(|series| series.moving_average(date, lookback, PriceField::Last)), + .and_then(|series| { + self.market_series_end_index_by_symbol_id(date, symbol_id, true) + .map(|end| series.moving_average_at_end(end, lookback, PriceField::Last)) + .unwrap_or_else(|| series.moving_average(date, lookback, PriceField::Last)) + }), other => self.factor_moving_average(date, symbol, other, lookback), } } @@ -3954,6 +4092,52 @@ fn build_dense_row_positions( Some(positions_by_date) } +fn build_symbol_series_end_positions( + series_by_symbol_id: &[Option>], + calendar: &TradingCalendar, +) -> Option { + let entries = series_by_symbol_id.len().checked_mul(calendar.len())?; + let bytes = entries + .checked_mul(2)? + .checked_mul(std::mem::size_of::())?; + if bytes > MAX_SERIES_END_POSITION_INDEX_BYTES + || series_by_symbol_id.iter().flatten().any(|series| { + series.dates.len() > u32::MAX as usize || calendar.len() > u32::MAX as usize + }) + { + return None; + } + + let calendar_days = calendar.days(); + let positions = series_by_symbol_id + .par_iter() + .map(|series| { + let series = series.as_deref()?; + let mut decision = Vec::with_capacity(calendar_days.len()); + let mut current = Vec::with_capacity(calendar_days.len()); + let mut series_index = 0usize; + for date in calendar_days { + while series + .dates + .get(series_index) + .is_some_and(|series_date| *series_date < *date) + { + series_index += 1; + } + decision.push(series_index as u32); + let current_index = if series.dates.get(series_index) == Some(date) { + series_index + 1 + } else { + series_index + }; + current.push(current_index as u32); + } + Some(SymbolSeriesEndPositions { decision, current }) + }) + .collect::>(); + Some(positions) +} + fn dense_row_position( positions_by_date: &Option, date: NaiveDate, @@ -5145,6 +5329,54 @@ mod tests { ); } + #[test] + fn series_end_position_index_preserves_decision_and_current_boundaries() { + let data = volume_contract_data(Some([1.0, 1.0, 1.0])); + let symbol_id = data.symbol_id("000001.SZ").expect("symbol id"); + let dates = data.calendar().days(); + + assert!(data + .market_series_end_positions_by_symbol_id + .as_ref() + .is_some()); + assert_eq!( + data.market_series_end_index_by_symbol_id(dates[0], symbol_id, false), + Some(0) + ); + assert_eq!( + data.market_series_end_index_by_symbol_id(dates[0], symbol_id, true), + Some(1) + ); + assert_eq!( + data.market_series_end_index_by_symbol_id(dates[2], symbol_id, false), + Some(2) + ); + assert_eq!( + data.market_series_end_index_by_symbol_id(dates[2], symbol_id, true), + Some(3) + ); + + let extended = data + .clone() + .with_additional_trading_dates([NaiveDate::from_ymd_opt(2025, 1, 7).unwrap()]); + assert_eq!( + extended.market_series_end_index_by_symbol_id( + NaiveDate::from_ymd_opt(2025, 1, 7).unwrap(), + symbol_id, + false, + ), + Some(3) + ); + assert_eq!( + extended.market_series_end_index_by_symbol_id( + NaiveDate::from_ymd_opt(2025, 1, 7).unwrap(), + symbol_id, + true, + ), + Some(3) + ); + } + #[test] fn source_volume_contract_rejects_windows_containing_missing_values() { let data = volume_contract_data(Some([1.0, 0.0, 1.0]));