seaweedfs/weed/storage/store_ec.go

408 lines
13 KiB
Go
Raw Permalink Normal View History

package storage
import (
"context"
2019-05-26 06:23:19 +00:00
"fmt"
"github.com/seaweedfs/seaweedfs/weed/pb"
"golang.org/x/exp/slices"
"io"
"os"
2019-05-29 06:48:39 +00:00
"sync"
2019-05-28 05:54:58 +00:00
"time"
2019-05-26 06:23:19 +00:00
2019-12-23 20:48:20 +00:00
"github.com/klauspost/reedsolomon"
"github.com/seaweedfs/seaweedfs/weed/glog"
"github.com/seaweedfs/seaweedfs/weed/operation"
"github.com/seaweedfs/seaweedfs/weed/pb/master_pb"
"github.com/seaweedfs/seaweedfs/weed/pb/volume_server_pb"
"github.com/seaweedfs/seaweedfs/weed/stats"
"github.com/seaweedfs/seaweedfs/weed/storage/erasure_coding"
"github.com/seaweedfs/seaweedfs/weed/storage/needle"
"github.com/seaweedfs/seaweedfs/weed/storage/types"
)
func (s *Store) CollectErasureCodingHeartbeat() *master_pb.Heartbeat {
var ecShardMessages []*master_pb.VolumeEcShardInformationMessage
collectionEcShardSize := make(map[string]int64)
for _, location := range s.Locations {
2019-05-28 04:40:51 +00:00
location.ecVolumesLock.RLock()
for _, ecShards := range location.ecVolumes {
ecShardMessages = append(ecShardMessages, ecShards.ToVolumeEcShardInformationMessage()...)
2019-06-16 09:44:20 +00:00
for _, ecShard := range ecShards.Shards {
collectionEcShardSize[ecShards.Collection] += ecShard.Size()
2019-06-16 09:44:20 +00:00
}
}
2019-05-28 04:40:51 +00:00
location.ecVolumesLock.RUnlock()
}
for col, size := range collectionEcShardSize {
stats.VolumeServerDiskSizeGauge.WithLabelValues(col, "ec").Set(float64(size))
}
2019-06-16 09:44:20 +00:00
return &master_pb.Heartbeat{
2019-06-05 20:32:33 +00:00
EcShards: ecShardMessages,
HasNoEcShards: len(ecShardMessages) == 0,
}
}
2019-05-26 06:23:19 +00:00
func (s *Store) MountEcShards(collection string, vid needle.VolumeId, shardId erasure_coding.ShardId) error {
for _, location := range s.Locations {
if err := location.LoadEcShard(collection, vid, shardId); err == nil {
glog.V(0).Infof("MountEcShards %d.%d", vid, shardId)
var shardBits erasure_coding.ShardBits
s.NewEcShardsChan <- master_pb.VolumeEcShardInformationMessage{
Id: uint32(vid),
Collection: collection,
EcIndexBits: uint32(shardBits.AddShardId(shardId)),
2021-02-16 10:47:02 +00:00
DiskType: string(location.DiskType),
2019-05-26 06:23:19 +00:00
}
return nil
} else if err == os.ErrNotExist {
continue
2019-06-03 09:26:31 +00:00
} else {
return fmt.Errorf("%s load ec shard %d.%d: %v", location.Directory, vid, shardId, err)
2019-05-26 06:23:19 +00:00
}
}
return fmt.Errorf("MountEcShards %d.%d not found on disk", vid, shardId)
}
func (s *Store) UnmountEcShards(vid needle.VolumeId, shardId erasure_coding.ShardId) error {
ecShard, found := s.findEcShard(vid, shardId)
if !found {
return nil
}
var shardBits erasure_coding.ShardBits
message := master_pb.VolumeEcShardInformationMessage{
Id: uint32(vid),
Collection: ecShard.Collection,
EcIndexBits: uint32(shardBits.AddShardId(shardId)),
2021-02-16 10:47:02 +00:00
DiskType: string(ecShard.DiskType),
2019-05-26 06:23:19 +00:00
}
for _, location := range s.Locations {
if deleted := location.UnloadEcShard(vid, shardId); deleted {
glog.V(0).Infof("UnmountEcShards %d.%d", vid, shardId)
s.DeletedEcShardsChan <- message
return nil
}
}
return fmt.Errorf("UnmountEcShards %d.%d not found on disk", vid, shardId)
}
func (s *Store) findEcShard(vid needle.VolumeId, shardId erasure_coding.ShardId) (*erasure_coding.EcVolumeShard, bool) {
2019-05-26 06:23:19 +00:00
for _, location := range s.Locations {
if v, found := location.FindEcShard(vid, shardId); found {
return v, found
}
}
return nil, false
}
2019-05-28 04:40:51 +00:00
func (s *Store) FindEcVolume(vid needle.VolumeId) (*erasure_coding.EcVolume, bool) {
for _, location := range s.Locations {
2019-05-28 04:40:51 +00:00
if s, found := location.FindEcVolume(vid); found {
return s, true
}
}
return nil, false
}
// shardFiles is a list of shard files, which is used to return the shard locations
func (s *Store) CollectEcShards(vid needle.VolumeId, shardFileNames []string) (ecVolume *erasure_coding.EcVolume, found bool) {
for _, location := range s.Locations {
if s, foundShards := location.CollectEcShards(vid, shardFileNames); foundShards {
ecVolume = s
found = true
}
}
return
}
2019-06-01 08:51:28 +00:00
func (s *Store) DestroyEcVolume(vid needle.VolumeId) {
for _, location := range s.Locations {
location.DestroyEcVolume(vid)
}
}
func (s *Store) ReadEcShardNeedle(vid needle.VolumeId, n *needle.Needle, onReadSizeFn func(size types.Size)) (int, error) {
for _, location := range s.Locations {
2019-05-28 04:40:51 +00:00
if localEcVolume, found := location.FindEcVolume(vid); found {
2019-05-27 18:59:03 +00:00
2019-12-28 20:44:59 +00:00
offset, size, intervals, err := localEcVolume.LocateEcShardNeedle(n.Id, localEcVolume.Version)
2019-05-27 18:59:03 +00:00
if err != nil {
2019-06-06 06:20:26 +00:00
return 0, fmt.Errorf("locate in local ec volume: %v", err)
2019-05-27 18:59:03 +00:00
}
2020-08-19 00:39:29 +00:00
if size.IsDeleted() {
return 0, ErrorDeleted
2019-06-20 05:57:14 +00:00
}
2019-05-27 18:59:03 +00:00
if onReadSizeFn != nil {
onReadSizeFn(size)
}
glog.V(3).Infof("read ec volume %d offset %d size %d intervals:%+v", vid, offset.ToActualOffset(), size, intervals)
2019-05-29 04:29:07 +00:00
if len(intervals) > 1 {
2019-06-21 08:14:10 +00:00
glog.V(3).Infof("ReadEcShardNeedle needle id %s intervals:%+v", n.String(), intervals)
2019-05-29 04:29:07 +00:00
}
bytes, isDeleted, err := s.readEcShardIntervals(vid, n.Id, localEcVolume, intervals)
2019-05-27 18:59:03 +00:00
if err != nil {
return 0, fmt.Errorf("ReadEcShardIntervals: %v", err)
}
2019-06-21 08:14:10 +00:00
if isDeleted {
return 0, ErrorDeleted
2019-06-21 08:14:10 +00:00
}
2019-05-27 18:59:03 +00:00
err = n.ReadBytes(bytes, offset.ToActualOffset(), size, localEcVolume.Version)
2019-05-27 18:59:03 +00:00
if err != nil {
return 0, fmt.Errorf("readbytes: %v", err)
}
return len(bytes), nil
}
}
return 0, fmt.Errorf("ec shard %d not found", vid)
}
2019-05-27 18:59:03 +00:00
func (s *Store) readEcShardIntervals(vid needle.VolumeId, needleId types.NeedleId, ecVolume *erasure_coding.EcVolume, intervals []erasure_coding.Interval) (data []byte, is_deleted bool, err error) {
2019-05-28 05:54:58 +00:00
if err = s.cachedLookupEcShardLocations(ecVolume); err != nil {
2019-06-21 08:14:10 +00:00
return nil, false, fmt.Errorf("failed to locate shard via master grpc %s: %v", s.MasterAddress, err)
}
2019-05-27 18:59:03 +00:00
for i, interval := range intervals {
if d, isDeleted, e := s.readOneEcShardInterval(needleId, ecVolume, interval); e != nil {
2019-06-21 08:14:10 +00:00
return nil, isDeleted, e
2019-05-27 18:59:03 +00:00
} else {
2019-06-21 08:14:10 +00:00
if isDeleted {
is_deleted = true
}
2019-05-27 18:59:03 +00:00
if i == 0 {
data = d
} else {
data = append(data, d...)
}
}
}
return
}
func (s *Store) readOneEcShardInterval(needleId types.NeedleId, ecVolume *erasure_coding.EcVolume, interval erasure_coding.Interval) (data []byte, is_deleted bool, err error) {
2019-05-27 18:59:03 +00:00
shardId, actualOffset := interval.ToShardIdAndOffset(erasure_coding.ErasureCodingLargeBlockSize, erasure_coding.ErasureCodingSmallBlockSize)
2019-05-29 04:29:07 +00:00
data = make([]byte, interval.Size)
2019-05-28 04:40:51 +00:00
if shard, found := ecVolume.FindEcVolumeShard(shardId); found {
2019-05-27 18:59:03 +00:00
if _, err = shard.ReadAt(data, actualOffset); err != nil {
2020-06-25 05:07:53 +00:00
glog.V(0).Infof("read local ec shard %d.%d offset %d: %v", ecVolume.VolumeId, shardId, actualOffset, err)
2019-05-27 18:59:03 +00:00
return
}
} else {
2019-05-28 05:54:58 +00:00
ecVolume.ShardLocationsLock.RLock()
2019-05-31 07:19:13 +00:00
sourceDataNodes, hasShardIdLocation := ecVolume.ShardLocations[shardId]
2019-05-28 05:54:58 +00:00
ecVolume.ShardLocationsLock.RUnlock()
2019-05-29 06:48:39 +00:00
// try reading directly
2019-05-31 07:19:13 +00:00
if hasShardIdLocation {
_, is_deleted, err = s.readRemoteEcShardInterval(sourceDataNodes, needleId, ecVolume.VolumeId, shardId, data, actualOffset)
2019-05-31 07:19:13 +00:00
if err == nil {
return
}
glog.V(0).Infof("clearing ec shard %d.%d locations: %v", ecVolume.VolumeId, shardId, err)
2019-05-29 06:48:39 +00:00
}
// try reading by recovering from other shards
_, is_deleted, err = s.recoverOneRemoteEcShardInterval(needleId, ecVolume, shardId, data, actualOffset)
2019-05-29 06:48:39 +00:00
if err == nil {
return
}
2019-05-29 06:48:39 +00:00
glog.V(0).Infof("recover ec shard %d.%d : %v", ecVolume.VolumeId, shardId, err)
2019-05-27 18:59:03 +00:00
}
return
}
2019-05-31 07:19:13 +00:00
func forgetShardId(ecVolume *erasure_coding.EcVolume, shardId erasure_coding.ShardId) {
// failed to access the source data nodes, clear it up
ecVolume.ShardLocationsLock.Lock()
delete(ecVolume.ShardLocations, shardId)
ecVolume.ShardLocationsLock.Unlock()
}
func (s *Store) cachedLookupEcShardLocations(ecVolume *erasure_coding.EcVolume) (err error) {
2019-05-28 05:54:58 +00:00
shardCount := len(ecVolume.ShardLocations)
if shardCount < erasure_coding.DataShardsCount &&
2019-06-23 03:05:25 +00:00
ecVolume.ShardLocationsRefreshTime.Add(11*time.Second).After(time.Now()) ||
shardCount == erasure_coding.TotalShardsCount &&
2019-06-23 03:05:25 +00:00
ecVolume.ShardLocationsRefreshTime.Add(37*time.Minute).After(time.Now()) ||
shardCount >= erasure_coding.DataShardsCount &&
2019-06-23 03:05:25 +00:00
ecVolume.ShardLocationsRefreshTime.Add(7*time.Minute).After(time.Now()) {
2019-05-28 05:54:58 +00:00
// still fresh
return nil
}
glog.V(3).Infof("lookup and cache ec volume %d locations", ecVolume.VolumeId)
2019-05-28 05:54:58 +00:00
err = operation.WithMasterServerClient(false, s.MasterAddress, s.grpcDialOption, func(masterClient master_pb.SeaweedClient) error {
req := &master_pb.LookupEcVolumeRequest{
VolumeId: uint32(ecVolume.VolumeId),
}
resp, err := masterClient.LookupEcVolume(context.Background(), req)
if err != nil {
return fmt.Errorf("lookup ec volume %d: %v", ecVolume.VolumeId, err)
}
2019-05-31 07:58:51 +00:00
if len(resp.ShardIdLocations) < erasure_coding.DataShardsCount {
return fmt.Errorf("only %d shards found but %d required", len(resp.ShardIdLocations), erasure_coding.DataShardsCount)
}
ecVolume.ShardLocationsLock.Lock()
for _, shardIdLocations := range resp.ShardIdLocations {
shardId := erasure_coding.ShardId(shardIdLocations.ShardId)
delete(ecVolume.ShardLocations, shardId)
for _, loc := range shardIdLocations.Locations {
ecVolume.ShardLocations[shardId] = append(ecVolume.ShardLocations[shardId], pb.NewServerAddressFromLocation(loc))
}
}
2019-05-29 04:29:07 +00:00
ecVolume.ShardLocationsRefreshTime = time.Now()
ecVolume.ShardLocationsLock.Unlock()
2019-05-28 05:54:58 +00:00
return nil
})
return
}
func (s *Store) readRemoteEcShardInterval(sourceDataNodes []pb.ServerAddress, needleId types.NeedleId, vid needle.VolumeId, shardId erasure_coding.ShardId, buf []byte, offset int64) (n int, is_deleted bool, err error) {
2019-05-29 06:48:39 +00:00
if len(sourceDataNodes) == 0 {
2019-06-21 08:14:10 +00:00
return 0, false, fmt.Errorf("failed to find ec shard %d.%d", vid, shardId)
}
2019-05-29 06:48:39 +00:00
for _, sourceDataNode := range sourceDataNodes {
2019-11-09 06:40:55 +00:00
glog.V(3).Infof("read remote ec shard %d.%d from %s", vid, shardId, sourceDataNode)
n, is_deleted, err = s.doReadRemoteEcShardInterval(sourceDataNode, needleId, vid, shardId, buf, offset)
2019-05-29 06:48:39 +00:00
if err == nil {
return
}
glog.V(1).Infof("read remote ec shard %d.%d from %s: %v", vid, shardId, sourceDataNode, err)
}
return
}
func (s *Store) doReadRemoteEcShardInterval(sourceDataNode pb.ServerAddress, needleId types.NeedleId, vid needle.VolumeId, shardId erasure_coding.ShardId, buf []byte, offset int64) (n int, is_deleted bool, err error) {
err = operation.WithVolumeServerClient(false, sourceDataNode, s.grpcDialOption, func(client volume_server_pb.VolumeServerClient) error {
// copy data slice
shardReadClient, err := client.VolumeEcShardRead(context.Background(), &volume_server_pb.VolumeEcShardReadRequest{
VolumeId: uint32(vid),
ShardId: uint32(shardId),
Offset: offset,
Size: int64(len(buf)),
2019-06-21 08:14:10 +00:00
FileKey: uint64(needleId),
})
if err != nil {
return fmt.Errorf("failed to start reading ec shard %d.%d from %s: %v", vid, shardId, sourceDataNode, err)
}
2019-05-27 18:59:03 +00:00
for {
resp, receiveErr := shardReadClient.Recv()
if receiveErr == io.EOF {
break
}
if receiveErr != nil {
return fmt.Errorf("receiving ec shard %d.%d from %s: %v", vid, shardId, sourceDataNode, receiveErr)
}
2019-06-21 08:14:10 +00:00
if resp.IsDeleted {
is_deleted = true
}
copy(buf[n:n+len(resp.Data)], resp.Data)
n += len(resp.Data)
}
2019-05-27 18:59:03 +00:00
return nil
})
if err != nil {
2019-06-21 08:14:10 +00:00
return 0, is_deleted, fmt.Errorf("read ec shard %d.%d from %s: %v", vid, shardId, sourceDataNode, err)
}
2019-05-27 18:59:03 +00:00
return
}
func (s *Store) recoverOneRemoteEcShardInterval(needleId types.NeedleId, ecVolume *erasure_coding.EcVolume, shardIdToRecover erasure_coding.ShardId, buf []byte, offset int64) (n int, is_deleted bool, err error) {
2019-11-09 06:40:55 +00:00
glog.V(3).Infof("recover ec shard %d.%d from other locations", ecVolume.VolumeId, shardIdToRecover)
2019-05-29 06:48:39 +00:00
enc, err := reedsolomon.New(erasure_coding.DataShardsCount, erasure_coding.ParityShardsCount)
if err != nil {
2019-06-21 08:14:10 +00:00
return 0, false, fmt.Errorf("failed to create encoder: %v", err)
2019-05-29 06:48:39 +00:00
}
bufs := make([][]byte, erasure_coding.TotalShardsCount)
var wg sync.WaitGroup
ecVolume.ShardLocationsLock.RLock()
for shardId, locations := range ecVolume.ShardLocations {
// skip current shard or empty shard
2019-05-29 06:48:39 +00:00
if shardId == shardIdToRecover {
continue
}
if len(locations) == 0 {
glog.V(3).Infof("readRemoteEcShardInterval missing %d.%d from %+v", ecVolume.VolumeId, shardId, locations)
continue
}
// read from remote locations
wg.Add(1)
go func(shardId erasure_coding.ShardId, locations []pb.ServerAddress) {
2019-05-29 06:48:39 +00:00
defer wg.Done()
data := make([]byte, len(buf))
nRead, isDeleted, readErr := s.readRemoteEcShardInterval(locations, needleId, ecVolume.VolumeId, shardId, data, offset)
if readErr != nil {
2019-05-31 07:19:13 +00:00
glog.V(3).Infof("recover: readRemoteEcShardInterval %d.%d %d bytes from %+v: %v", ecVolume.VolumeId, shardId, nRead, locations, readErr)
forgetShardId(ecVolume, shardId)
2019-05-29 06:48:39 +00:00
}
2019-06-21 08:14:10 +00:00
if isDeleted {
is_deleted = true
}
if nRead == len(buf) {
2019-05-29 06:48:39 +00:00
bufs[shardId] = data
}
}(shardId, locations)
}
ecVolume.ShardLocationsLock.RUnlock()
wg.Wait()
if err = enc.ReconstructData(bufs); err != nil {
glog.V(3).Infof("recovered ec shard %d.%d failed: %v", ecVolume.VolumeId, shardIdToRecover, err)
2019-06-21 08:14:10 +00:00
return 0, false, err
2019-05-29 06:48:39 +00:00
}
2019-05-31 07:19:13 +00:00
glog.V(4).Infof("recovered ec shard %d.%d from other locations", ecVolume.VolumeId, shardIdToRecover)
2019-05-29 06:48:39 +00:00
copy(buf, bufs[shardIdToRecover])
2019-06-21 08:14:10 +00:00
return len(buf), is_deleted, nil
}
2019-06-05 04:52:37 +00:00
func (s *Store) EcVolumes() (ecVolumes []*erasure_coding.EcVolume) {
for _, location := range s.Locations {
location.ecVolumesLock.RLock()
for _, v := range location.ecVolumes {
ecVolumes = append(ecVolumes, v)
}
location.ecVolumesLock.RUnlock()
}
slices.SortFunc(ecVolumes, func(a, b *erasure_coding.EcVolume) bool {
return a.VolumeId > b.VolumeId
2019-06-05 04:52:37 +00:00
})
return ecVolumes
}