> ## Documentation Index
> Fetch the complete documentation index at: https://www.helius.dev/docs/llms.txt
> Use this file to discover all available pages before exploring further.

# Đo độ trễ

> Tìm hiểu cách đo và phân tích chính xác độ trễ của các luồng gRPC bằng nhiều phương pháp kiểm thử.

<Warning>
  **TUYÊN BỐ MIỄN TRỪ TRÁCH NHIỆM QUAN TRỌNG: Chỉ dành cho trung tâm dữ liệu production**

  Các bài kiểm thử độ trễ này **chỉ được thiết kế cho môi trường trung tâm dữ liệu production**. KHÔNG chạy các bài kiểm thử này trên máy cục bộ hoặc kết nối Internet dân dụng. Băng thông cục bộ không thể xử lý các subscription Solana nặng và sẽ cho ra kết quả vô nghĩa, không phản ánh hiệu năng thực tế.
</Warning>

<Warning>
  **YÊU CẦU ĐỒNG VỊ TRÍ: Triển khai gần endpoint LaserStream của bạn**

  Để có số liệu đo độ trễ có ý nghĩa, bạn **phải** đặt hạ tầng kiểm thử trong cùng khu vực với endpoint LaserStream đã chọn. Khoảng cách mạng sẽ chi phối kết quả đo — việc kiểm thử từ một châu lục khác sẽ phản ánh độ trễ mạng chứ không phải hiệu năng của LaserStream.
</Warning>

***

## Tìm hiểu độ trễ trong các hệ thống blockchain phân tán

Khi làm việc với các dịch vụ truyền phát blockchain, việc đo độ trễ trở nên phức tạp vì các hệ thống phân tán không có đồng hồ chung. Không giống các hệ thống truyền thống, nơi bạn có thể đo thời gian khứ hồi đến một máy chủ duy nhất, mạng blockchain gồm nhiều trình xác thực, mỗi trình xác thực nhận và xử lý cùng một giao dịch tại các thời điểm khác nhau.

**Thách thức cơ bản:** Các blockchain như Solana không có khái niệm thời gian tuyệt đối. Mỗi nút xác thực trên toàn cầu sẽ nhận cùng một giao dịch tại các thời điểm khác nhau và việc xác nhận phụ thuộc vào tỷ lệ phần trăm của cụm đạt được đồng thuận. Điều này khiến việc đo độ trễ theo cách tất định trở nên bất khả thi theo nghĩa truyền thống.

## Các mức cam kết và mức độ ưu tiên về độ trễ

Solana cung cấp ba mức cam kết, mỗi mức có đặc điểm độ trễ khác nhau:

* **Đã xử lý**: Nhanh nhất, xác nhận bởi một trình xác thực (\~400ms)
* **Đã xác nhận**: Trung bình, xác nhận bởi siêu đa số (\~2-3 giây)
* **Đã hoàn tất**: Chậm nhất, hoàn tất trên toàn mạng (\~15-30 giây)

Đối với các ứng dụng nhạy cảm với độ trễ, **mức cam kết đã xử lý** thường là mục tiêu. Tất cả bài kiểm thử trong hướng dẫn này đều sử dụng mức cam kết đã xử lý vì hầu hết trường hợp sử dụng tần suất cao đều ưu tiên tốc độ hơn tính hoàn tất tuyệt đối.

## Ba phương pháp đo độ trễ

### 1. So sánh các luồng gRPC song song

**Phương pháp đáng tin cậy nhất** — So sánh hai luồng độc lập đến cùng một nguồn dữ liệu, đo xem luồng nào nhận được các sự kiện giống nhau trước.

**Ưu điểm:**

* Loại bỏ các vấn đề đồng bộ hóa đồng hồ
* Cung cấp phép so sánh hiệu năng tương đối
* Chính xác nhất khi so sánh các dịch vụ

### 2. So sánh dấu thời gian cục bộ với created\_at

