add io.Reader based xdb input supports
This commit is contained in:
parent
f4fcd5e900
commit
0ce74b8873
|
|
@ -116,7 +116,7 @@ func Edit(sCmd string) {
|
|||
// quit directly
|
||||
break
|
||||
} else if cmd == "save" {
|
||||
err = editor.Save()
|
||||
err = editor.SaveToFile(srcFile)
|
||||
if err != nil {
|
||||
fmt.Printf("failed to save the changes: %s\n", err)
|
||||
continue
|
||||
|
|
|
|||
|
|
@ -9,6 +9,7 @@ package xdb
|
|||
import (
|
||||
"container/list"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
|
|
@ -18,8 +19,7 @@ type Editor struct {
|
|||
verison *Version
|
||||
|
||||
// source ip file
|
||||
srcPath string
|
||||
srcHandle *os.File
|
||||
srcHandle io.ReadCloser
|
||||
toSave bool
|
||||
|
||||
// segments list
|
||||
|
|
@ -41,17 +41,20 @@ func NewEditor(version *Version, srcFile string) (*Editor, error) {
|
|||
return nil, err
|
||||
}
|
||||
|
||||
return INewEditor(version, srcHandle)
|
||||
}
|
||||
|
||||
func INewEditor(version *Version, srcReader io.ReadCloser) (*Editor, error) {
|
||||
e := &Editor{
|
||||
verison: version,
|
||||
srcPath: srcPath,
|
||||
srcHandle: srcHandle,
|
||||
srcHandle: srcReader,
|
||||
toSave: false,
|
||||
segments: list.New(),
|
||||
rgCache: NewRegionCache(),
|
||||
}
|
||||
|
||||
// load the segments
|
||||
if err = e.loadSegments(); err != nil {
|
||||
if err := e.loadSegments(); err != nil {
|
||||
return nil, fmt.Errorf("failed to load segments: %s", err)
|
||||
}
|
||||
|
||||
|
|
@ -372,11 +375,6 @@ func (e *Editor) PutFile(src string, cb func(newSeg *Segment, oldList []*Segment
|
|||
return oldRows, newRows, nil
|
||||
}
|
||||
|
||||
// save the changes to the source file.
|
||||
func (e *Editor) Save() error {
|
||||
return e.SaveToFile(e.srcPath)
|
||||
}
|
||||
|
||||
func (e *Editor) SaveToFile(dstFile string) error {
|
||||
// check the to-save flag
|
||||
if !e.toSave {
|
||||
|
|
|
|||
|
|
@ -55,6 +55,7 @@ package xdb
|
|||
import (
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"math"
|
||||
"os"
|
||||
|
|
@ -72,11 +73,17 @@ const (
|
|||
VectorIndexLength = VectorIndexRows * VectorIndexCols * VectorIndexSize
|
||||
)
|
||||
|
||||
type WriteSeekCloser interface {
|
||||
io.Writer
|
||||
io.Seeker
|
||||
io.Closer
|
||||
}
|
||||
|
||||
type Maker struct {
|
||||
version *Version
|
||||
|
||||
srcHandle *os.File
|
||||
dstHandle *os.File
|
||||
srcHandle io.ReadCloser
|
||||
dstHandle WriteSeekCloser
|
||||
|
||||
// self-define field index
|
||||
fields []int
|
||||
|
|
@ -107,11 +114,15 @@ func NewMaker(version *Version, policy IndexPolicy, srcFile string, dstFile stri
|
|||
return nil, fmt.Errorf("open target file `%s`: %w", dstFile, err)
|
||||
}
|
||||
|
||||
return INewMaker(version, policy, srcHandle, dstHandle, fields), nil
|
||||
}
|
||||
|
||||
func INewMaker(version *Version, policy IndexPolicy, srcReader io.ReadCloser, dstWriter WriteSeekCloser, fields []int) *Maker {
|
||||
return &Maker{
|
||||
version: version,
|
||||
|
||||
srcHandle: srcHandle,
|
||||
dstHandle: dstHandle,
|
||||
srcHandle: srcReader,
|
||||
dstHandle: dstWriter,
|
||||
|
||||
// fields filter index
|
||||
fields: fields,
|
||||
|
|
@ -121,7 +132,7 @@ func NewMaker(version *Version, policy IndexPolicy, srcFile string, dstFile stri
|
|||
regionPool: map[string]uint32{},
|
||||
regionCache: NewRegionCache(),
|
||||
vectorIndex: make([]byte, VectorIndexLength),
|
||||
}, nil
|
||||
}
|
||||
}
|
||||
|
||||
func (m *Maker) initDbHeader() error {
|
||||
|
|
|
|||
|
|
@ -8,6 +8,7 @@ package xdb
|
|||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
"log/slog"
|
||||
"os"
|
||||
"sort"
|
||||
|
|
@ -16,8 +17,8 @@ import (
|
|||
)
|
||||
|
||||
type Processor struct {
|
||||
srcHandle *os.File
|
||||
dstHandle *os.File
|
||||
srcReader io.ReadCloser
|
||||
dstWriter io.WriteCloser
|
||||
|
||||
// value clear
|
||||
clearBasedIndex int
|
||||
|
|
@ -45,9 +46,17 @@ func NewProcessor(srcFile string, dstFile string, fields []int,
|
|||
return nil, fmt.Errorf("open target file `%s`: %w", dstFile, err)
|
||||
}
|
||||
|
||||
return INewProcessor(
|
||||
srcHandle, dstHandle, fields,
|
||||
clearBasedIndex, clearValueEqual, clearValueExcept,
|
||||
), nil
|
||||
}
|
||||
|
||||
func INewProcessor(srcReader io.ReadCloser, dstWriter io.WriteCloser, fields []int,
|
||||
clearBasedIndex int, clearValueEqual string, clearValueExcept string) *Processor {
|
||||
return &Processor{
|
||||
srcHandle: srcHandle,
|
||||
dstHandle: dstHandle,
|
||||
srcReader: srcReader,
|
||||
dstWriter: dstWriter,
|
||||
|
||||
// clear
|
||||
clearBasedIndex: clearBasedIndex,
|
||||
|
|
@ -59,14 +68,14 @@ func NewProcessor(srcFile string, dstFile string, fields []int,
|
|||
|
||||
segments: []*Segment{},
|
||||
rgCache: NewRegionCache(),
|
||||
}, nil
|
||||
}
|
||||
}
|
||||
|
||||
func (p *Processor) loadSegments() error {
|
||||
slog.Info("try to load the segments ... ")
|
||||
var tStart = time.Now()
|
||||
|
||||
_, mergeCount, iErr := IterateSegments(p.srcHandle, true, func(l string) {
|
||||
_, mergeCount, iErr := IterateSegments(p.srcReader, true, func(l string) {
|
||||
slog.Debug("loaded", "segment", l)
|
||||
}, func(region string) (string, error) {
|
||||
if p.clearBasedIndex > -1 {
|
||||
|
|
@ -135,7 +144,7 @@ func (p *Processor) Start() error {
|
|||
|
||||
slog.Info("try to write all segments to target file ...")
|
||||
for _, seg := range p.segments {
|
||||
_, err := fmt.Fprintln(p.dstHandle, seg.String())
|
||||
_, err := fmt.Fprintln(p.dstWriter, seg.String())
|
||||
if err != nil {
|
||||
return fmt.Errorf("write segment index for '%s': %w", seg.String(), err)
|
||||
}
|
||||
|
|
@ -146,12 +155,12 @@ func (p *Processor) Start() error {
|
|||
}
|
||||
|
||||
func (p *Processor) End() error {
|
||||
err := p.dstHandle.Close()
|
||||
err := p.dstWriter.Close()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
err = p.srcHandle.Close()
|
||||
err = p.srcReader.Close()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
|
|
|||
|
|
@ -13,13 +13,14 @@ package xdb
|
|||
import (
|
||||
"encoding/binary"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
)
|
||||
|
||||
type Searcher struct {
|
||||
version *Version
|
||||
|
||||
handle *os.File
|
||||
handle io.ReadSeekCloser
|
||||
|
||||
// header info
|
||||
header []byte
|
||||
|
|
@ -36,14 +37,16 @@ func NewSearcher(version *Version, dbFile string) (*Searcher, error) {
|
|||
return nil, err
|
||||
}
|
||||
|
||||
return INewSearcher(version, handle), nil
|
||||
}
|
||||
|
||||
func INewSearcher(version *Version, handle io.ReadSeekCloser) *Searcher {
|
||||
return &Searcher{
|
||||
version: version,
|
||||
|
||||
handle: handle,
|
||||
header: nil,
|
||||
|
||||
version: version,
|
||||
handle: handle,
|
||||
header: nil,
|
||||
vectorIndex: nil,
|
||||
}, nil
|
||||
}
|
||||
}
|
||||
|
||||
func (s *Searcher) Close() {
|
||||
|
|
|
|||
|
|
@ -8,10 +8,10 @@ import (
|
|||
"bufio"
|
||||
"bytes"
|
||||
"fmt"
|
||||
"io"
|
||||
"math/big"
|
||||
"net"
|
||||
"net/netip"
|
||||
"os"
|
||||
"strings"
|
||||
)
|
||||
|
||||
|
|
@ -180,10 +180,10 @@ func IPMiddle(sip, eip []byte) ([]byte, error) {
|
|||
return IPHalf(buf), nil
|
||||
}
|
||||
|
||||
func IterateSegments(handle *os.File, autoMerge bool, before func(l string), filter func(region string) (string, error), cRegion func(string) *Region, done func(seg *Segment) error) (int, int, error) {
|
||||
func IterateSegments(reader io.Reader, autoMerge bool, before func(l string), filter func(region string) (string, error), cRegion func(string) *Region, done func(seg *Segment) error) (int, int, error) {
|
||||
var last *Segment = nil
|
||||
var totalCount, mergeCount = 0, 0
|
||||
var scanner = bufio.NewScanner(handle)
|
||||
var scanner = bufio.NewScanner(reader)
|
||||
scanner.Split(bufio.ScanLines)
|
||||
for scanner.Scan() {
|
||||
var l = strings.TrimSpace(strings.TrimSuffix(scanner.Text(), "\n"))
|
||||
|
|
|
|||
Loading…
Reference in New Issue