perf: 默认启用高吞吐 DHT 抓取策略

This commit is contained in:
chuan
2026-08-10 10:18:03 +08:00
parent 07e44c401b
commit bed4458a79
7 changed files with 199 additions and 48 deletions
+10 -5
View File
@@ -245,13 +245,17 @@
- [x] 增加 `find_node` Peer Lookup 新目标和 Metadata 建连的显式配置
- [x] 为主动 `find_node` `get_peers``sample_infohashes` 增加共享 UDP 查询总预算
- [x] 为 Metadata TCP 建连增加独立每秒速率限制
- [x] 将桌面默认配置调整为保守网络预算
- [x] 根据本机首次验证将主动 UDP 从 `40/s` 下调至 `10/s` 并将 Metadata 建连从 `5/s` 下调至 `2/s`
- [x] 对照 Bitmagnet 默认并发将运行模板调整为受全局预算约束的高吞吐配置
- [x] 使用保守网络预算定位早期本机断网问题
- [x] 根据本机首次验证将主动 UDP 从 `40/s` 下调至 `10/s` 并将 Metadata 建连从 `5/s` 下调至 `2/s` 完成故障隔离
- [x] 对照 Bitmagnet 默认并发建立受全局预算和有界队列保护的激进配置
- [x] 根据两轮一分钟资源测试将 Bitmagnet 等效激进配置设为应用和运行模板默认值
- [x] 修复 Windows 临时索引文件占用导致整个服务退出的问题
- [x] 验证极保守配置运行三分钟不影响同机代理网络并安全退出
- [x] 将 BEP-51 采样准入压力反向传递到采样查询调度
- [x] 实现样本来源节点单点 `get_peers` 优先和失败后有限递归降级
- [x] 将 BEP-51 最大在途请求和采样失败后的迭代回退暴露为应用配置
- [x] 使用 Bitmagnet 等效并发完成一分钟资源测试并确认主网卡无丢包无错误且持久化无积压
- [ ] 评估将新发现但尚未验证的节点直接作为 BEP-51 候选并避免与 `find_node` 重复探测
- [x] 完成首轮三分钟对比并验证 Peer Lookup UDP 从 `278` 降至 `254` 且网络稳定
- [ ] 通过多轮或更长时间运行评估随机 DHT 样本下的 Metadata 成功率
- [ ] 统计按需验证的 Peer 发现率握手成功率和平均验证耗时
@@ -271,6 +275,7 @@
- [ ] 增加数据库备份检查点和恢复验证
- [ ] 增加日志轮转和保留策略
- [x] 验证间歇运行和正常退出恢复
- [x] 完成本机约七小时真实持续运行并确认采集索引和搜索服务可用
- [ ] 验证二十四小时和七天连续运行
- [ ] 根据规模决定是否继续使用 RocksDB
@@ -317,6 +322,6 @@
- [x] 将规模基准拆分为参数数据集工作负载采样报告和编排模块
- [x] 明确单元组件集成端到端和性能测试层级并增加公开 API 集成测试
使用真实采集数据进行一小时持续运行并记录内存队列磁盘网络和验证指标
对比 Bitmagnet 新节点直接采样策略并测量有效样本数每次成功 Metadata 的 UDP 和 TCP 成本
随后扩展为二十四小时持续运行并根据数据决定资源参数和正则查询优化优先级
随后完成二十四小时持续运行并根据数据决定资源参数和正则查询优化优先级
+118 -2
View File
@@ -56,6 +56,7 @@ struct LookupState {
outstanding: usize,
deadline: Instant,
preferred_phase: bool,
allow_iterative_fallback: bool,
completion: Option<oneshot::Sender<PeerLookupResult>>,
}
@@ -90,6 +91,7 @@ struct LookupResponse {
pub(crate) struct PeerLookupRequest {
pub(crate) info_hash: [u8; 20],
pub(crate) preferred_node: Option<NodeTuple>,
pub(crate) allow_iterative_fallback: bool,
pub(crate) completion: Option<oneshot::Sender<PeerLookupResult>>,
}
@@ -98,6 +100,7 @@ impl PeerLookupRequest {
Self {
info_hash,
preferred_node: None,
allow_iterative_fallback: true,
completion: None,
}
}
@@ -131,6 +134,7 @@ impl PeerLookupHandle {
.send(PeerLookupRequest {
info_hash,
preferred_node: None,
allow_iterative_fallback: true,
completion: Some(completion),
})
.await
@@ -364,6 +368,7 @@ impl PeerLookupActor {
outstanding: 0,
deadline: now + LOOKUP_TIMEOUT,
preferred_phase: preferred.is_some(),
allow_iterative_fallback: request.allow_iterative_fallback,
completion: request.completion,
},
);
@@ -418,6 +423,10 @@ impl PeerLookupActor {
self.outbound_query_budget.refund_one();
self.runtime_stats.peer_lookup_send_failed();
if preferred_phase && let Some(state) = self.active.get_mut(&lookup_id) {
if !state.allow_iterative_fallback {
self.finish_lookup(lookup_id);
return;
}
state.preferred_phase = false;
state.deadline = now + LOOKUP_TIMEOUT;
self.runtime_stats.peer_lookup_fallback();
@@ -485,6 +494,7 @@ impl PeerLookupActor {
let hash = state.info_hash_hex.clone();
let lookup_id = pending.lookup_id;
let preferred_succeeded = pending.preferred_phase && !discovered.is_empty();
let preferred_without_fallback = pending.preferred_phase && !state.allow_iterative_fallback;
let _ = state;
for peer in discovered {
@@ -509,6 +519,10 @@ impl PeerLookupActor {
self.finish_lookup(lookup_id);
return;
}
if preferred_without_fallback {
self.finish_lookup(lookup_id);
return;
}
if pending.preferred_phase {
self.runtime_stats.peer_lookup_fallback();
}
@@ -517,6 +531,7 @@ impl PeerLookupActor {
async fn expire(&mut self, now: Instant) {
let mut affected = AHashSet::new();
let mut finish_without_fallback = AHashSet::new();
while let Some((deadline, key)) = self.pending_expiry.front().copied() {
if deadline > now {
break;
@@ -534,13 +549,21 @@ impl PeerLookupActor {
state.outstanding = state.outstanding.saturating_sub(1);
if pending.preferred_phase {
state.preferred_phase = false;
state.deadline = now + LOOKUP_TIMEOUT;
self.runtime_stats.peer_lookup_fallback();
if state.allow_iterative_fallback {
state.deadline = now + LOOKUP_TIMEOUT;
self.runtime_stats.peer_lookup_fallback();
} else {
finish_without_fallback.insert(pending.lookup_id);
}
}
affected.insert(pending.lookup_id);
}
self.runtime_stats.peer_lookup_timeout();
}
for lookup_id in finish_without_fallback {
affected.remove(&lookup_id);
self.finish_lookup(lookup_id);
}
for lookup_id in affected {
self.dispatch_more(lookup_id, now).await;
}
@@ -642,6 +665,7 @@ mod tests {
id: [9; 20],
addr: remote_addr,
}),
allow_iterative_fallback: true,
completion: Some(completion),
})
.await
@@ -738,6 +762,7 @@ mod tests {
id: [9; 20],
addr: preferred_addr,
}),
allow_iterative_fallback: true,
completion: None,
})
.await
@@ -780,4 +805,95 @@ mod tests {
assert_eq!(snapshot.peer_lookup_fallbacks, 1);
shutdown.cancel();
}
#[tokio::test]
async fn preferred_node_without_peers_can_finish_without_iterative_lookup() {
let local_socket = Arc::new(UdpSocket::bind("127.0.0.1:0").await.unwrap());
let preferred_socket = UdpSocket::bind("127.0.0.1:0").await.unwrap();
let fallback_socket = UdpSocket::bind("127.0.0.1:0").await.unwrap();
let local_addr = local_socket.local_addr().unwrap();
let preferred_addr = preferred_socket.local_addr().unwrap();
let fallback_addr = fallback_socket.local_addr().unwrap();
let sockets = std::collections::HashMap::from([(local_addr, local_socket)]);
let snapshot = Arc::new(ArcSwap::from_pointee(RoutingSnapshot::from_nodes(
vec![NodeTuple {
id: [8; 20],
addr: fallback_addr,
}],
1,
)));
let (hash_tx, _hash_rx) = mpsc::channel(4);
let stats = DhtRuntimeStats::default();
let shutdown = CancellationToken::new();
let handle = spawn_peer_lookup(
NetMode::Ipv4Only,
[7; 20],
&sockets,
snapshot,
hash_tx,
PeerLookupRuntime {
options: PeerLookupOptions::default(),
stats: stats.clone(),
outbound_query_budget: SharedRateBudget::per_second(10_000, 10_000, true),
shutdown: shutdown.clone(),
},
);
let (completion, result_rx) = oneshot::channel();
handle
.request_tx
.send(PeerLookupRequest {
info_hash: [3; 20],
preferred_node: Some(NodeTuple {
id: [9; 20],
addr: preferred_addr,
}),
allow_iterative_fallback: false,
completion: Some(completion),
})
.await
.unwrap();
let mut buffer = [0u8; 512];
let (len, _) = tokio::time::timeout(
Duration::from_secs(1),
preferred_socket.recv_from(&mut buffer),
)
.await
.unwrap()
.unwrap();
let query: DhtMessage = serde_bencode::from_bytes(&buffer[..len]).unwrap();
let tid: TransactionId = query.t.as_ref().try_into().unwrap();
handle.route_response(
preferred_addr,
tid,
DhtResponse {
id: Some(serde_bytes::ByteBuf::from(vec![9; 20])),
nodes: None,
nodes6: None,
values: None,
samples: None,
num: None,
interval: None,
},
);
let result = tokio::time::timeout(Duration::from_secs(1), result_rx)
.await
.unwrap()
.unwrap();
assert!(result.peers.is_empty());
assert_eq!(result.queries, 1);
assert!(
tokio::time::timeout(
Duration::from_millis(200),
fallback_socket.recv_from(&mut buffer),
)
.await
.is_err()
);
let snapshot = stats.snapshot();
assert_eq!(snapshot.peer_lookup_queries, 1);
assert_eq!(snapshot.peer_lookup_fallbacks, 0);
shutdown.cancel();
}
}
+5
View File
@@ -130,6 +130,7 @@ pub(crate) fn spawn_sample_infohashes(
admission_rx,
hash_admission,
peer_lookup.request_sender(),
options.fallback_to_iterative,
stats.clone(),
shutdown.clone(),
);
@@ -200,6 +201,7 @@ fn spawn_sample_admission(
mut receiver: mpsc::Receiver<SampleAdmissionBatch>,
hash_admission: Arc<ArcSwapOption<SampleHashAdmissionCallback>>,
peer_lookup_tx: mpsc::Sender<PeerLookupRequest>,
fallback_to_iterative: bool,
runtime_stats: DhtRuntimeStats,
shutdown: CancellationToken,
) {
@@ -222,6 +224,7 @@ fn spawn_sample_admission(
let request = PeerLookupRequest {
info_hash,
preferred_node: Some(batch.preferred_node),
allow_iterative_fallback: fallback_to_iterative,
completion: None,
};
let sent = tokio::select! {
@@ -520,6 +523,7 @@ mod tests {
batch_rx,
admission,
lookup_tx,
false,
stats.clone(),
shutdown.clone(),
);
@@ -539,6 +543,7 @@ mod tests {
.unwrap()
.unwrap();
assert_eq!(request.info_hash, [2; 20]);
assert!(!request.allow_iterative_fallback);
let snapshot = stats.snapshot();
assert_eq!(snapshot.sample_infohashes_hashes_filtered, 1);
assert_eq!(snapshot.sample_infohashes_hashes_discovered, 1);
+4
View File
@@ -179,6 +179,9 @@ pub struct SampleInfohashesOptions {
pub burst: u32,
/// Maximum outstanding BEP-51 requests.
pub max_in_flight: usize,
/// Whether a sampled Hash should fall back to iterative get_peers after the sampling node
/// returns no Peer.
pub fallback_to_iterative: bool,
/// Per-request timeout in milliseconds.
pub request_timeout_millis: u64,
/// Retry delay for timeouts or nodes that do not return samples, in seconds.
@@ -333,6 +336,7 @@ impl Default for SampleInfohashesOptions {
max_queries_per_second: 1,
burst: 1,
max_in_flight: 4,
fallback_to_iterative: true,
request_timeout_millis: 1_500,
unsupported_backoff_secs: 300,
dedup_capacity: 1_000_000,
+12 -10
View File
@@ -17,18 +17,20 @@ max_path_depth = 64
port = 12313
netmode = "ipv4-only"
hash_queue_capacity = 20000
max_outbound_queries_per_second = 50
outbound_query_burst = 10
max_outbound_queries_per_second = 1000
outbound_query_burst = 200
metadata_timeout_secs = 6
metadata_queue_capacity = 20000
metadata_workers = 64
metadata_connects_per_second = 16
sample_queries_per_second = 10
peer_lookups_per_second = 8
peer_lookup_max_active = 32
find_node_queries_per_second = 30
find_node_max_in_flight = 64
new_destinations_per_minute = 1200
metadata_workers = 400
metadata_connects_per_second = 400
sample_queries_per_second = 60
sample_max_in_flight = 100
sample_fallback_to_iterative = false
peer_lookups_per_second = 200
peer_lookup_max_active = 200
find_node_queries_per_second = 10
find_node_max_in_flight = 100
new_destinations_per_minute = 12000
[http]
listen = "127.0.0.1:8080"
+13 -11
View File
@@ -38,20 +38,22 @@ cargo run -p dht-search -- --data-dir D:\data\dht-search --run-duration-secs 360
### 当前运行模板网络配置
`dht-search.example.toml` 使用低于 Bitmagnet 默认并发的高吞吐配置 主动 DHT 查询仍共享 `max_outbound_queries_per_second` 总预算 因此 `find_node` `get_peers``sample_infohashes` 的总发送速率不会各自叠加后失控
`dht-search.example.toml` 默认使用经过本机一分钟资源测试的 Bitmagnet 等效激进配置 主动 DHT 查询仍共享 `max_outbound_queries_per_second` 总预算 所有队列保持有界 可以按设备和网络条件主动下调
| 配置项 | 运行模板值 | 作用 |
|---|---:|---|
| `max_outbound_queries_per_second` | `50` | 三类主动 DHT UDP 查询的合计每秒速率 |
| `outbound_query_burst` | `10` | 空闲后允许立即消费的 UDP 查询数 |
| `find_node_queries_per_second` | `30` | `find_node` 自身速率上限 |
| `find_node_max_in_flight` | `64` | 同时等待响应的 `find_node` 数量 |
| `new_destinations_per_minute` | `1200` | 每分钟首次探测的新 UDP 目标数量 |
| `peer_lookups_per_second` | `8` | 每秒启动的 infohash Peer 查找数量 |
| `peer_lookup_max_active` | `32` | 同时运行的 Peer 查找数量 |
| `sample_queries_per_second` | `10` | BEP-51 采样查询速率 |
| `metadata_workers` | `64` | 同时处理的 Metadata 任务数量 |
| `metadata_connects_per_second` | `16` | 每秒真正开始的 Peer TCP 连接数量 |
| `max_outbound_queries_per_second` | `1000` | 三类主动 DHT UDP 查询的合计每秒速率硬上限 |
| `outbound_query_burst` | `200` | 空闲后允许立即消费的 UDP 查询数 |
| `find_node_queries_per_second` | `10` | `find_node` 自身速率上限 |
| `find_node_max_in_flight` | `100` | 同时等待响应的 `find_node` 数量 |
| `new_destinations_per_minute` | `12000` | 每分钟首次探测的新 UDP 目标数量 |
| `peer_lookups_per_second` | `200` | 每秒启动的 infohash Peer 查找数量 |
| `peer_lookup_max_active` | `200` | 同时运行的 Peer 查找数量 |
| `sample_queries_per_second` | `60` | BEP-51 采样查询速率 |
| `sample_max_in_flight` | `100` | 同时等待响应的 BEP-51 采样请求数量 |
| `sample_fallback_to_iterative` | `false` | 采样节点没有返回 Peer 时快速结束以扩大覆盖面 |
| `metadata_workers` | `400` | 同时处理的 Metadata 任务数量 |
| `metadata_connects_per_second` | `400` | 每秒真正开始的 Peer TCP 连接硬上限 |
Metadata 下载和可用性握手共用 `metadata_connects_per_second` 预算不会各自叠加
+37 -20
View File
@@ -51,6 +51,8 @@ pub(crate) struct DhtConfig {
pub(crate) metadata_workers: usize,
pub(crate) metadata_connects_per_second: u32,
pub(crate) sample_queries_per_second: u32,
pub(crate) sample_max_in_flight: usize,
pub(crate) sample_fallback_to_iterative: bool,
pub(crate) peer_lookups_per_second: u32,
pub(crate) peer_lookup_max_active: usize,
pub(crate) find_node_queries_per_second: u32,
@@ -151,6 +153,9 @@ impl AppConfig {
},
sample_infohashes: SampleInfohashesOptions {
max_queries_per_second: self.dht.sample_queries_per_second,
burst: self.dht.sample_queries_per_second.max(1),
max_in_flight: self.dht.sample_max_in_flight,
fallback_to_iterative: self.dht.sample_fallback_to_iterative,
..defaults.sample_infohashes
},
crawl: CrawlOptions {
@@ -235,6 +240,7 @@ impl AppConfig {
if self.dht.max_outbound_queries_per_second == 0
|| self.dht.outbound_query_burst == 0
|| self.dht.metadata_connects_per_second == 0
|| self.dht.sample_max_in_flight == 0
|| self.dht.peer_lookup_max_active == 0
|| self.dht.find_node_max_in_flight == 0
{
@@ -248,10 +254,10 @@ impl Default for AppConfig {
fn default() -> Self {
Self {
data_dir: PathBuf::from("data"),
persistence_queue_capacity: 4_096,
persistence_queue_capacity: 8_192,
stats_interval_secs: 10,
run_duration_secs: None,
index_batch_size: 512,
index_batch_size: 1_024,
index_interval_millis: 5_000,
metadata_limits: MetadataLimitsConfig::default(),
dht: DhtConfig::default(),
@@ -278,19 +284,21 @@ impl Default for DhtConfig {
Self {
port: 12_313,
netmode: NetworkMode::Ipv4Only,
hash_queue_capacity: 10_000,
max_outbound_queries_per_second: 10,
outbound_query_burst: 2,
metadata_timeout_secs: 4,
metadata_queue_capacity: 10_000,
metadata_workers: 8,
metadata_connects_per_second: 2,
sample_queries_per_second: 1,
peer_lookups_per_second: 1,
peer_lookup_max_active: 4,
find_node_queries_per_second: 6,
find_node_max_in_flight: 12,
new_destinations_per_minute: 60,
hash_queue_capacity: 20_000,
max_outbound_queries_per_second: 1_000,
outbound_query_burst: 200,
metadata_timeout_secs: 6,
metadata_queue_capacity: 20_000,
metadata_workers: 400,
metadata_connects_per_second: 400,
sample_queries_per_second: 60,
sample_max_in_flight: 100,
sample_fallback_to_iterative: false,
peer_lookups_per_second: 200,
peer_lookup_max_active: 200,
find_node_queries_per_second: 10,
find_node_max_in_flight: 100,
new_destinations_per_minute: 12_000,
}
}
}
@@ -309,7 +317,7 @@ impl Default for VerificationConfig {
Self {
enabled: true,
queue_capacity: 10_000,
max_active: 2,
max_active: 8,
max_peer_attempts: 3,
lease_secs: 60,
poll_interval_millis: 250,
@@ -354,11 +362,20 @@ mod tests {
config.validate().unwrap();
let options = config.dht_options();
assert_eq!(options.port, 12_313);
assert_eq!(options.metadata.max_worker_count, 8);
assert_eq!(options.metadata.max_connects_per_second, 2);
assert_eq!(options.metadata.max_worker_count, 400);
assert_eq!(options.metadata.max_connects_per_second, 400);
assert_eq!(options.metadata.max_metadata_size_bytes, 10 * 1024 * 1024);
assert_eq!(options.max_outbound_queries_per_second, 10);
assert_eq!(options.crawl.rate_limit.max_find_node_rate_per_sec, 6);
assert_eq!(options.max_outbound_queries_per_second, 1_000);
assert_eq!(options.sample_infohashes.max_queries_per_second, 60);
assert_eq!(options.sample_infohashes.max_in_flight, 100);
assert!(!options.sample_infohashes.fallback_to_iterative);
assert_eq!(options.peer_lookup.max_lookups_per_second, 200);
assert_eq!(options.peer_lookup.max_active_lookups, 200);
assert_eq!(options.crawl.rate_limit.max_find_node_rate_per_sec, 10);
assert_eq!(
options.crawl.rate_limit.max_new_destinations_per_minute,
12_000
);
}
#[test]