Revert "以运行覆盖层隔离补充行情"

This reverts commit 757b5665ca.
This commit is contained in:
boris
2026-09-07 11:28:33 +08:00
parent 757b5665ca
commit 728ed7998d
+41 -308
View File
@@ -340,12 +340,9 @@ pub struct ExecutionQuoteIterator<'a> {
heap: BinaryHeap<Reverse<(NaiveDateTime, usize, usize)>>, heap: BinaryHeap<Reverse<(NaiveDateTime, usize, usize)>>,
} }
type ExecutionQuotesBySymbol = HashMap<String, Vec<IntradayExecutionQuote>>;
type ExecutionQuotesByDate = HashMap<NaiveDate, ExecutionQuotesBySymbol>;
impl<'a> ExecutionQuoteIterator<'a> { impl<'a> ExecutionQuoteIterator<'a> {
fn new( fn new(
rows_by_symbol: Option<&'a ExecutionQuotesBySymbol>, rows_by_symbol: Option<&'a HashMap<String, Vec<IntradayExecutionQuote>>>,
symbols: Option<&BTreeSet<String>>, symbols: Option<&BTreeSet<String>>,
) -> Self { ) -> Self {
let mut streams = rows_by_symbol let mut streams = rows_by_symbol
@@ -368,34 +365,6 @@ impl<'a> ExecutionQuoteIterator<'a> {
} }
Self { streams, heap } Self { streams, heap }
} }
fn new_with_overlay(
base_rows_by_symbol: Option<&'a ExecutionQuotesBySymbol>,
overlay_rows_by_symbol: &'a ExecutionQuotesBySymbol,
symbols: Option<&BTreeSet<String>>,
) -> Self {
let mut streams = base_rows_by_symbol
.into_iter()
.flat_map(|rows_by_symbol| rows_by_symbol.iter())
.filter(|(symbol, _)| !overlay_rows_by_symbol.contains_key(*symbol))
.chain(overlay_rows_by_symbol.iter())
.filter(|(symbol, _)| {
symbols
.map(|allowed_symbols| allowed_symbols.contains(*symbol))
.unwrap_or(true)
})
.map(|(symbol, rows)| (symbol.as_str(), rows.as_slice()))
.collect::<Vec<_>>();
streams.sort_by_key(|(symbol, _)| *symbol);
let mut heap = BinaryHeap::with_capacity(streams.len());
for (stream_index, (_, rows)) in streams.iter().enumerate() {
if let Some(first) = rows.first() {
heap.push(Reverse((first.timestamp, stream_index, 0)));
}
}
Self { streams, heap }
}
} }
impl<'a> Iterator for ExecutionQuoteIterator<'a> { impl<'a> Iterator for ExecutionQuoteIterator<'a> {
@@ -1345,8 +1314,7 @@ pub struct DataSet {
candidate_symbol_ids_by_date: Arc<BTreeMap<NaiveDate, Vec<u32>>>, candidate_symbol_ids_by_date: Arc<BTreeMap<NaiveDate, Vec<u32>>>,
candidate_row_positions_by_date: Arc<Option<DenseRowPositionIndex>>, candidate_row_positions_by_date: Arc<Option<DenseRowPositionIndex>>,
corporate_actions_by_date: Arc<BTreeMap<NaiveDate, Vec<CorporateAction>>>, corporate_actions_by_date: Arc<BTreeMap<NaiveDate, Vec<CorporateAction>>>,
execution_quotes_by_date: Arc<ExecutionQuotesByDate>, execution_quotes_by_date: Arc<HashMap<NaiveDate, HashMap<String, Vec<IntradayExecutionQuote>>>>,
execution_quote_overlays_by_date: Arc<ExecutionQuotesByDate>,
execution_quote_dates: Arc<Vec<NaiveDate>>, execution_quote_dates: Arc<Vec<NaiveDate>>,
order_book_depth_index: Arc<HashMap<(NaiveDate, String), Vec<IntradayOrderBookDepthLevel>>>, order_book_depth_index: Arc<HashMap<(NaiveDate, String), Vec<IntradayOrderBookDepthLevel>>>,
benchmark_by_date: Arc<BTreeMap<NaiveDate, BenchmarkSnapshot>>, benchmark_by_date: Arc<BTreeMap<NaiveDate, BenchmarkSnapshot>>,
@@ -1882,7 +1850,6 @@ impl DataSet {
candidate_row_positions_by_date: Arc::new(candidate_row_positions_by_date), candidate_row_positions_by_date: Arc::new(candidate_row_positions_by_date),
corporate_actions_by_date: Arc::new(corporate_actions_by_date), corporate_actions_by_date: Arc::new(corporate_actions_by_date),
execution_quotes_by_date: Arc::new(execution_quotes_by_date), execution_quotes_by_date: Arc::new(execution_quotes_by_date),
execution_quote_overlays_by_date: Arc::new(HashMap::new()),
execution_quote_dates: Arc::new(execution_quote_dates), execution_quote_dates: Arc::new(execution_quote_dates),
order_book_depth_index: Arc::new(order_book_depth_index), order_book_depth_index: Arc::new(order_book_depth_index),
benchmark_by_date: Arc::new(benchmark_by_date), benchmark_by_date: Arc::new(benchmark_by_date),
@@ -2213,46 +2180,23 @@ impl DataSet {
} }
pub fn execution_quotes_on(&self, date: NaiveDate, symbol: &str) -> &[IntradayExecutionQuote] { pub fn execution_quotes_on(&self, date: NaiveDate, symbol: &str) -> &[IntradayExecutionQuote] {
if self.execution_quote_overlays_by_date.is_empty() { self.execution_quotes_by_date
return self
.execution_quotes_by_date
.get(&date)
.and_then(|rows_by_symbol| rows_by_symbol.get(symbol))
.map(Vec::as_slice)
.unwrap_or(&[]);
}
self.execution_quote_overlays_by_date
.get(&date) .get(&date)
.and_then(|rows_by_symbol| rows_by_symbol.get(symbol)) .and_then(|rows_by_symbol| rows_by_symbol.get(symbol))
.or_else(|| {
self.execution_quotes_by_date
.get(&date)
.and_then(|rows_by_symbol| rows_by_symbol.get(symbol))
})
.map(Vec::as_slice) .map(Vec::as_slice)
.unwrap_or(&[]) .unwrap_or(&[])
} }
pub fn has_execution_quotes_on_date(&self, date: NaiveDate) -> bool { pub fn has_execution_quotes_on_date(&self, date: NaiveDate) -> bool {
if self.execution_quote_overlays_by_date.is_empty() { self.execution_quotes_by_date
return self
.execution_quotes_by_date
.get(&date)
.is_some_and(|rows_by_symbol| !rows_by_symbol.is_empty());
}
self.execution_quote_overlays_by_date
.get(&date) .get(&date)
.is_some_and(|rows_by_symbol| !rows_by_symbol.is_empty()) .map(|rows_by_symbol| !rows_by_symbol.is_empty())
|| self .unwrap_or(false)
.execution_quotes_by_date
.get(&date)
.is_some_and(|rows_by_symbol| !rows_by_symbol.is_empty())
} }
pub fn execution_quote_key_set(&self) -> HashSet<(NaiveDate, String)> { pub fn execution_quote_key_set(&self) -> HashSet<(NaiveDate, String)> {
self.execution_quotes_by_date self.execution_quotes_by_date
.iter() .iter()
.chain(self.execution_quote_overlays_by_date.iter())
.flat_map(|(date, rows_by_symbol)| { .flat_map(|(date, rows_by_symbol)| {
rows_by_symbol rows_by_symbol
.keys() .keys()
@@ -2262,51 +2206,15 @@ impl DataSet {
} }
pub fn execution_quote_count(&self) -> usize { pub fn execution_quote_count(&self) -> usize {
let mut count = self self.execution_quotes_by_date
.execution_quotes_by_date
.values() .values()
.flat_map(|rows_by_symbol| rows_by_symbol.values()) .flat_map(|rows_by_symbol| rows_by_symbol.values())
.map(Vec::len) .map(Vec::len)
.sum::<usize>(); .sum()
for (date, rows_by_symbol) in self.execution_quote_overlays_by_date.iter() {
for (symbol, rows) in rows_by_symbol {
count = count
.saturating_sub(
self.execution_quotes_by_date
.get(date)
.and_then(|base| base.get(symbol))
.map(Vec::len)
.unwrap_or(0),
)
.saturating_add(rows.len());
}
}
count
}
fn execution_quote_count_on_date(&self, date: NaiveDate) -> usize {
let base = self.execution_quotes_by_date.get(&date);
let mut count = base
.into_iter()
.flat_map(|rows_by_symbol| rows_by_symbol.values())
.map(Vec::len)
.sum::<usize>();
if let Some(overlay) = self.execution_quote_overlays_by_date.get(&date) {
for (symbol, rows) in overlay {
count = count
.saturating_sub(
base.and_then(|rows_by_symbol| rows_by_symbol.get(symbol))
.map(Vec::len)
.unwrap_or(0),
)
.saturating_add(rows.len());
}
}
count
} }
pub fn add_execution_quotes(&mut self, quotes: Vec<IntradayExecutionQuote>) -> usize { pub fn add_execution_quotes(&mut self, quotes: Vec<IntradayExecutionQuote>) -> usize {
let mut grouped = ExecutionQuotesByDate::new(); let mut grouped = HashMap::<NaiveDate, HashMap<String, Vec<IntradayExecutionQuote>>>::new();
for quote in quotes { for quote in quotes {
grouped grouped
.entry(quote.date) .entry(quote.date)
@@ -2317,30 +2225,17 @@ impl DataSet {
} }
let mut added = 0usize; let mut added = 0usize;
let mut new_dates = Vec::new(); let mut new_dates = Vec::new();
let base_quotes = Arc::clone(&self.execution_quotes_by_date); let execution_quotes_by_date = Arc::make_mut(&mut self.execution_quotes_by_date);
let execution_quote_overlays_by_date =
Arc::make_mut(&mut self.execution_quote_overlays_by_date);
for (date, rows_by_symbol) in grouped { for (date, rows_by_symbol) in grouped {
let date_is_new = !base_quotes.contains_key(&date) let date_is_new = !execution_quotes_by_date.contains_key(&date);
&& !execution_quote_overlays_by_date.contains_key(&date); let target_by_symbol = execution_quotes_by_date.entry(date).or_default();
let target_by_symbol = execution_quote_overlays_by_date.entry(date).or_default();
if date_is_new { if date_is_new {
new_dates.push(date); new_dates.push(date);
} }
for (symbol, mut incoming) in rows_by_symbol { for (symbol, mut incoming) in rows_by_symbol {
incoming.sort_by_key(|quote| quote.timestamp); incoming.sort_by_key(|quote| quote.timestamp);
incoming.dedup_by(|left, right| left.timestamp == right.timestamp); incoming.dedup_by(|left, right| left.timestamp == right.timestamp);
let target = match target_by_symbol.entry(symbol) { let target = target_by_symbol.entry(symbol).or_default();
std::collections::hash_map::Entry::Occupied(entry) => entry.into_mut(),
std::collections::hash_map::Entry::Vacant(entry) => {
let base_rows = base_quotes
.get(&date)
.and_then(|rows_by_symbol| rows_by_symbol.get(entry.key()))
.cloned()
.unwrap_or_default();
entry.insert(base_rows)
}
};
if target.is_empty() { if target.is_empty() {
added = added.saturating_add(incoming.len()); added = added.saturating_add(incoming.len());
*target = incoming; *target = incoming;
@@ -2405,13 +2300,6 @@ impl DataSet {
date: NaiveDate, date: NaiveDate,
symbols: Option<&BTreeSet<String>>, symbols: Option<&BTreeSet<String>>,
) -> ExecutionQuoteIterator<'_> { ) -> ExecutionQuoteIterator<'_> {
if let Some(overlay) = self.execution_quote_overlays_by_date.get(&date) {
return ExecutionQuoteIterator::new_with_overlay(
self.execution_quotes_by_date.get(&date),
overlay,
symbols,
);
}
ExecutionQuoteIterator::new(self.execution_quotes_by_date.get(&date), symbols) ExecutionQuoteIterator::new(self.execution_quotes_by_date.get(&date), symbols)
} }
@@ -2426,33 +2314,29 @@ impl DataSet {
} }
pub fn remove_execution_quotes_on_date(&mut self, date: NaiveDate) -> usize { pub fn remove_execution_quotes_on_date(&mut self, date: NaiveDate) -> usize {
let row_count = self.execution_quote_count_on_date(date); let removed = Arc::make_mut(&mut self.execution_quotes_by_date).remove(&date);
Arc::make_mut(&mut self.execution_quote_overlays_by_date).remove(&date); let Some(rows_by_symbol) = removed else {
Arc::make_mut(&mut self.execution_quotes_by_date).remove(&date); return 0;
};
let dates = Arc::make_mut(&mut self.execution_quote_dates); let dates = Arc::make_mut(&mut self.execution_quote_dates);
if let Ok(index) = dates.binary_search(&date) { if let Ok(index) = dates.binary_search(&date) {
dates.remove(index); dates.remove(index);
} }
row_count rows_by_symbol.into_values().map(|rows| rows.len()).sum()
} }
pub fn release_execution_quotes_on_date(&mut self, date: NaiveDate) -> usize { pub fn release_execution_quotes_on_date(&mut self, date: NaiveDate) -> usize {
let row_count = self.execution_quote_count_on_date(date); let row_count = self
if self.execution_quote_overlays_by_date.contains_key(&date) { .execution_quotes_by_date
Arc::make_mut(&mut self.execution_quote_overlays_by_date).remove(&date); .get(&date)
} .map(|rows_by_symbol| rows_by_symbol.values().map(Vec::len).sum())
.unwrap_or(0);
// Run data shares this immutable map with the prepared-data cache. Arc::make_mut here // Run data shares this immutable map with the prepared-data cache. Arc::make_mut here
// would clone every date just to remove one entry and would not release the cached base. // would clone every date just to remove one entry and would not release the cached base.
if Arc::strong_count(&self.execution_quotes_by_date) == 1 { if row_count == 0 || Arc::strong_count(&self.execution_quotes_by_date) > 1 {
Arc::make_mut(&mut self.execution_quotes_by_date).remove(&date); return row_count;
} }
if !self.execution_quotes_by_date.contains_key(&date) self.remove_execution_quotes_on_date(date)
&& !self.execution_quote_overlays_by_date.contains_key(&date)
&& let Ok(index) = self.execution_quote_dates.binary_search(&date)
{
Arc::make_mut(&mut self.execution_quote_dates).remove(index);
}
row_count
} }
pub fn snapshot_components(&self) -> DataSetSnapshotComponents { pub fn snapshot_components(&self) -> DataSetSnapshotComponents {
@@ -2480,31 +2364,12 @@ impl DataSet {
.values() .values()
.flat_map(|rows| rows.iter().cloned()) .flat_map(|rows| rows.iter().cloned())
.collect::<Vec<_>>(); .collect::<Vec<_>>();
let execution_quotes = if self.execution_quote_overlays_by_date.is_empty() { let execution_quotes = self
self.execution_quotes_by_date .execution_quotes_by_date
.values() .values()
.flat_map(|rows_by_symbol| rows_by_symbol.values()) .flat_map(|rows_by_symbol| rows_by_symbol.values())
.flat_map(|rows| rows.iter().cloned()) .flat_map(|rows| rows.iter().cloned())
.collect::<Vec<_>>() .collect::<Vec<_>>();
} else {
let mut quotes = Vec::with_capacity(self.execution_quote_count());
for (date, rows_by_symbol) in self.execution_quotes_by_date.iter() {
let overlay = self.execution_quote_overlays_by_date.get(date);
for (symbol, rows) in rows_by_symbol {
if overlay.is_some_and(|overlay| overlay.contains_key(symbol)) {
continue;
}
quotes.extend(rows.iter().cloned());
}
}
quotes.extend(
self.execution_quote_overlays_by_date
.values()
.flat_map(|rows_by_symbol| rows_by_symbol.values())
.flat_map(|rows| rows.iter().cloned()),
);
quotes
};
DataSetSnapshotComponents { DataSetSnapshotComponents {
instruments, instruments,
@@ -2625,53 +2490,6 @@ impl DataSet {
if bar_count == 0 { if bar_count == 0 {
return Vec::new(); return Vec::new();
} }
if self.execution_quote_overlays_by_date.is_empty() {
return self.history_intraday_quotes_from_base(
date,
active_datetime,
symbol,
bar_count,
include_now,
);
}
let end = self
.execution_quote_dates
.partition_point(|quote_date| *quote_date <= date);
let mut quotes = Vec::with_capacity(bar_count);
'dates: for quote_date in self.execution_quote_dates[..end].iter().rev() {
let Some(rows) = self
.execution_quote_overlays_by_date
.get(quote_date)
.and_then(|rows_by_symbol| rows_by_symbol.get(symbol))
.or_else(|| {
self.execution_quotes_by_date
.get(quote_date)
.and_then(|rows_by_symbol| rows_by_symbol.get(symbol))
})
else {
continue;
};
for quote in rows.iter().rev() {
if intraday_quote_visible(quote, date, active_datetime, include_now) {
quotes.push(quote.clone());
if quotes.len() == bar_count {
break 'dates;
}
}
}
}
quotes.reverse();
quotes
}
fn history_intraday_quotes_from_base(
&self,
date: NaiveDate,
active_datetime: Option<NaiveDateTime>,
symbol: &str,
bar_count: usize,
include_now: bool,
) -> Vec<IntradayExecutionQuote> {
let end = self let end = self
.execution_quote_dates .execution_quote_dates
.partition_point(|quote_date| *quote_date <= date); .partition_point(|quote_date| *quote_date <= date);
@@ -3187,22 +3005,14 @@ impl DataSet {
.map(daily_market_price_bar) .map(daily_market_price_bar)
.collect(), .collect(),
Some("1m") => { Some("1m") => {
let mut bars = if self.execution_quote_overlays_by_date.is_empty() { let mut bars = self
self.execution_quotes_by_date .execution_quotes_by_date
.iter() .iter()
.filter(|(date, _)| **date >= start && **date <= end) .filter(|(date, _)| **date >= start && **date <= end)
.filter_map(|(_, rows_by_symbol)| rows_by_symbol.get(symbol)) .filter_map(|(_, rows_by_symbol)| rows_by_symbol.get(symbol))
.flat_map(|rows| rows.iter()) .flat_map(|rows| rows.iter())
.map(intraday_quote_price_bar) .map(intraday_quote_price_bar)
.collect::<Vec<_>>() .collect::<Vec<_>>();
} else {
self.execution_quote_dates
.iter()
.filter(|date| **date >= start && **date <= end)
.flat_map(|date| self.execution_quotes_on(*date, symbol).iter())
.map(intraday_quote_price_bar)
.collect::<Vec<_>>()
};
bars.sort_by(|left, right| { bars.sort_by(|left, right| {
left.date left.date
.cmp(&right.date) .cmp(&right.date)
@@ -5007,10 +4817,6 @@ mod tests {
&data.execution_quotes_by_date, &data.execution_quotes_by_date,
&run_data.execution_quotes_by_date &run_data.execution_quotes_by_date
)); ));
assert!(Arc::ptr_eq(
&data.execution_quote_overlays_by_date,
&run_data.execution_quote_overlays_by_date
));
assert!(Arc::ptr_eq( assert!(Arc::ptr_eq(
&data.execution_quote_dates, &data.execution_quote_dates,
&run_data.execution_quote_dates &run_data.execution_quote_dates
@@ -5033,14 +4839,10 @@ mod tests {
assert_eq!(data.execution_quote_count(), 0); assert_eq!(data.execution_quote_count(), 0);
assert_eq!(run_data.execution_quote_count(), 1); assert_eq!(run_data.execution_quote_count(), 1);
assert!(Arc::ptr_eq( assert!(!Arc::ptr_eq(
&data.execution_quotes_by_date, &data.execution_quotes_by_date,
&run_data.execution_quotes_by_date &run_data.execution_quotes_by_date
)); ));
assert!(!Arc::ptr_eq(
&data.execution_quote_overlays_by_date,
&run_data.execution_quote_overlays_by_date
));
assert!(!Arc::ptr_eq( assert!(!Arc::ptr_eq(
&data.execution_quote_dates, &data.execution_quote_dates,
&run_data.execution_quote_dates &run_data.execution_quote_dates
@@ -5917,75 +5719,6 @@ mod tests {
assert_eq!(run_data.execution_quote_count(), 0); assert_eq!(run_data.execution_quote_count(), 0);
} }
#[test]
fn execution_quote_overlay_preserves_base_precedence_and_merged_iteration() {
let date = NaiveDate::parse_from_str("2025-01-02", "%Y-%m-%d").unwrap();
let quote = |symbol: &str, time: &str, last_price: f64| IntradayExecutionQuote {
date,
timestamp: NaiveDateTime::parse_from_str(
&format!("2025-01-02 {time}"),
"%Y-%m-%d %H:%M:%S",
)
.unwrap(),
symbol: symbol.to_string(),
last_price,
bid1: last_price,
ask1: last_price,
bid1_volume: 10_000,
ask1_volume: 10_000,
volume_delta: 10_000,
amount_delta: last_price * 10_000.0,
trading_phase: Some("continuous".to_string()),
};
let data = DataSet::from_components_with_actions_and_quotes(
Vec::new(),
Vec::new(),
Vec::new(),
Vec::new(),
vec![benchmark_row("2025-01-02", 12.0)],
Vec::new(),
vec![
quote("000001.SZ", "09:30:00", 10.0),
quote("000001.SZ", "09:31:00", 10.1),
quote("000002.SZ", "09:30:00", 20.0),
],
)
.unwrap();
let mut run_data = data.clone();
assert_eq!(
run_data.add_execution_quotes(vec![
quote("000001.SZ", "09:31:00", 99.0),
quote("000001.SZ", "09:32:00", 10.2),
]),
1
);
let symbol_rows = run_data.execution_quotes_on(date, "000001.SZ");
assert_eq!(symbol_rows.len(), 3);
assert_eq!(symbol_rows[1].last_price, 10.1);
assert_eq!(data.execution_quotes_on(date, "000001.SZ").len(), 2);
let merged = run_data
.execution_quotes_iter_on_date_for_symbols(date, None)
.map(|row| (row.timestamp.time().to_string(), row.symbol.clone()))
.collect::<Vec<_>>();
assert_eq!(
merged,
vec![
("09:30:00".to_string(), "000001.SZ".to_string()),
("09:30:00".to_string(), "000002.SZ".to_string()),
("09:31:00".to_string(), "000001.SZ".to_string()),
("09:32:00".to_string(), "000001.SZ".to_string()),
]
);
assert_eq!(run_data.execution_quote_count(), 4);
assert_eq!(run_data.snapshot_components().execution_quotes.len(), 4);
assert_eq!(run_data.release_execution_quotes_on_date(date), 4);
assert_eq!(run_data.execution_quotes_on(date, "000001.SZ").len(), 2);
assert_eq!(run_data.execution_quote_count(), 3);
}
#[test] #[test]
fn shared_execution_quote_release_does_not_clone_the_base_map() { fn shared_execution_quote_release_does_not_clone_the_base_map() {
let date = NaiveDate::parse_from_str("2025-01-02", "%Y-%m-%d").unwrap(); let date = NaiveDate::parse_from_str("2025-01-02", "%Y-%m-%d").unwrap();