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