**Độ tin cậy trung bình** — Đo chênh lệch giữa thời điểm hệ thống của bạn nhận được thông báo và dấu thời gian do dịch vụ LaserStream nhúng trong thông báo.

**Hạn chế:**

* Chỉ thể hiện thời điểm LaserStream tạo thông báo trong nội bộ
* Không ghi nhận được độ trễ ở thượng nguồn trước khi đến LaserStream
* Kém chính xác hơn Phương pháp 1 khi đo độ trễ đầu cuối thực sự

### 3. Phân tích dấu thời gian khối (không khuyến nghị)

**Không khuyến nghị** — So sánh thời gian nhận cục bộ với dấu thời gian khối của Solana.

**Hạn chế đáng kể:**

* Dấu thời gian khối chỉ có độ chi tiết đến giây
* Solana tạo khối sau mỗi 400ms
* Cung cấp rất ít thông tin hữu ích

***

## Yêu cầu thiết lập

### Đồng vị trí theo khu vực

Để có số liệu đo độ trễ có ý nghĩa, hãy triển khai hạ tầng kiểm thử trong cùng trung tâm dữ liệu hoặc khu vực với endpoint LaserStream của bạn.

**Các khu vực LaserStream hiện có:**

* **ewr**: New York, Hoa Kỳ (Bờ Đông) - `https://laserstream-mainnet-ewr.helius-rpc.com`
* **pitt**: Pittsburgh, Hoa Kỳ (Miền Trung) - `https://laserstream-mainnet-pitt.helius-rpc.com`
* **slc**: Salt Lake City, Hoa Kỳ (Bờ Tây) - `https://laserstream-mainnet-slc.helius-rpc.com`
* **ams**: Amsterdam, Châu Âu - `https://laserstream-mainnet-ams.helius-rpc.com`
* **fra**: Frankfurt, Châu Âu - `https://laserstream-mainnet-fra.helius-rpc.com`
* **tyo**: Tokyo, Châu Á - `https://laserstream-mainnet-tyo.helius-rpc.com`
* **sgp**: Singapore, Châu Á - `https://laserstream-mainnet-sgp.helius-rpc.com`

Để kiểm thử trên devnet, hãy sử dụng: `https://laserstream-devnet-ewr.helius-rpc.com`

Xem [tài liệu gRPC của LaserStream](/docs/vi/laserstream/grpc) để biết hướng dẫn thiết lập đầy đủ và nguyên tắc chọn endpoint.

### Thiết lập môi trường Rust

Tất cả script đo lường đều sử dụng Rust với Cargo. Thiết lập cơ bản:

```bash theme={"system"}
# Install Rust
curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh

# Setup project
cargo new latency-testing
cd latency-testing

# Add dependencies to Cargo.toml
[dependencies]
# Async runtime
tokio = { version = "1", features = ["full"] }
# Futures utilities for StreamExt, SinkExt
futures = "0.3"
# Environment variable loading
dotenvy = "0.15"
# Logging
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["fmt", "env-filter"] }
# Yellowstone gRPC client and protocol
yellowstone-grpc-client = "13.3.0"
yellowstone-grpc-proto = "12.6.0"
# For timestamp analysis
prost-types = "0.14"
```

Tạo tệp `.env` chứa thông tin xác thực của bạn:

```bash theme={"system"}
YS_GRPC_URL=your-comparison-endpoint
YS_API_KEY=your-comparison-api-key
LS_GRPC_URL=your-laserstream-endpoint  
LS_API_KEY=your-helius-api-key
```

