kawa-storage 0.1.0

High-performance storage engine for Kawa message broker
Documentation

๐Ÿ’พ Kawa Storage - ้ซ˜ๆ€ง่ƒฝใ‚คใƒ™ใƒณใƒˆใ‚นใƒˆใƒฌใƒผใ‚ธใ‚จใƒณใ‚ธใƒณ

kawa-storageใฏใ€Kawaใƒ—ใƒญใ‚ธใ‚งใ‚ฏใƒˆใฎๆ ธใจใชใ‚‹้ซ˜ๆ€ง่ƒฝใ‚นใƒˆใƒฌใƒผใ‚ธใ‚จใƒณใ‚ธใƒณใงใ™ใ€‚Write-Ahead Log๏ผˆWAL๏ผ‰ใจใ‚ปใ‚ฐใƒกใƒณใƒˆๅŒ–ใซใ‚ˆใ‚‹ๆฐธ็ถšๅŒ–ใ€ใ‚คใƒ™ใƒณใƒˆใ‚ฝใƒผใ‚ทใƒณใ‚ฐใƒ‘ใ‚ฟใƒผใƒณใ‚’ๅฎŸ่ฃ…ใ—ใฆใ„ใพใ™ใ€‚

๐ŸŽฏ ่จญ่จˆ็›ฎๆจ™

  • โšก ้ซ˜ๆ€ง่ƒฝ: 274K+ writes/sec, 1.2M+ reads/sec
  • ๐Ÿ”’ ๆ•ดๅˆๆ€ง: CRC32ใƒใ‚งใƒƒใ‚ฏใ‚ตใƒ ใซใ‚ˆใ‚‹ใƒ‡ใƒผใ‚ฟๆ•ดๅˆๆ€งไฟ่จผ
  • ๐Ÿ“ˆ ใ‚นใ‚ฑใƒผใƒฉใƒ“ใƒชใƒ†ใ‚ฃ: ใ‚ปใ‚ฐใƒกใƒณใƒˆๅŒ–ใซใ‚ˆใ‚‹ๅŠน็އ็š„ใชใƒ‡ใƒผใ‚ฟ็ฎก็†
  • ๐Ÿ›ก๏ธ ไฟก้ ผๆ€ง: ้šœๅฎณ่€ๆ€งใจใƒ‡ใƒผใ‚ฟๅพฉๆ—งๆฉŸ่ƒฝ

๐Ÿ—๏ธ ใ‚ขใƒผใ‚ญใƒ†ใ‚ฏใƒใƒฃ

โ”Œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”
โ”‚                 Kawa Storage Engine                         โ”‚
โ”œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ค
โ”‚  ๐Ÿ“ค Public API                                             โ”‚
โ”‚   โ”œโ”€โ”€ StorageEngine                                        โ”‚
โ”‚   โ”œโ”€โ”€ append_event() / read_events()                       โ”‚
โ”‚   โ””โ”€โ”€ get_latest_offset()                                  โ”‚
โ”œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ค
โ”‚  ๐Ÿ“ Write-Ahead Log (WAL)                                  โ”‚
โ”‚   โ”œโ”€โ”€ ใƒ‘ใƒผใƒ†ใ‚ฃใ‚ทใƒงใƒณๅˆฅใ‚ชใƒ•ใ‚ปใƒƒใƒˆ็ฎก็†                          โ”‚
โ”‚   โ”œโ”€โ”€ ใ‚ปใ‚ฐใƒกใƒณใƒˆ่‡ชๅ‹•ใƒญใƒผใƒ†ใƒผใ‚ทใƒงใƒณ                           โ”‚
โ”‚   โ””โ”€โ”€ ้žๅŒๆœŸๆ›ธใ่พผใฟใƒป่ชญใฟๅ–ใ‚Š                             โ”‚
โ”œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ค
โ”‚  ๐Ÿ“ฆ Segment Management                                      โ”‚
โ”‚   โ”œโ”€โ”€ ๅ›บๅฎšใ‚ตใ‚คใ‚บใ‚ปใ‚ฐใƒกใƒณใƒˆ๏ผˆ1GB default๏ผ‰                   โ”‚
โ”‚   โ”œโ”€โ”€ ใƒกใƒขใƒชใƒžใƒƒใƒ—ใƒ‰I/O                                     โ”‚
โ”‚   โ”œโ”€โ”€ CRC32ๆ•ดๅˆๆ€งใƒใ‚งใƒƒใ‚ฏ                                   โ”‚
โ”‚   โ””โ”€โ”€ ไธฆ่กŒใ‚ขใ‚ฏใ‚ปใ‚นๅˆถๅพก                                     โ”‚
โ”œโ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”ค
โ”‚  ๐Ÿ” Event Management                                        โ”‚
โ”‚   โ”œโ”€โ”€ EventId๏ผˆUUID v4๏ผ‰                                   โ”‚
โ”‚   โ”œโ”€โ”€ ใ‚คใƒ™ใƒณใƒˆใƒกใ‚ฟใƒ‡ใƒผใ‚ฟ                                     โ”‚
โ”‚   โ”œโ”€โ”€ JSONใƒปใƒใ‚คใƒŠใƒชใƒ‡ใƒผใ‚ฟใ‚ตใƒใƒผใƒˆ                          โ”‚
โ”‚   โ””โ”€โ”€ ใ‚ฟใ‚คใƒ ใ‚นใ‚ฟใƒณใƒ—็ฎก็†                                     โ”‚
โ””โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”€โ”˜

