mirror of
https://bitbucket.org/s_l_teichmann/mtsatellite
synced 2025-06-28 22:26:47 +02:00
Moved sub programs into folder cmd to clean up project structure.
This commit is contained in:
137
cmd/mtdbconverter/leveldb.go
Normal file
137
cmd/mtdbconverter/leveldb.go
Normal file
@ -0,0 +1,137 @@
|
||||
// Copyright 2014 by Sascha L. Teichmann
|
||||
// Use of this source code is governed by the MIT license
|
||||
// that can be found in the LICENSE file.
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
"os"
|
||||
|
||||
"bitbucket.org/s_l_teichmann/mtredisalize/common"
|
||||
|
||||
leveldb "github.com/jmhodges/levigo"
|
||||
)
|
||||
|
||||
type (
|
||||
LevelDBBlockProducer struct {
|
||||
db *leveldb.DB
|
||||
opts *leveldb.Options
|
||||
ro *leveldb.ReadOptions
|
||||
iterator *leveldb.Iterator
|
||||
splitter common.KeySplitter
|
||||
decoder common.KeyDecoder
|
||||
}
|
||||
|
||||
LevelDBBlockConsumer struct {
|
||||
db *leveldb.DB
|
||||
opts *leveldb.Options
|
||||
wo *leveldb.WriteOptions
|
||||
joiner common.KeyJoiner
|
||||
encoder common.KeyEncoder
|
||||
}
|
||||
)
|
||||
|
||||
func NewLevelDBBlockProducer(path string,
|
||||
splitter common.KeySplitter,
|
||||
decoder common.KeyDecoder) (ldbp *LevelDBBlockProducer, err error) {
|
||||
|
||||
// check if we can stat it -> exists.
|
||||
if _, err = os.Stat(path); err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
opts := leveldb.NewOptions()
|
||||
opts.SetCreateIfMissing(false)
|
||||
|
||||
var db *leveldb.DB
|
||||
if db, err = leveldb.Open(path, opts); err != nil {
|
||||
opts.Close()
|
||||
return
|
||||
}
|
||||
|
||||
ro := leveldb.NewReadOptions()
|
||||
ro.SetFillCache(false)
|
||||
|
||||
iterator := db.NewIterator(ro)
|
||||
iterator.SeekToFirst()
|
||||
|
||||
ldbp = &LevelDBBlockProducer{
|
||||
db: db,
|
||||
opts: opts,
|
||||
ro: ro,
|
||||
iterator: iterator,
|
||||
splitter: splitter,
|
||||
decoder: decoder}
|
||||
return
|
||||
}
|
||||
|
||||
func (ldbp *LevelDBBlockProducer) Close() error {
|
||||
if ldbp.iterator != nil {
|
||||
ldbp.iterator.Close()
|
||||
}
|
||||
ldbp.ro.Close()
|
||||
ldbp.db.Close()
|
||||
ldbp.opts.Close()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (ldbp *LevelDBBlockProducer) Next(block *common.Block) (err error) {
|
||||
if ldbp.iterator == nil {
|
||||
err = common.ErrNoMoreBlocks
|
||||
return
|
||||
}
|
||||
if !ldbp.iterator.Valid() {
|
||||
if err = ldbp.iterator.GetError(); err == nil {
|
||||
err = common.ErrNoMoreBlocks
|
||||
}
|
||||
ldbp.iterator.Close()
|
||||
ldbp.iterator = nil
|
||||
return
|
||||
}
|
||||
var key int64
|
||||
if key, err = ldbp.decoder(ldbp.iterator.Key()); err != nil {
|
||||
return
|
||||
}
|
||||
block.Coord = ldbp.splitter(key)
|
||||
block.Data = ldbp.iterator.Value()
|
||||
ldbp.iterator.Next()
|
||||
return
|
||||
}
|
||||
|
||||
func NewLevelDBBlockConsumer(
|
||||
path string,
|
||||
joiner common.KeyJoiner,
|
||||
encoder common.KeyEncoder) (ldbc *LevelDBBlockConsumer, err error) {
|
||||
|
||||
opts := leveldb.NewOptions()
|
||||
opts.SetCreateIfMissing(true)
|
||||
|
||||
var db *leveldb.DB
|
||||
if db, err = leveldb.Open(path, opts); err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
ldbc = &LevelDBBlockConsumer{
|
||||
db: db,
|
||||
opts: opts,
|
||||
wo: leveldb.NewWriteOptions(),
|
||||
joiner: joiner,
|
||||
encoder: encoder}
|
||||
return
|
||||
}
|
||||
|
||||
func (ldbc *LevelDBBlockConsumer) Close() error {
|
||||
ldbc.wo.Close()
|
||||
ldbc.db.Close()
|
||||
ldbc.opts.Close()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (ldbc *LevelDBBlockConsumer) Consume(block *common.Block) (err error) {
|
||||
var encodedKey []byte
|
||||
if encodedKey, err = ldbc.encoder(ldbc.joiner(block.Coord)); err != nil {
|
||||
return
|
||||
}
|
||||
err = ldbc.db.Put(ldbc.wo, encodedKey, block.Data)
|
||||
return
|
||||
}
|
169
cmd/mtdbconverter/main.go
Normal file
169
cmd/mtdbconverter/main.go
Normal file
@ -0,0 +1,169 @@
|
||||
// Copyright 2014 by Sascha L. Teichmann
|
||||
// Use of this source code is governed by the MIT license
|
||||
// that can be found in the LICENSE file.
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
"flag"
|
||||
"fmt"
|
||||
"log"
|
||||
"os"
|
||||
"sync"
|
||||
|
||||
"bitbucket.org/s_l_teichmann/mtredisalize/common"
|
||||
)
|
||||
|
||||
func usage() {
|
||||
fmt.Fprintf(os.Stderr,
|
||||
"Usage: %s [<options>] <source database> <dest database>\n", os.Args[0])
|
||||
fmt.Fprintln(os.Stderr, "Options:")
|
||||
flag.PrintDefaults()
|
||||
}
|
||||
|
||||
func selectKeySplitter(interleaved bool) common.KeySplitter {
|
||||
if interleaved {
|
||||
return common.InterleavedToCoord
|
||||
}
|
||||
return common.PlainToCoord
|
||||
}
|
||||
|
||||
func selectKeyJoiner(interleaved bool) common.KeyJoiner {
|
||||
if interleaved {
|
||||
return common.CoordToInterleaved
|
||||
}
|
||||
return common.CoordToPlain
|
||||
}
|
||||
|
||||
func selectKeyDecoder(interleaved bool) common.KeyDecoder {
|
||||
if interleaved {
|
||||
return common.DecodeFromBigEndian
|
||||
}
|
||||
return common.DecodeStringFromBytes
|
||||
}
|
||||
|
||||
func selectKeyEncoder(interleaved bool) common.KeyEncoder {
|
||||
if interleaved {
|
||||
return common.EncodeToBigEndian
|
||||
}
|
||||
return common.EncodeStringToBytes
|
||||
}
|
||||
|
||||
func copyProducerToConsumer(producer common.BlockProducer, consumer common.BlockConsumer) error {
|
||||
|
||||
blocks := make(chan *common.Block)
|
||||
done := make(chan struct{})
|
||||
defer close(done)
|
||||
|
||||
pool := sync.Pool{New: func() interface{} { return new(common.Block) }}
|
||||
|
||||
go func() {
|
||||
defer close(blocks)
|
||||
for {
|
||||
block := pool.Get().(*common.Block)
|
||||
if err := producer.Next(block); err != nil {
|
||||
if err != common.ErrNoMoreBlocks {
|
||||
log.Printf("Reading failed: %s\n", err)
|
||||
}
|
||||
return
|
||||
}
|
||||
select {
|
||||
case blocks <- block:
|
||||
case <-done:
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
i := 0
|
||||
for block := range blocks {
|
||||
if err := consumer.Consume(block); err != nil {
|
||||
return err
|
||||
}
|
||||
block.Data = nil
|
||||
pool.Put(block)
|
||||
i++
|
||||
if i%1000 == 0 {
|
||||
log.Printf("%d blocks transferred.\n", i)
|
||||
}
|
||||
}
|
||||
log.Printf("%d blocks transferred in total.\n", i)
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func main() {
|
||||
var (
|
||||
srcBackend string
|
||||
dstBackend string
|
||||
srcInterleaved bool
|
||||
dstInterleaved bool
|
||||
)
|
||||
|
||||
flag.Usage = usage
|
||||
|
||||
flag.StringVar(&srcBackend, "source-backend", "sqlite",
|
||||
"type of source database (leveldb, sqlite)")
|
||||
flag.StringVar(&srcBackend, "sb", "sqlite",
|
||||
"type of source database (leveldb, sqlite). Shorthand")
|
||||
flag.StringVar(&dstBackend, "dest-backend", "leveldb",
|
||||
"type of destination database (leveldb, sqlite)")
|
||||
flag.StringVar(&dstBackend, "db", "leveldb",
|
||||
"type of destination database (leveldb, sqlite). Shorthand")
|
||||
flag.BoolVar(&srcInterleaved, "source-interleaved", false,
|
||||
"Is source database interleaved?")
|
||||
flag.BoolVar(&srcInterleaved, "si", false,
|
||||
"Is source database interleaved? Shorthand")
|
||||
flag.BoolVar(&dstInterleaved, "dest-interleaved", true,
|
||||
"Should dest database be interleaved?")
|
||||
flag.BoolVar(&dstInterleaved, "di", true,
|
||||
"Should source database be interleaved? Shorthand")
|
||||
|
||||
flag.Parse()
|
||||
|
||||
if flag.NArg() < 2 {
|
||||
log.Fatal("Missing source and/or destination database.")
|
||||
}
|
||||
|
||||
var (
|
||||
producer common.BlockProducer
|
||||
consumer common.BlockConsumer
|
||||
err error
|
||||
)
|
||||
|
||||
if srcBackend == "sqlite" {
|
||||
if producer, err = NewSQLiteBlockProducer(
|
||||
flag.Arg(0),
|
||||
selectKeySplitter(srcInterleaved)); err != nil {
|
||||
log.Fatalf("Cannot open '%s': %s", flag.Arg(0), err)
|
||||
}
|
||||
} else { // LevelDB
|
||||
if producer, err = NewLevelDBBlockProducer(
|
||||
flag.Arg(0),
|
||||
selectKeySplitter(srcInterleaved),
|
||||
selectKeyDecoder(srcInterleaved)); err != nil {
|
||||
log.Fatalf("Cannot open '%s': %s", flag.Arg(0), err)
|
||||
}
|
||||
}
|
||||
defer producer.Close()
|
||||
|
||||
if dstBackend == "sqlite" {
|
||||
if consumer, err = NewSQLiteBlockConsumer(
|
||||
flag.Arg(1),
|
||||
selectKeyJoiner(dstInterleaved)); err != nil {
|
||||
log.Fatalf("Cannot open '%s': %s", flag.Arg(1), err)
|
||||
}
|
||||
} else { // LevelDB
|
||||
if consumer, err = NewLevelDBBlockConsumer(
|
||||
flag.Arg(1),
|
||||
selectKeyJoiner(dstInterleaved),
|
||||
selectKeyEncoder(dstInterleaved)); err != nil {
|
||||
log.Fatalf("Cannot open '%s': %s", flag.Arg(1), err)
|
||||
}
|
||||
}
|
||||
defer consumer.Close()
|
||||
|
||||
if err = copyProducerToConsumer(producer, consumer); err != nil {
|
||||
log.Fatalf("Database transfer failed: %s\n", err)
|
||||
}
|
||||
}
|
175
cmd/mtdbconverter/sqlite.go
Normal file
175
cmd/mtdbconverter/sqlite.go
Normal file
@ -0,0 +1,175 @@
|
||||
// Copyright 2014 by Sascha L. Teichmann
|
||||
// Use of this source code is governed by the MIT license
|
||||
// that can be found in the LICENSE file.
|
||||
|
||||
package main
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"errors"
|
||||
"os"
|
||||
|
||||
"bitbucket.org/s_l_teichmann/mtredisalize/common"
|
||||
|
||||
_ "github.com/mattn/go-sqlite3"
|
||||
)
|
||||
|
||||
const (
|
||||
createSql = "CREATE TABLE blocks (pos INT NOT NULL PRIMARY KEY, data BLOB)"
|
||||
insertSql = "INSERT INTO blocks (pos, data) VALUES (?, ?)"
|
||||
deleteSql = "DELETE FROM blocks"
|
||||
selectSql = "SELECT pos, data FROM blocks"
|
||||
)
|
||||
|
||||
var ErrDatabaseNotExists = errors.New("Database does not exists.")
|
||||
|
||||
const blocksPerTx = 128
|
||||
|
||||
type (
|
||||
SQLiteBlockProducer struct {
|
||||
db *sql.DB
|
||||
rows *sql.Rows
|
||||
splitter common.KeySplitter
|
||||
}
|
||||
|
||||
SQLiteBlockConsumer struct {
|
||||
db *sql.DB
|
||||
insertStmt *sql.Stmt
|
||||
tx *sql.Tx
|
||||
txCounter int
|
||||
joiner common.KeyJoiner
|
||||
}
|
||||
)
|
||||
|
||||
func fileExists(path string) bool {
|
||||
_, err := os.Stat(path)
|
||||
return !os.IsNotExist(err)
|
||||
}
|
||||
|
||||
func NewSQLiteBlockConsumer(
|
||||
path string,
|
||||
joiner common.KeyJoiner) (sbc *SQLiteBlockConsumer, err error) {
|
||||
|
||||
createNew := !fileExists(path)
|
||||
|
||||
var db *sql.DB
|
||||
if db, err = sql.Open("sqlite3", path); err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
if createNew {
|
||||
if _, err = db.Exec(createSql); err != nil {
|
||||
db.Close()
|
||||
return
|
||||
}
|
||||
} else {
|
||||
if _, err = db.Exec(deleteSql); err != nil {
|
||||
db.Close()
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
var insertStmt *sql.Stmt
|
||||
if insertStmt, err = db.Prepare(insertSql); err != nil {
|
||||
db.Close()
|
||||
return
|
||||
}
|
||||
|
||||
var tx *sql.Tx
|
||||
if tx, err = db.Begin(); err != nil {
|
||||
insertStmt.Close()
|
||||
db.Close()
|
||||
return
|
||||
}
|
||||
|
||||
sbc = &SQLiteBlockConsumer{
|
||||
db: db,
|
||||
insertStmt: insertStmt,
|
||||
tx: tx,
|
||||
joiner: joiner}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (sbc *SQLiteBlockConsumer) Close() error {
|
||||
sbc.tx.Commit()
|
||||
sbc.insertStmt.Close()
|
||||
return sbc.db.Close()
|
||||
}
|
||||
|
||||
func (sbc *SQLiteBlockConsumer) getTx() (tx *sql.Tx, err error) {
|
||||
if sbc.txCounter >= blocksPerTx {
|
||||
sbc.txCounter = 0
|
||||
if err = sbc.tx.Commit(); err != nil {
|
||||
return
|
||||
}
|
||||
if sbc.tx, err = sbc.db.Begin(); err != nil {
|
||||
return
|
||||
}
|
||||
}
|
||||
sbc.txCounter++
|
||||
tx = sbc.tx
|
||||
return
|
||||
}
|
||||
|
||||
func (sbc *SQLiteBlockConsumer) Consume(block *common.Block) (err error) {
|
||||
var tx *sql.Tx
|
||||
if tx, err = sbc.getTx(); err != nil {
|
||||
return
|
||||
}
|
||||
_, err = tx.Stmt(sbc.insertStmt).Exec(sbc.joiner(block.Coord), block.Data)
|
||||
return
|
||||
}
|
||||
|
||||
func NewSQLiteBlockProducer(
|
||||
path string,
|
||||
splitter common.KeySplitter) (sbp *SQLiteBlockProducer, err error) {
|
||||
|
||||
if !fileExists(path) {
|
||||
err = ErrDatabaseNotExists
|
||||
return
|
||||
}
|
||||
|
||||
var db *sql.DB
|
||||
if db, err = sql.Open("sqlite3", path); err != nil {
|
||||
return
|
||||
}
|
||||
|
||||
var rows *sql.Rows
|
||||
if rows, err = db.Query(selectSql); err != nil {
|
||||
db.Close()
|
||||
return
|
||||
}
|
||||
|
||||
sbp = &SQLiteBlockProducer{
|
||||
db: db,
|
||||
rows: rows,
|
||||
splitter: splitter}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
func (sbp *SQLiteBlockProducer) Next(block *common.Block) (err error) {
|
||||
if sbp.rows == nil {
|
||||
err = common.ErrNoMoreBlocks
|
||||
return
|
||||
}
|
||||
if sbp.rows.Next() {
|
||||
var key int64
|
||||
if err = sbp.rows.Scan(&key, &block.Data); err == nil {
|
||||
block.Coord = sbp.splitter(key)
|
||||
}
|
||||
} else {
|
||||
sbp.rows.Close()
|
||||
sbp.rows = nil
|
||||
err = common.ErrNoMoreBlocks
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func (sbp *SQLiteBlockProducer) Close() error {
|
||||
if sbp.rows != nil {
|
||||
sbp.rows.Close()
|
||||
}
|
||||
return sbp.db.Close()
|
||||
}
|
Reference in New Issue
Block a user