ledisdb/store/rocksdb/db.go

313 lines
6.2 KiB
Go
Raw Normal View History

2014-07-26 14:39:54 +04:00
// +build rocksdb
// Package rocksdb is a wrapper for c++ rocksdb
package rocksdb
/*
#cgo LDFLAGS: -lrocksdb
#include <rocksdb/c.h>
#include <stdlib.h>
#include "rocksdb_ext.h"
*/
import "C"
import (
"github.com/siddontang/ledisdb/config"
2014-07-26 14:39:54 +04:00
"github.com/siddontang/ledisdb/store/driver"
"os"
"unsafe"
)
const defaultFilterBits int = 10
type Store struct {
}
2014-07-26 14:39:54 +04:00
func (s Store) String() string {
2014-08-15 20:08:01 +04:00
return DBName
2014-07-26 14:39:54 +04:00
}
func (s Store) Open(path string, cfg *config.Config) (driver.IDB, error) {
2014-09-18 17:27:43 +04:00
if err := os.MkdirAll(path, 0755); err != nil {
2014-07-26 14:39:54 +04:00
return nil, err
}
db := new(DB)
db.path = path
2014-10-24 09:01:13 +04:00
db.cfg = &cfg.RocksDB
2014-07-26 14:39:54 +04:00
if err := db.open(); err != nil {
return nil, err
}
return db, nil
}
func (s Store) Repair(path string, cfg *config.Config) error {
2014-07-26 14:39:54 +04:00
db := new(DB)
db.path = path
2014-10-24 09:01:13 +04:00
db.cfg = &cfg.RocksDB
2014-07-26 14:39:54 +04:00
err := db.open()
defer db.Close()
//open ok, do not need repair
if err == nil {
return nil
}
var errStr *C.char
ldbname := C.CString(path)
2014-07-26 14:39:54 +04:00
defer C.free(unsafe.Pointer(ldbname))
C.rocksdb_repair_db(db.opts.Opt, ldbname, &errStr)
if errStr != nil {
return saveError(errStr)
}
return nil
}
type DB struct {
path string
2014-10-24 09:01:13 +04:00
cfg *config.RocksDBConfig
2014-07-26 14:39:54 +04:00
db *C.rocksdb_t
env *Env
2014-09-08 05:55:55 +04:00
opts *Options
blockOpts *BlockBasedTableOptions
2014-07-26 14:39:54 +04:00
//for default read and write options
readOpts *ReadOptions
writeOpts *WriteOptions
iteratorOpts *ReadOptions
2014-10-09 09:05:55 +04:00
syncOpts *WriteOptions
2014-07-26 14:39:54 +04:00
cache *Cache
filter *FilterPolicy
}
func (db *DB) open() error {
db.initOptions(db.cfg)
var errStr *C.char
ldbname := C.CString(db.path)
2014-07-26 14:39:54 +04:00
defer C.free(unsafe.Pointer(ldbname))
db.db = C.rocksdb_open(db.opts.Opt, ldbname, &errStr)
if errStr != nil {
db.db = nil
return saveError(errStr)
}
return nil
}
2014-10-24 09:01:13 +04:00
func (db *DB) initOptions(cfg *config.RocksDBConfig) {
2014-07-26 14:39:54 +04:00
opts := NewOptions()
2014-09-08 05:55:55 +04:00
blockOpts := NewBlockBasedTableOptions()
2014-07-26 14:39:54 +04:00
opts.SetCreateIfMissing(true)
db.env = NewDefaultEnv()
2014-10-24 09:01:13 +04:00
db.env.SetBackgroundThreads(cfg.BackgroundThreads)
db.env.SetHighPriorityBackgroundThreads(cfg.HighPriorityBackgroundThreads)
2014-07-26 14:39:54 +04:00
opts.SetEnv(db.env)
db.cache = NewLRUCache(cfg.CacheSize)
2014-09-08 05:55:55 +04:00
blockOpts.SetCache(db.cache)
2014-07-26 14:39:54 +04:00
//we must use bloomfilter
db.filter = NewBloomFilter(defaultFilterBits)
2014-09-08 05:55:55 +04:00
blockOpts.SetFilterPolicy(db.filter)
blockOpts.SetBlockSize(cfg.BlockSize)
2014-10-24 09:01:13 +04:00
opts.SetBlockBasedTableFactory(blockOpts)
2014-07-26 14:39:54 +04:00
2014-10-24 09:01:13 +04:00
opts.SetCompression(CompressionOpt(cfg.Compression))
2014-07-26 14:39:54 +04:00
opts.SetWriteBufferSize(cfg.WriteBufferSize)
opts.SetMaxOpenFiles(cfg.MaxOpenFiles)
2014-10-24 09:01:13 +04:00
opts.SetMaxBackgroundCompactions(cfg.MaxBackgroundCompactions)
opts.SetMaxBackgroundFlushes(cfg.MaxBackgroundFlushes)
opts.SetLevel0SlowdownWritesTrigger(cfg.Level0SlowdownWritesTrigger)
opts.SetLevel0StopWritesTrigger(cfg.Level0StopWritesTrigger)
opts.SetTargetFileSizeBase(cfg.TargetFileSizeBase)
opts.SetTargetFileSizeMultiplier(cfg.TargetFileSizeMultiplier)
opts.SetMaxBytesForLevelBase(cfg.MaxBytesForLevelBase)
opts.SetMaxBytesForLevelMultiplier(cfg.MaxBytesForLevelMultiplier)
opts.DisableDataSync(cfg.DisableDataSync)
opts.SetMinWriteBufferNumberToMerge(cfg.MinWriteBufferNumberToMerge)
opts.DisableAutoCompactions(cfg.DisableAutoCompactions)
opts.EnableStatistics(cfg.EnableStatistics)
opts.UseFsync(cfg.UseFsync)
opts.AllowOsBuffer(cfg.AllowOsBuffer)
opts.SetStatsDumpPeriodSec(cfg.StatsDumpPeriodSec)
2014-09-08 05:55:55 +04:00
2014-07-26 14:39:54 +04:00
db.opts = opts
2014-09-08 05:55:55 +04:00
db.blockOpts = blockOpts
2014-07-26 14:39:54 +04:00
db.readOpts = NewReadOptions()
db.writeOpts = NewWriteOptions()
2014-10-09 09:05:55 +04:00
db.syncOpts = NewWriteOptions()
db.syncOpts.SetSync(true)
2014-07-26 14:39:54 +04:00
db.iteratorOpts = NewReadOptions()
db.iteratorOpts.SetFillCache(false)
}
func (db *DB) Close() error {
if db.db != nil {
C.rocksdb_close(db.db)
db.db = nil
}
2014-09-08 05:55:55 +04:00
if db.filter != nil {
db.filter.Close()
}
2014-07-26 14:39:54 +04:00
if db.cache != nil {
db.cache.Close()
}
if db.env != nil {
db.env.Close()
}
2014-09-08 05:55:55 +04:00
//db.blockOpts.Close()
db.opts.Close()
2014-07-26 14:39:54 +04:00
db.readOpts.Close()
db.writeOpts.Close()
db.iteratorOpts.Close()
return nil
}
func (db *DB) Put(key, value []byte) error {
return db.put(db.writeOpts, key, value)
}
func (db *DB) Get(key []byte) ([]byte, error) {
return db.get(db.readOpts, key)
}
func (db *DB) Delete(key []byte) error {
return db.delete(db.writeOpts, key)
}
2014-10-09 09:05:55 +04:00
func (db *DB) SyncPut(key []byte, value []byte) error {
return db.put(db.syncOpts, key, value)
}
func (db *DB) SyncDelete(key []byte) error {
return db.delete(db.syncOpts, key)
}
2014-07-26 14:39:54 +04:00
func (db *DB) NewWriteBatch() driver.IWriteBatch {
wb := &WriteBatch{
db: db,
wbatch: C.rocksdb_writebatch_create(),
}
2014-07-26 14:39:54 +04:00
return wb
}
func (db *DB) NewIterator() driver.IIterator {
it := new(Iterator)
it.it = C.rocksdb_create_iterator(db.db, db.iteratorOpts.Opt)
return it
}
2014-08-25 10:18:23 +04:00
func (db *DB) NewSnapshot() (driver.ISnapshot, error) {
snap := &Snapshot{
db: db,
snap: C.rocksdb_create_snapshot(db.db),
readOpts: NewReadOptions(),
iteratorOpts: NewReadOptions(),
}
snap.readOpts.SetSnapshot(snap)
snap.iteratorOpts.SetSnapshot(snap)
snap.iteratorOpts.SetFillCache(false)
return snap, nil
}
2014-07-26 14:39:54 +04:00
func (db *DB) put(wo *WriteOptions, key, value []byte) error {
var errStr *C.char
var k, v *C.char
if len(key) != 0 {
k = (*C.char)(unsafe.Pointer(&key[0]))
}
if len(value) != 0 {
v = (*C.char)(unsafe.Pointer(&value[0]))
}
lenk := len(key)
lenv := len(value)
C.rocksdb_put(
db.db, wo.Opt, k, C.size_t(lenk), v, C.size_t(lenv), &errStr)
if errStr != nil {
return saveError(errStr)
}
return nil
}
func (db *DB) get(ro *ReadOptions, key []byte) ([]byte, error) {
var errStr *C.char
var vallen C.size_t
var k *C.char
if len(key) != 0 {
k = (*C.char)(unsafe.Pointer(&key[0]))
}
value := C.rocksdb_get(
db.db, ro.Opt, k, C.size_t(len(key)), &vallen, &errStr)
if errStr != nil {
return nil, saveError(errStr)
}
if value == nil {
return nil, nil
}
defer C.free(unsafe.Pointer(value))
return C.GoBytes(unsafe.Pointer(value), C.int(vallen)), nil
}
func (db *DB) delete(wo *WriteOptions, key []byte) error {
var errStr *C.char
var k *C.char
if len(key) != 0 {
k = (*C.char)(unsafe.Pointer(&key[0]))
}
C.rocksdb_delete(
db.db, wo.Opt, k, C.size_t(len(key)), &errStr)
if errStr != nil {
return saveError(errStr)
}
return nil
}
2014-07-29 13:29:51 +04:00
func (db *DB) Begin() (driver.Tx, error) {
return nil, driver.ErrTxSupport
}
2014-08-15 20:08:01 +04:00
2014-09-12 11:06:36 +04:00
func (db *DB) Compact() error {
C.rocksdb_compact_range(db.db, nil, 0, nil, 0)
return nil
}
2014-08-15 20:08:01 +04:00
func init() {
driver.Register(Store{})
}