๐Ÿš€ ใ‚ฏใ‚คใƒƒใ‚ฏใ‚นใ‚ฟใƒผใƒˆ

ๅŸบๆœฌ็š„ใชไฝฟ็”จๆ–นๆณ•

use kawa_storage::{
    StorageEngine, StorageConfig, 
    Topic, Partition, EventData, Offset
};

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    // ใ‚นใƒˆใƒฌใƒผใ‚ธใ‚จใƒณใ‚ธใƒณใ‚’ๅˆๆœŸๅŒ–
    let config = StorageConfig {
        data_dir: "./kawa-data".into(),
        segment_size: 1024 * 1024 * 1024, // 1GB
        sync_interval_ms: 1000,
        enable_compression: false,
    };
    
    let storage = StorageEngine::new(config).await?;
    
    // ใ‚คใƒ™ใƒณใƒˆใ‚’่ฟฝๅŠ 
    let topic = Topic::new("user-events");
    let partition = Partition::new(0);
    let event_data = EventData::from_json(r#"{"user_id": "123", "action": "login"}"#);
    
    let (event_id, offset) = storage.append_event(&topic, partition, event_data).await?;
    println!("Event {} stored at offset {}", event_id, offset);
    
    // ใ‚คใƒ™ใƒณใƒˆใ‚’่ชญใฟๅ–ใ‚Š
    let events = storage.read_events(&topic, partition, Offset::new(0), 10).await?;
    for event in events {
        println!("Event: {:?}", event);
    }
    
    // ๆœ€ๆ–ฐใ‚ชใƒ•ใ‚ปใƒƒใƒˆใ‚’ๅ–ๅพ—
    if let Some(latest) = storage.get_latest_offset(&topic, partition).await? {
        println!("Latest offset: {}", latest);
    }
    
    Ok(())
}

ใƒ‘ใƒ•ใ‚ฉใƒผใƒžใƒณใ‚นๆœ€้ฉๅŒ–

// ๅคง้‡ใƒ‡ใƒผใ‚ฟใฎๅŠน็އ็š„ใชๆ›ธใ่พผใฟ
for i in 0..10000 {
    let event_data = EventData::from_bytes(format!("Event {}", i).into_bytes());
    storage.append_event(&topic, partition, event_data).await?;
}

// ใƒใƒƒใƒ่ชญใฟๅ–ใ‚Š
let events = storage.read_events(&topic, partition, Offset::new(0), 1000).await?;

๐Ÿ“Š ๆ€ง่ƒฝ็‰นๆ€ง

ใƒ™ใƒณใƒใƒžใƒผใ‚ฏ็ตๆžœ