Lấy khóa API Helius từ [Bảng điều khiển Helius](https://dashboard.helius.dev/). LaserStream devnet có sẵn trong tất cả các gói. Quyền truy cập mainnet yêu cầu gói Business hoặc Professional.

***

## Phương pháp 1: So sánh các luồng song song

Script này thiết lập hai kết nối độc lập đến các endpoint gRPC khác nhau và đo xem kết nối nào nhận được cùng một thông báo `BlockMeta` trước. Phương pháp này loại bỏ các vấn đề đồng bộ hóa đồng hồ bằng cách sử dụng thời gian tương đối.

```rust [expandable] theme={"system"}
use std::collections::HashMap;
use std::time::SystemTime;

use dotenvy::dotenv;
use futures::{StreamExt, SinkExt};
use tokio::sync::mpsc;
use tracing::{debug, error, info};

use yellowstone_grpc_client::{ClientTlsConfig, GeyserGrpcClient};
use yellowstone_grpc_proto::prelude::{
    subscribe_update::UpdateOneof, CommitmentLevel, SubscribeRequest, 
    SubscribeRequestFilterBlocksMeta, SubscribeUpdate,
};

#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
enum Source {
    Yellowstone,
    Laserstream,
}

#[derive(Default, Debug)]
struct SlotTimings {
    ys_recv_ms: Option<i128>,
    ls_recv_ms: Option<i128>,
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let _ = dotenv();
    tracing_subscriber::fmt().with_env_filter("info").init();

    let ys_url = std::env::var("YS_GRPC_URL").expect("YS_GRPC_URL env variable not set");
    let ls_url = std::env::var("LS_GRPC_URL").expect("LS_GRPC_URL env variable not set");
    let ys_api_key = std::env::var("YS_API_KEY").ok();
    let ls_api_key = std::env::var("LS_API_KEY").ok();

    let commitment = CommitmentLevel::Processed;
    info!(?commitment, "Starting latency comparison");

    // Establish both clients
    let mut ys_client = GeyserGrpcClient::build_from_shared(ys_url.clone())
        .expect("invalid YS url")
        .x_token(ys_api_key.clone())?
        .tls_config(ClientTlsConfig::new().with_native_roots())?
        .max_decoding_message_size(10 * 1024 * 1024)
        .connect()
        .await?;

    let mut ls_client = GeyserGrpcClient::build_from_shared(ls_url.clone())
        .expect("invalid LS url")
        .x_token(ls_api_key.clone())?
        .tls_config(ClientTlsConfig::new().with_native_roots())?
        .max_decoding_message_size(10 * 1024 * 1024)
        .connect()
        .await?;

    let (mut ys_tx, mut ys_rx) = ys_client.subscribe().await?;
    let (mut ls_tx, mut ls_rx) = ls_client.subscribe().await?;
    
    let subscribe_request = SubscribeRequest {
        blocks_meta: {
            let mut m = HashMap::<String, SubscribeRequestFilterBlocksMeta>::new();
            m.insert("all".to_string(), SubscribeRequestFilterBlocksMeta::default());
            m
        },
        commitment: Some(commitment as i32),
        ..Default::default()
    };

    ys_tx.send(subscribe_request.clone()).await?;
    ls_tx.send(subscribe_request).await?;

    let (agg_tx, mut agg_rx) = mpsc::unbounded_channel::<(Source, u64, i128)>();

    // Spawn task for Yellowstone stream
    {
        let agg_tx = agg_tx.clone();
        tokio::spawn(async move {
            while let Some(update_res) = ys_rx.next().await {
                match update_res {
                    Ok(update) => handle_update(Source::Yellowstone, update, &agg_tx).await,
                    Err(e) => {
                        error!(target: "ys", "stream error: {:?}", e);
                        break;
                    }
                }
            }
        });
    }

    // Spawn task for Laserstream stream
    {
        let agg_tx = agg_tx.clone();
        tokio::spawn(async move {
            while let Some(update_res) = ls_rx.next().await {
                match update_res {
                    Ok(update) => handle_update(Source::Laserstream, update, &agg_tx).await,
                    Err(e) => {
                        error!(target: "ls", "stream error: {:?}", e);
                        break;
                    }
                }
            }
        });
    }

    // Aggregator – collect latencies per slot and print once we have both sources
    let mut timings: HashMap<u64, SlotTimings> = HashMap::new();
    let mut deltas: Vec<i128> = Vec::new();
    let mut count = 0;

    println!("slot,ys_recv_ms,ls_recv_ms,delta_ms");

    while let Some((source, slot, latency_ms)) = agg_rx.recv().await {
        let entry = timings.entry(slot).or_default();
        match source {
            Source::Yellowstone => entry.ys_recv_ms = Some(latency_ms),
            Source::Laserstream => entry.ls_recv_ms = Some(latency_ms),
        }

        if let (Some(ys), Some(ls)) = (entry.ys_recv_ms, entry.ls_recv_ms) {
            let delta = ys - ls; // positive => YS arrived later
            println!("{slot},{ys},{ls},{delta}");
            
            deltas.push(delta);
            count += 1;
            
            if count % 100 == 0 {
                print_statistics(&deltas, count);
            }
            
            timings.remove(&slot);
        }
    }

    Ok(())
}

async fn handle_update(
    source: Source,
    update: SubscribeUpdate,
    agg_tx: &mpsc::UnboundedSender<(Source, u64, i128)>,
) {
    if let Some(UpdateOneof::BlockMeta(block_meta)) = update.update_oneof {
        let slot = block_meta.slot;
        let recv_ms = system_time_to_millis(SystemTime::now());
        debug!(?source, slot, recv_ms, "BlockMeta received");
        let _ = agg_tx.send((source, slot, recv_ms));
    }
}

fn print_statistics(deltas: &[i128], count: usize) {
    if deltas.is_empty() {
        return;
    }
    
    let mut sorted_deltas = deltas.to_vec();
    sorted_deltas.sort();
    
    let median = if sorted_deltas.len() % 2 == 0 {
        let mid = sorted_deltas.len() / 2;
        (sorted_deltas[mid - 1] + sorted_deltas[mid]) / 2
    } else {
        sorted_deltas[sorted_deltas.len() / 2]
    };
    
    let min = *sorted_deltas.first().unwrap();
    let max = *sorted_deltas.last().unwrap();
    let sum: i128 = sorted_deltas.iter().sum();
    let mean = sum / sorted_deltas.len() as i128;
    
    let p25_idx = (sorted_deltas.len() as f64 * 0.25) as usize;
    let p75_idx = (sorted_deltas.len() as f64 * 0.75) as usize;
    let p95_idx = (sorted_deltas.len() as f64 * 0.95) as usize;
    
    let p25 = sorted_deltas[p25_idx.min(sorted_deltas.len() - 1)];
    let p75 = sorted_deltas[p75_idx.min(sorted_deltas.len() - 1)];
    let p95 = sorted_deltas[p95_idx.min(sorted_deltas.len() - 1)];
    
    eprintln!("--- Statistics after {} slots ---", count);
    eprintln!("Delta (YS - LS) in milliseconds:");
    eprintln!("  Min: {}, Max: {}", min, max);
    eprintln!("  Mean: {}, Median: {}", mean, median);
    eprintln!("  P25: {}, P75: {}, P95: {}", p25, p75, p95);
    eprintln!("  Positive deltas (YS slower): {}/{} ({:.1}%)", 
              sorted_deltas.iter().filter(|&&x| x > 0).count(),
              sorted_deltas.len(),
              sorted_deltas.iter().filter(|&&x| x > 0).count() as f64 / sorted_deltas.len() as f64 * 100.0);
    eprintln!("---");
}

fn system_time_to_millis(st: SystemTime) -> i128 {
    st.duration_since(SystemTime::UNIX_EPOCH)
        .unwrap()
        .as_millis() as i128
}
```

**Nội dung được đo:** Chênh lệch hiệu năng tương đối giữa hai dịch vụ truyền phát. Giá trị delta cho biết dịch vụ nào phân phối cùng một thông tin slot trước.

**Các chỉ số chính:**

* **Delta dương**: Dịch vụ thứ nhất (YS) chậm hơn dịch vụ thứ hai (LS) — LaserStream nhanh hơn
* **Delta âm**: Dịch vụ thứ nhất (YS) nhanh hơn dịch vụ thứ hai (LS) — LaserStream chậm hơn
* **Trung bình/Trung vị**: Chênh lệch hiệu năng trung bình
* **P95**: Chênh lệch độ trễ tại phân vị thứ 95

**Chạy bài kiểm thử:**

```bash theme={"system"}
cargo run --bin latency-comparison
```

**Đầu ra mẫu:**

```
slot,ys_recv_ms,ls_recv_ms,delta_ms
352416939,1752168399141,1752168399140,1
352416940,1752168399526,1752168399512,14
352416941,1752168399890,1752168399877,13
```

Đầu ra hiển thị chênh lệch độ trễ theo thời gian thực và số liệu thống kê định kỳ. Delta trung bình dương cho biết dịch vụ thứ hai (LaserStream) luôn phân phối dữ liệu nhanh hơn.

***

## Phương pháp 2: Phân tích dấu thời gian tạo

Phương pháp này so sánh dấu thời gian `created_at` được nhúng trong thông báo với thời gian hệ thống cục bộ khi bạn nhận được thông báo.

```rust [expandable] theme={"system"}
use std::time::{Duration, SystemTime};
use dotenvy::dotenv;
use tracing::{debug, error, info};
use yellowstone_grpc_proto::prost_types::Timestamp;
use futures::StreamExt;
use futures::SinkExt;
use std::collections::HashMap;

use yellowstone_grpc_client::{ClientTlsConfig, GeyserGrpcClient};
use yellowstone_grpc_proto::prelude::{
    subscribe_update::UpdateOneof, CommitmentLevel, SubscribeRequest,
    SubscribeRequestFilterTransactions, SubscribeUpdate,
};

const ACCOUNTS_INCLUDE: &[&str] = &["BB5dnY55FXS1e1NXqZDwCzgdYJdMCj3B92PU6Q5Fb6DT"];
const COMMITMENT_LEVEL: &str = "processed";

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let _ = dotenv();
    tracing_subscriber::fmt().with_env_filter("info").init();

    let grpc_url = std::env::var("YS_GRPC_URL").expect("GRPC_URL env variable not set");
    let api_key = std::env::var("YS_API_KEY").ok();

    info!("Connecting to {} …", grpc_url);
    debug!(accounts_include = ?ACCOUNTS_INCLUDE, "Subscribing with accountsInclude filter");

    let mut client = GeyserGrpcClient::build_from_shared(grpc_url.clone())
        .expect("invalid URL")
        .x_token(api_key.clone())?
        .tls_config(ClientTlsConfig::new().with_native_roots())?
        .connect()
        .await?;

    let (mut subscribe_tx, mut subscribe_rx) = client.subscribe().await?;

    let mut tx_filter_map = HashMap::new();
    tx_filter_map.insert(
        "latency".to_string(),
        SubscribeRequestFilterTransactions {
            account_include: ACCOUNTS_INCLUDE.iter().map(|s| s.to_string()).collect(),
            vote: Some(false),
            failed: Some(false),
            ..Default::default()
        },
    );

    subscribe_tx
        .send(SubscribeRequest {
            transactions: tx_filter_map,
            commitment: Some(CommitmentLevel::Processed as i32),
            ..Default::default()
        })
        .await?;

    let mut latencies: Vec<i64> = Vec::new();
    let mut count = 0;

    while let Some(update_res) = subscribe_rx.next().await {
        match update_res {
            Ok(update) => {
                if let Some(UpdateOneof::Transaction(tx)) = update.update_oneof {
                    let recv_time = SystemTime::now();
                    
                    if let Some(created_at) = update.created_at {
                        let created_time = SystemTime::UNIX_EPOCH + Duration::new(
                            created_at.seconds as u64,
                            created_at.nanos as u32,
                        );
                        
                        if let Ok(latency) = recv_time.duration_since(created_time) {
                            let latency_ms = latency.as_millis() as i64;
                            latencies.push(latency_ms);
                            count += 1;
                            
                            println!("Transaction latency: {}ms", latency_ms);
                            
                            if count % 100 == 0 {
                                print_statistics(&latencies, count);
                            }
                        }
                    }
                }
            }
            Err(e) => {
                error!("Stream error: {:?}", e);
                break;
            }
        }
    }

    Ok(())
}

fn print_statistics(latencies: &[i64], count: usize) {
    if latencies.is_empty() {
        return;
    }
    
    let mut sorted = latencies.to_vec();
    sorted.sort();
    
    let median = if sorted.len() % 2 == 0 {
        let mid = sorted.len() / 2;
        (sorted[mid - 1] + sorted[mid]) / 2
    } else {
        sorted[sorted.len() / 2]
    };
    
    let min = *sorted.first().unwrap();
    let max = *sorted.last().unwrap();
    let sum: i64 = sorted.iter().sum();
    let mean = sum / sorted.len() as i64;
    
    let p95_idx = (sorted.len() as f64 * 0.95) as usize;
    let p95 = sorted[p95_idx.min(sorted.len() - 1)];
    
    eprintln!("--- Statistics after {} transactions ---", count);
    eprintln!("Latency (created_at to receive) in milliseconds:");
    eprintln!("  Min: {}, Max: {}", min, max);
    eprintln!("  Mean: {}, Median: {}", mean, median);
    eprintln!("  P95: {}", p95);
    eprintln!("---");
}
```

**Hạn chế quan trọng:** Phương pháp này chỉ đo khoảng thời gian từ lúc LaserStream tạo thông báo đến lúc bạn nhận được thông báo. Phương pháp này không tính đến bất kỳ độ trễ thượng nguồn nào giữa sự kiện blockchain và quá trình xử lý của LaserStream.

**Chạy bài kiểm thử:**

```bash theme={"system"}
cargo run --bin timestamp-analysis
```

**Đầu ra mẫu:**

```
Transaction latency: 45ms
Transaction latency: 52ms
Transaction latency: 38ms
--- Statistics after 100 transactions ---
Latency (created_at to receive) in milliseconds:
  Min: 28, Max: 89
  Mean: 47, Median: 45
  P95: 72
---
```

Phương pháp này cung cấp thông tin chuyên sâu về độ trễ mạng và độ trễ xử lý giữa LaserStream với ứng dụng của bạn, nhưng nên được sử dụng cùng Phương pháp 1 để có kết quả phân tích toàn diện.

***

## Phương pháp hay nhất để kiểm thử độ trễ

### Các nguyên tắc chính

* **Đồng vị trí**: Triển khai các bài kiểm thử trong cùng khu vực với endpoint LaserStream để giảm thiểu độ trễ mạng
* **Nhiều phương pháp**: Sử dụng phép so sánh các luồng song song (Phương pháp 1) làm chỉ số chính và bổ sung bằng phân tích dấu thời gian
* **Giám sát dài hạn**: Chạy kiểm thử trong thời gian dài để ghi nhận các điều kiện mạng và mức độ tắc nghẽn blockchain khác nhau
* **Phân tích thống kê**: Tập trung vào các phân vị (P95, P99) thay vì chỉ dùng giá trị trung bình để hiểu độ trễ đuôi

### Diễn giải kết quả

1. **Thiết lập đường cơ sở**: Chạy kiểm thử trong ít nhất 1 giờ để thiết lập hiệu năng cơ sở trong điều kiện bình thường
2. **Xác định quy luật**: Tìm quy luật trong các đợt tăng đột biến về độ trễ — chúng có tương quan với hoạt động blockchain cao hoặc tình trạng tắc nghẽn mạng không?
3. **So sánh các phân vị**: Độ trễ P95 thường quan trọng hơn độ trễ trung bình đối với trải nghiệm người dùng
4. **Theo dõi tính nhất quán**: Hiệu năng ổn định thường có giá trị hơn độ trễ tối thiểu tuyệt đối

Hãy nhớ rằng độ trễ blockchain vốn biến thiên do các yêu cầu về đồng thuận mạng. Hãy tập trung vào chênh lệch hiệu năng tương đối và tính nhất quán thay vì các con số tuyệt đối.
