mirror of
https://github.com/cubefs/cubefs.git
synced 2026-08-02 02:00:56 +00:00
feat(objectnode): read/write cold volume buffer pool
Signed-off-by: tangdeyi <tangdeyi@oppo.com>
This commit is contained in:
parent
b036d1fa7c
commit
0bb8a4b6ea
@ -13,7 +13,5 @@
|
||||
"logDir": "/cfs/log",
|
||||
"walDir": "/cfs/data/wal",
|
||||
"storeDir": "/cfs/data/store",
|
||||
"metaNodeReservedMem": "67108864",
|
||||
"ebsAddr": "10.177.40.215:8500",
|
||||
"ebsServicePath": "access"
|
||||
"metaNodeReservedMem": "67108864"
|
||||
}
|
||||
|
||||
@ -13,7 +13,5 @@
|
||||
"logDir": "/cfs/log",
|
||||
"walDir": "/cfs/data/wal",
|
||||
"storeDir": "/cfs/data/store",
|
||||
"metaNodeReservedMem": "67108864",
|
||||
"ebsAddr": "10.177.40.215:8500",
|
||||
"ebsServicePath": "access"
|
||||
"metaNodeReservedMem": "67108864"
|
||||
}
|
||||
|
||||
@ -13,7 +13,5 @@
|
||||
"logDir": "/cfs/log",
|
||||
"walDir": "/cfs/data/wal",
|
||||
"storeDir": "/cfs/data/store",
|
||||
"metaNodeReservedMem": "67108864",
|
||||
"ebsAddr": "10.177.40.215:8500",
|
||||
"ebsServicePath": "access"
|
||||
"metaNodeReservedMem": "67108864"
|
||||
}
|
||||
|
||||
@ -21,6 +21,7 @@ import (
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"net/http"
|
||||
"regexp"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
@ -29,6 +30,13 @@ import (
|
||||
"github.com/gorilla/mux"
|
||||
)
|
||||
|
||||
const (
|
||||
DefaultMinBucketLength = 3
|
||||
DefaultMaxBucketLength = 63
|
||||
)
|
||||
|
||||
var regexBucketName = regexp.MustCompile(`^[0-9a-z][-0-9a-z]+[0-9a-z]$`)
|
||||
|
||||
// Head bucket
|
||||
// API reference: https://docs.aws.amazon.com/AmazonS3/latest/API/API_HeadBucket.html
|
||||
func (o *ObjectNode) headBucketHandler(w http.ResponseWriter, r *http.Request) {
|
||||
@ -53,6 +61,11 @@ func (o *ObjectNode) createBucketHandler(w http.ResponseWriter, r *http.Request)
|
||||
return
|
||||
}
|
||||
|
||||
if !IsValidBucketName(bucket, DefaultMinBucketLength, DefaultMaxBucketLength) {
|
||||
errorCode = InvalidBucketName
|
||||
return
|
||||
}
|
||||
|
||||
if vol, _ := o.vm.VolumeWithoutBlacklist(bucket); vol != nil {
|
||||
log.LogInfof("createBucketHandler: duplicated bucket name: requestID(%v) bucket(%v)", GetRequestID(r), bucket)
|
||||
errorCode = BucketAlreadyOwnedByYou
|
||||
@ -104,7 +117,7 @@ func (o *ObjectNode) createBucketHandler(w http.ResponseWriter, r *http.Request)
|
||||
GetRequestID(r), bucket, auth.accessKey, err)
|
||||
return
|
||||
}
|
||||
w.Header()[HeaderNameLocation] = []string{bucket}
|
||||
w.Header().Set(HeaderNameLocation, "/"+bucket)
|
||||
|
||||
vol, err1 := o.vm.VolumeWithoutBlacklist(bucket)
|
||||
if err1 != nil {
|
||||
@ -432,3 +445,13 @@ func (o *ObjectNode) getUserInfoByAccessKeyV2(accessKey string) (userInfo *proto
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
func IsValidBucketName(bucketName string, minBucketLength, maxBucketLength int) bool {
|
||||
if len(bucketName) < minBucketLength || len(bucketName) > maxBucketLength {
|
||||
return false
|
||||
}
|
||||
if !regexBucketName.MatchString(bucketName) {
|
||||
return false
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
@ -220,7 +220,7 @@ func (o *ObjectNode) uploadPartHandler(w http.ResponseWriter, r *http.Request) {
|
||||
errorCode = NoSuchUpload
|
||||
return
|
||||
}
|
||||
if err == syscall.EAGAIN {
|
||||
if err == syscall.EEXIST {
|
||||
errorCode = ConflictUploadRequest
|
||||
return
|
||||
}
|
||||
|
||||
@ -38,6 +38,7 @@ import (
|
||||
"github.com/cubefs/cubefs/sdk/master"
|
||||
"github.com/cubefs/cubefs/sdk/meta"
|
||||
"github.com/cubefs/cubefs/util"
|
||||
"github.com/cubefs/cubefs/util/buf"
|
||||
"github.com/cubefs/cubefs/util/exporter"
|
||||
"github.com/cubefs/cubefs/util/log"
|
||||
)
|
||||
@ -1504,7 +1505,8 @@ func (v *Volume) readEbs(inode, inodeSize uint64, path string, writer io.Writer,
|
||||
reader := v.getEbsReader(inode)
|
||||
var n int
|
||||
var rest uint64
|
||||
var tmp = make([]byte, 2*v.ebsBlockSize)
|
||||
var tmp = buf.ReadBufPool.Get().([]byte)
|
||||
defer buf.ReadBufPool.Put(tmp)
|
||||
|
||||
for {
|
||||
if rest = upper - offset; rest <= 0 {
|
||||
|
||||
@ -125,8 +125,8 @@ var (
|
||||
objMetaCache *ObjMetaCache
|
||||
blockCache *bcache.BcacheClient
|
||||
ebsClient *blobstore.BlobStoreClient
|
||||
writeThreads int = 4
|
||||
readThreads int = 4
|
||||
writeThreads = 4
|
||||
readThreads = 4
|
||||
enableBlockcache bool
|
||||
)
|
||||
|
||||
@ -261,26 +261,18 @@ func handleStart(s common.Server, cfg *config.Config) (err error) {
|
||||
return
|
||||
}
|
||||
// Get cluster info from master
|
||||
|
||||
var ci *proto.ClusterInfo
|
||||
if ci, err = o.mc.AdminAPI().GetClusterInfo(); err != nil {
|
||||
return
|
||||
}
|
||||
o.updateRegion(ci.Cluster)
|
||||
log.LogInfof("handleStart: get cluster information: region(%v)", o.region)
|
||||
ebsClient, err = blobstore.NewEbsClient(access.Config{
|
||||
ConnMode: access.NoLimitConnMode,
|
||||
Consul: access.ConsulConfig{
|
||||
Address: ci.EbsAddr,
|
||||
},
|
||||
//ServicePath: ci.ServicePath,
|
||||
MaxSizePutOnce: MaxSizePutOnce,
|
||||
Logger: &access.Logger{
|
||||
Filename: path.Join(cfg.GetString("logDir"), "ebs.log"),
|
||||
},
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
if ci.EbsAddr != "" {
|
||||
err = newEbsClient(ci, cfg)
|
||||
if err != nil {
|
||||
log.LogWarnf("handleStart: new ebsClient err(%v)", err)
|
||||
return err
|
||||
}
|
||||
wt := cfg.GetInt(ebsWriteThreads)
|
||||
if wt != 0 {
|
||||
writeThreads = wt
|
||||
@ -304,6 +296,20 @@ func handleStart(s common.Server, cfg *config.Config) (err error) {
|
||||
return
|
||||
}
|
||||
|
||||
func newEbsClient(ci *proto.ClusterInfo, cfg *config.Config) (err error) {
|
||||
ebsClient, err = blobstore.NewEbsClient(access.Config{
|
||||
ConnMode: access.NoLimitConnMode,
|
||||
Consul: access.ConsulConfig{
|
||||
Address: ci.EbsAddr,
|
||||
},
|
||||
MaxSizePutOnce: MaxSizePutOnce,
|
||||
Logger: &access.Logger{
|
||||
Filename: path.Join(cfg.GetString("logDir"), "ebs.log"),
|
||||
},
|
||||
})
|
||||
return err
|
||||
}
|
||||
|
||||
func handleShutdown(s common.Server) {
|
||||
o, ok := s.(*ObjectNode)
|
||||
if !ok {
|
||||
|
||||
@ -207,10 +207,11 @@ func (writer *Writer) cacheLevel2(wSlice *rwSlice) {
|
||||
|
||||
func (writer *Writer) WriteFromReader(ctx context.Context, reader io.Reader, h hash.Hash) (size uint64, err error) {
|
||||
var (
|
||||
buf = make([]byte, 2*writer.blockSize)
|
||||
tmp = buf.ReadBufPool.Get().([]byte)
|
||||
exec = NewExecutor(writer.wConcurrency)
|
||||
leftToWrite int
|
||||
)
|
||||
defer buf.ReadBufPool.Put(tmp)
|
||||
|
||||
writer.fileOffset = 0
|
||||
writer.err = make(chan *wSliceErr)
|
||||
@ -268,7 +269,7 @@ func (writer *Writer) WriteFromReader(ctx context.Context, reader io.Reader, h h
|
||||
LOOP:
|
||||
for {
|
||||
position := 0
|
||||
leftToWrite, err = reader.Read(buf)
|
||||
leftToWrite, err = reader.Read(tmp)
|
||||
if err != nil && err != io.EOF {
|
||||
return
|
||||
}
|
||||
@ -284,7 +285,7 @@ LOOP:
|
||||
|
||||
freeSize := writer.blockSize - len(writer.buf)
|
||||
writeSize := util.Min(leftToWrite, freeSize)
|
||||
writer.buf = append(writer.buf, buf[position:position+writeSize]...)
|
||||
writer.buf = append(writer.buf, tmp[position:position+writeSize]...)
|
||||
position += writeSize
|
||||
leftToWrite -= writeSize
|
||||
writer.fileOffset += writeSize
|
||||
|
||||
@ -360,7 +360,7 @@ func statusToErrno(status int) error {
|
||||
case statusTxTimeout:
|
||||
return syscall.EAGAIN
|
||||
case statusUploadPartConflict:
|
||||
return syscall.EAGAIN
|
||||
return syscall.EEXIST
|
||||
default:
|
||||
}
|
||||
return syscall.EIO
|
||||
|
||||
@ -15,6 +15,13 @@ const (
|
||||
InvalidLimit = 0
|
||||
)
|
||||
|
||||
var ReadBufPool = sync.Pool{
|
||||
New: func() interface{} {
|
||||
b := make([]byte, 32*1024)
|
||||
return b
|
||||
},
|
||||
}
|
||||
|
||||
var tinyBuffersTotalLimit int64 = 4096
|
||||
var NormalBuffersTotalLimit int64
|
||||
var HeadBuffersTotalLimit int64
|
||||
|
||||
Loading…
Reference in New Issue
Block a user