ๆ“ไฝœ ๆ€ง่ƒฝ ๆกไปถ
ๆ›ธใ่พผใฟ 274,816 events/sec Release build, 1000 events
่ชญใฟๅ–ใ‚Š 1,190,890 events/sec Release build, 1000 events
่ตทๅ‹•ๆ™‚้–“ < 100ms ็ฉบใฎใƒ‡ใƒผใ‚ฟใƒ‡ใ‚ฃใƒฌใ‚ฏใƒˆใƒช
ใƒกใƒขใƒชไฝฟ็”จ้‡ < 50MB ๅŸบๆœฌๆ“ไฝœๆ™‚

ใ‚นใ‚ฑใƒผใƒฉใƒ“ใƒชใƒ†ใ‚ฃ

  • ๆœ€ๅคงใ‚ปใ‚ฐใƒกใƒณใƒˆใ‚ตใ‚คใ‚บ: ๅˆถ้™ใชใ—๏ผˆๆŽจๅฅจ1GB๏ผ‰
  • ๆœ€ๅคงใƒ‘ใƒผใƒ†ใ‚ฃใ‚ทใƒงใƒณๆ•ฐ: ๅˆถ้™ใชใ—
  • ๆœ€ๅคงใƒˆใƒ”ใƒƒใ‚ฏๆ•ฐ: ๅˆถ้™ใชใ—
  • ไธฆ่กŒใ‚ขใ‚ฏใ‚ปใ‚น: ๅฎŒๅ…จๅฏพๅฟœ

๐Ÿ”ง ่จญๅฎšใ‚ชใƒ—ใ‚ทใƒงใƒณ

use kawa_storage::StorageConfig;

let config = StorageConfig {
    // ใƒ‡ใƒผใ‚ฟไฟๅญ˜ใƒ‡ใ‚ฃใƒฌใ‚ฏใƒˆใƒช
    data_dir: PathBuf::from("./kawa-data"),
    
    // ใ‚ปใ‚ฐใƒกใƒณใƒˆใƒ•ใ‚กใ‚คใƒซใ‚ตใ‚คใ‚บ๏ผˆใƒใ‚คใƒˆ๏ผ‰
    segment_size: 1024 * 1024 * 1024, // 1GB
    
    // ๅŒๆœŸ้–“้š”๏ผˆใƒŸใƒช็ง’๏ผ‰
    sync_interval_ms: 1000,
    
    // ๅœง็ธฎๆœ‰ๅŠนๅŒ–๏ผˆๅฐ†ๆฅๅฎŸ่ฃ…๏ผ‰
    enable_compression: false,
};

๐Ÿ—ƒ๏ธ ใƒ‡ใƒผใ‚ฟๆง‹้€ 

Event

pub struct Event {
    pub id: EventId,           // UUID v4
    pub topic: Topic,          // ใƒˆใƒ”ใƒƒใ‚ฏๅ
    pub partition: Partition,  // ใƒ‘ใƒผใƒ†ใ‚ฃใ‚ทใƒงใƒณ็•ชๅท
    pub data: EventData,       // ใ‚คใƒ™ใƒณใƒˆใƒ‡ใƒผใ‚ฟ
    pub timestamp: Timestamp,  // ไฝœๆˆๆ™‚ๅˆป
    pub offset: Option<Offset>, // ใ‚ชใƒ•ใ‚ปใƒƒใƒˆ๏ผˆ่ชญใฟๅ–ใ‚Šๆ™‚่จญๅฎš๏ผ‰
}

EventData

pub enum EventData {
    Json(String),      // JSONๆ–‡ๅญ—ๅˆ—
    Binary(Vec<u8>),   // ใƒใ‚คใƒŠใƒชใƒ‡ใƒผใ‚ฟ
    Text(String),      // ใƒ—ใƒฌใƒผใƒณใƒ†ใ‚ญใ‚นใƒˆ
}

