mirror of
https://github.com/0xfnzero/solana-streamer.git
synced 2026-08-24 14:28:10 +00:00
optimize EventParserFactory::create_parser
This commit is contained in:
@@ -1,6 +1,6 @@
|
|||||||
use anyhow::{anyhow, Result};
|
use anyhow::{anyhow, Result};
|
||||||
use solana_sdk::pubkey::Pubkey;
|
use solana_sdk::pubkey::Pubkey;
|
||||||
use std::sync::Arc;
|
use std::{collections::HashMap, sync::{Arc, LazyLock}};
|
||||||
|
|
||||||
use crate::streaming::event_parser::protocols::{
|
use crate::streaming::event_parser::protocols::{
|
||||||
bonk::parser::BONK_PROGRAM_ID, pumpfun::parser::PUMPFUN_PROGRAM_ID,
|
bonk::parser::BONK_PROGRAM_ID, pumpfun::parser::PUMPFUN_PROGRAM_ID,
|
||||||
@@ -63,19 +63,27 @@ impl std::str::FromStr for Protocol {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
static EVENT_PARSERS: LazyLock<HashMap<Protocol, Arc<dyn EventParser>>> = LazyLock::new(|| {
|
||||||
|
let mut parsers: HashMap<Protocol, Arc<dyn EventParser>> = HashMap::new();
|
||||||
|
parsers.insert(Protocol::PumpSwap, Arc::new(PumpSwapEventParser::new()));
|
||||||
|
parsers.insert(Protocol::PumpFun, Arc::new(PumpFunEventParser::new()));
|
||||||
|
parsers.insert(Protocol::Bonk, Arc::new(BonkEventParser::new()));
|
||||||
|
parsers.insert(Protocol::RaydiumCpmm, Arc::new(RaydiumCpmmEventParser::new()));
|
||||||
|
parsers.insert(Protocol::RaydiumClmm, Arc::new(RaydiumClmmEventParser::new()));
|
||||||
|
parsers
|
||||||
|
});
|
||||||
|
|
||||||
|
|
||||||
/// 事件解析器工厂 - 用于创建不同协议的事件解析器
|
/// 事件解析器工厂 - 用于创建不同协议的事件解析器
|
||||||
pub struct EventParserFactory;
|
pub struct EventParserFactory;
|
||||||
|
|
||||||
impl EventParserFactory {
|
impl EventParserFactory {
|
||||||
|
|
||||||
/// 创建指定协议的事件解析器
|
/// 创建指定协议的事件解析器
|
||||||
pub fn create_parser(protocol: Protocol) -> Arc<dyn EventParser> {
|
pub fn create_parser(protocol: Protocol) -> Arc<dyn EventParser> {
|
||||||
match protocol {
|
EVENT_PARSERS.get(&protocol).cloned().unwrap_or_else(|| {
|
||||||
Protocol::PumpSwap => Arc::new(PumpSwapEventParser::new()),
|
panic!("Parser for protocol {} not found", protocol);
|
||||||
Protocol::PumpFun => Arc::new(PumpFunEventParser::new()),
|
})
|
||||||
Protocol::Bonk => Arc::new(BonkEventParser::new()),
|
|
||||||
Protocol::RaydiumCpmm => Arc::new(RaydiumCpmmEventParser::new()),
|
|
||||||
Protocol::RaydiumClmm => Arc::new(RaydiumClmmEventParser::new()),
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// 创建所有协议的事件解析器
|
/// 创建所有协议的事件解析器
|
||||||
|
|||||||
@@ -261,7 +261,6 @@ impl YellowstoneGrpc {
|
|||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
/// 处理事件交易
|
|
||||||
async fn process_event_transaction<F>(
|
async fn process_event_transaction<F>(
|
||||||
transaction_pretty: TransactionPretty,
|
transaction_pretty: TransactionPretty,
|
||||||
callback: &F,
|
callback: &F,
|
||||||
@@ -271,24 +270,39 @@ impl YellowstoneGrpc {
|
|||||||
where
|
where
|
||||||
F: Fn(Box<dyn UnifiedEvent>) + Send + Sync,
|
F: Fn(Box<dyn UnifiedEvent>) + Send + Sync,
|
||||||
{
|
{
|
||||||
|
// let start_time = std::time::Instant::now();
|
||||||
let slot = transaction_pretty.slot;
|
let slot = transaction_pretty.slot;
|
||||||
let signature = transaction_pretty.signature.to_string();
|
let signature = transaction_pretty.signature.to_string();
|
||||||
|
let mut futures = Vec::new();
|
||||||
for protocol in protocols {
|
for protocol in protocols {
|
||||||
let parser = EventParserFactory::create_parser(protocol);
|
let parser = EventParserFactory::create_parser(protocol);
|
||||||
let events = parser
|
let tx_clone = transaction_pretty.tx.clone();
|
||||||
.parse_transaction(
|
let signature_clone = signature.clone();
|
||||||
transaction_pretty.tx.clone(),
|
let bot_wallet_clone = bot_wallet.clone();
|
||||||
&signature,
|
|
||||||
|
futures.push(tokio::spawn(async move {
|
||||||
|
parser.parse_transaction(
|
||||||
|
tx_clone,
|
||||||
|
&signature_clone,
|
||||||
Some(slot),
|
Some(slot),
|
||||||
transaction_pretty.block_time,
|
transaction_pretty.block_time,
|
||||||
bot_wallet.clone(),
|
bot_wallet_clone,
|
||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.unwrap_or_else(|_e| vec![]);
|
.unwrap_or_else(|_e| vec![])
|
||||||
for event in events {
|
}));
|
||||||
callback(event);
|
}
|
||||||
|
|
||||||
|
let results = futures::future::join_all(futures).await;
|
||||||
|
for result in results {
|
||||||
|
if let Ok(events) = result {
|
||||||
|
for event in events {
|
||||||
|
callback(event);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
// let elapsed = start_time.elapsed();
|
||||||
|
// println!("处理交易耗时: {:?}", elapsed);
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user