// ไฝฟ็”จไพ‹
let json_data = EventData::from_json(r#"{"key": "value"}"#);
let binary_data = EventData::from_bytes(vec![1, 2, 3, 4]);
let text_data = EventData::from_text("Hello, World!");

๐Ÿงช ใƒ†ใ‚นใƒˆ

ใƒ†ใ‚นใƒˆๅฎŸ่กŒ

# ๅ…จใƒ†ใ‚นใƒˆๅฎŸ่กŒ
cargo test

# ็ตฑๅˆใƒ†ใ‚นใƒˆ
cargo test --test integration_tests

# ใƒ™ใƒณใƒใƒžใƒผใ‚ฏ
cargo test bench_ --release -- --nocapture

# ใ‚ซใƒใƒฌใƒƒใ‚ธ
cargo tarpaulin --out Html

ใƒ†ใ‚นใƒˆๆง‹ๆˆ

  • Unit Tests: ๅ„ใƒขใ‚ธใƒฅใƒผใƒซใฎๅŸบๆœฌๆฉŸ่ƒฝ
  • Integration Tests: ใ‚จใƒณใƒ‰ใƒ„ใƒผใ‚จใƒณใƒ‰ใ‚ทใƒŠใƒชใ‚ช
  • Benchmark Tests: ๆ€ง่ƒฝๆธฌๅฎš
  • Property Tests: ใƒฉใƒณใƒ€ใƒ ใƒ‡ใƒผใ‚ฟใงใฎๆคœ่จผ

๐Ÿšง ๅฐ†ๆฅใฎๆฉŸ่ƒฝ

v0.2.0

  • ๅœง็ธฎใ‚ตใƒใƒผใƒˆ: LZ4/Snappyๅœง็ธฎ
  • ใ‚คใƒณใƒ‡ใƒƒใ‚ฏใ‚นๆง‹็ฏ‰: ้ซ˜้€Ÿๆคœ็ดขๆฉŸ่ƒฝ
  • ใƒกใƒขใƒชใƒ—ใƒผใƒซ: ใ‚ขใƒญใ‚ฑใƒผใ‚ทใƒงใƒณๆœ€้ฉๅŒ–
  • ใƒใƒƒใƒAPI: ่ค‡ๆ•ฐใ‚คใƒ™ใƒณใƒˆไธ€ๆ‹ฌๅ‡ฆ็†

v0.3.0

  • ใƒฌใƒ—ใƒชใ‚ฑใƒผใ‚ทใƒงใƒณ: ๅคšใƒŽใƒผใƒ‰ๅฏพๅฟœ
  • ๆš—ๅทๅŒ–: ไฟๅญ˜ๆ™‚ๆš—ๅทๅŒ–
  • ใ‚นใ‚ญใƒผใƒžๆคœ่จผ: ใ‚คใƒ™ใƒณใƒˆๆง‹้€ ๆคœ่จผ
  • ่‡ชๅ‹•ๅœง็ธฎ: ๅคใ„ใ‚ปใ‚ฐใƒกใƒณใƒˆๆ•ด็†

โš ๏ธ ๅˆถ้™ไบ‹้ …

  • ใƒˆใƒฉใƒณใ‚ถใ‚ฏใ‚ทใƒงใƒณ: ็พๅœจๆœชๅฏพๅฟœ
  • ๅ‰Š้™คๆ“ไฝœ: ใ‚คใƒ™ใƒณใƒˆใฎๅ‰Š้™คใฏๆœชๅฏพๅฟœ
  • ๅˆ†ๆ•ฃใ‚นใƒˆใƒฌใƒผใ‚ธ: ๅ˜ไธ€ใƒŽใƒผใƒ‰ใฎใฟ
  • ใ‚นใ‚ญใƒผใƒž้€ฒๅŒ–: ็พๅœจๆœชๅฏพๅฟœ

๐Ÿ”— ้–ข้€ฃใƒ‰ใ‚ญใƒฅใƒกใƒณใƒˆ

๐Ÿ“ ใ‚จใƒฉใƒผใƒใƒณใƒ‰ใƒชใƒณใ‚ฐ

use kawa_storage::{StorageError, StorageResult};

match storage.append_event(&topic, partition, event_data).await {
    Ok((event_id, offset)) => {
        println!("Success: {} at {}", event_id, offset);
    }
    Err(StorageError::Io { source, .. }) => {
        eprintln!("I/O Error: {}", source);
    }
    Err(StorageError::Serialization { source, .. }) => {
        eprintln!("Serialization Error: {}", source);
    }
    Err(e) => {
        eprintln!("Other Error: {}", e);
    }
}

kawa-storage - RustใงๅฎŸ่ฃ…ใ•ใ‚ŒใŸ้ซ˜ๆ€ง่ƒฝใ‚คใƒ™ใƒณใƒˆใ‚นใƒˆใƒฌใƒผใ‚ธ ๐Ÿš€