From 554aa13c5b9c2313f277463034c16b1705e23a3f Mon Sep 17 00:00:00 2001 From: Bruce Zou Date: Fri, 29 May 2026 21:44:52 +0800 Subject: [PATCH] expose whether a volume replica is backed by remote storage Volume locations returned by lookups do not indicate whether a replica has been tiered to remote storage. Readers cannot distinguish a local replica from a remote-backed one, so they may hit a remote-backed replica first even when a local replica is available. Add DataInRemote to the lookup location message, populate it from the master's volume info, and carry it through the wdclient vid map so clients can prefer local replicas when resolving chunk locations. --- other/java/client/src/main/proto/filer.proto | 2 + seaweed-volume/proto/master.proto | 2 + weed/filer/reader_at.go | 10 +++++ weed/operation/lookup.go | 35 ++++++++++++----- weed/pb/filer.proto | 1 + weed/pb/filer_pb/filer.pb.go | 13 ++++++- weed/pb/filer_pb/filer_vtproto.pb.go | 33 ++++++++++++++++ weed/pb/master.proto | 2 + weed/pb/master_pb/master.pb.go | 28 ++++++++++++-- weed/server/filer_grpc_server.go | 9 +++-- weed/server/master_grpc_server.go | 9 ++++- weed/server/master_grpc_server_volume.go | 9 +++-- weed/server/master_server_handlers.go | 25 +++++++----- weed/topology/topology_info.go | 10 ++++- weed/util/slice.go | 15 ++++++++ weed/util/slice_test.go | 40 ++++++++++++++++++++ weed/wdclient/masterclient.go | 14 +++++-- weed/wdclient/vid_map.go | 25 +++++++++--- 18 files changed, 238 insertions(+), 44 deletions(-) create mode 100644 weed/util/slice_test.go diff --git a/other/java/client/src/main/proto/filer.proto b/other/java/client/src/main/proto/filer.proto index b8beefd71..ca1c0114c 100644 --- a/other/java/client/src/main/proto/filer.proto +++ b/other/java/client/src/main/proto/filer.proto @@ -576,6 +576,7 @@ message Location { string public_url = 2; uint32 grpc_port = 3; string data_center = 4; + bool data_in_remote = 5; } message LookupVolumeResponse { map locations_map = 1; @@ -662,6 +663,7 @@ message SubscribeMetadataResponse { int64 ts_ns = 3; repeated SubscribeMetadataResponse events = 4; // batch of additional events (backlog catch-up) repeated LogFileChunkRef log_file_refs = 5; // log file chunk refs for client direct-read + int64 flushed_ts_ns = 6; // local log-buffer flush watermark: everything at or below it is on disk } message ListMetadataSubscribersRequest { repeated string client_types = 1; // optional filter by client type, e.g. "mount"; empty = all diff --git a/seaweed-volume/proto/master.proto b/seaweed-volume/proto/master.proto index 2dfd9caa0..eec385e4b 100644 --- a/seaweed-volume/proto/master.proto +++ b/seaweed-volume/proto/master.proto @@ -222,6 +222,7 @@ message VolumeLocation { uint32 grpc_port = 7; repeated uint32 new_ec_vids = 8; repeated uint32 deleted_ec_vids = 9; + repeated uint32 remote_vids = 10; } message ClusterNodeUpdate { @@ -267,6 +268,7 @@ message Location { string public_url = 2; uint32 grpc_port = 3; string data_center = 4; + bool data_in_remote = 5; } message AssignRequest { diff --git a/weed/filer/reader_at.go b/weed/filer/reader_at.go index 55231a193..7e2e2918f 100644 --- a/weed/filer/reader_at.go +++ b/weed/filer/reader_at.go @@ -117,9 +117,15 @@ func LookupFn(filerClient filer_pb.FilerClient) wdclient.LookupFileIdFunctionTyp fcDataCenter := filerClient.GetDataCenter() var sameDcTargetUrls, otherTargetUrls []string + localUrls := make(map[string]bool) for _, loc := range locations.Locations { volumeServerAddress := filerClient.AdjustedUrl(loc) targetUrl := fmt.Sprintf("http://%s/%s", volumeServerAddress, fileId) + glog.V(4).Infof("lookup %s => %s, data in remote storage tier: %v", fileId, targetUrl, loc.DataInRemote) + + if !loc.DataInRemote { + localUrls[targetUrl] = true + } if fcDataCenter == "" || fcDataCenter != loc.DataCenter { otherTargetUrls = append(otherTargetUrls, targetUrl) } else { @@ -132,6 +138,10 @@ func LookupFn(filerClient filer_pb.FilerClient) wdclient.LookupFileIdFunctionTyp rand.Shuffle(len(otherTargetUrls), func(i, j int) { otherTargetUrls[i], otherTargetUrls[j] = otherTargetUrls[j], otherTargetUrls[i] }) + if len(localUrls) > 0 { + sameDcTargetUrls = util.ReorderToFront(localUrls, sameDcTargetUrls) + otherTargetUrls = util.ReorderToFront(localUrls, otherTargetUrls) + } // Prefer same data center targetUrls = append(sameDcTargetUrls, otherTargetUrls...) return diff --git a/weed/operation/lookup.go b/weed/operation/lookup.go index eabea475e..c8d83ea2d 100644 --- a/weed/operation/lookup.go +++ b/weed/operation/lookup.go @@ -15,10 +15,11 @@ import ( ) type Location struct { - Url string `json:"url,omitempty"` - PublicUrl string `json:"publicUrl,omitempty"` - DataCenter string `json:"dataCenter,omitempty"` - GrpcPort int `json:"grpcPort,omitempty"` + Url string `json:"url,omitempty"` + PublicUrl string `json:"publicUrl,omitempty"` + DataCenter string `json:"dataCenter,omitempty"` + GrpcPort int `json:"grpcPort,omitempty"` + DataInRemote bool `json:"dataInRemote,omitempty"` } func (l *Location) ServerAddress() pb.ServerAddress { @@ -42,6 +43,8 @@ var ( ) func LookupFileId(masterFn GetMasterFn, grpcDialOption grpc.DialOption, fileId string) (fullUrl string, jwt string, err error) { + var location string + parts := strings.Split(fileId, ",") if len(parts) != 2 { return "", jwt, errors.New("Invalid fileId " + fileId) @@ -53,7 +56,20 @@ func LookupFileId(masterFn GetMasterFn, grpcDialOption grpc.DialOption, fileId s if len(lookup.Locations) == 0 { return "", jwt, errors.New("File Not Found") } - return "http://" + lookup.Locations[rand.IntN(len(lookup.Locations))].Url + "/" + fileId, lookup.Jwt, nil + + localUrls := make([]string, 0, len(lookup.Locations)) + for _, loc := range lookup.Locations { + if !loc.DataInRemote { + localUrls = append(localUrls, loc.Url) + } + } + if len(localUrls) > 0 { + location = "http://" + localUrls[rand.IntN(len(localUrls))] + "/" + fileId + } else { + location = "http://" + lookup.Locations[rand.IntN(len(lookup.Locations))].Url + "/" + fileId + } + + return location, lookup.Jwt, nil } func LookupVolumeId(masterFn GetMasterFn, grpcDialOption grpc.DialOption, vid string) (*LookupResult, error) { @@ -102,10 +118,11 @@ func LookupVolumeIds(masterFn GetMasterFn, grpcDialOption grpc.DialOption, vids var locations []Location for _, loc := range vidLocations.Locations { locations = append(locations, Location{ - Url: loc.Url, - PublicUrl: loc.PublicUrl, - DataCenter: loc.DataCenter, - GrpcPort: int(loc.GrpcPort), + Url: loc.Url, + PublicUrl: loc.PublicUrl, + DataCenter: loc.DataCenter, + GrpcPort: int(loc.GrpcPort), + DataInRemote: loc.DataInRemote, }) } if vidLocations.Error == "" { diff --git a/weed/pb/filer.proto b/weed/pb/filer.proto index adef26743..dbd18dfaf 100644 --- a/weed/pb/filer.proto +++ b/weed/pb/filer.proto @@ -579,6 +579,7 @@ message Location { string public_url = 2; uint32 grpc_port = 3; string data_center = 4; + bool data_in_remote = 5; } message LookupVolumeResponse { map locations_map = 1; diff --git a/weed/pb/filer_pb/filer.pb.go b/weed/pb/filer_pb/filer.pb.go index b95a2cfe5..14e532d49 100644 --- a/weed/pb/filer_pb/filer.pb.go +++ b/weed/pb/filer_pb/filer.pb.go @@ -3433,6 +3433,7 @@ type Location struct { PublicUrl string `protobuf:"bytes,2,opt,name=public_url,json=publicUrl,proto3" json:"public_url,omitempty"` GrpcPort uint32 `protobuf:"varint,3,opt,name=grpc_port,json=grpcPort,proto3" json:"grpc_port,omitempty"` DataCenter string `protobuf:"bytes,4,opt,name=data_center,json=dataCenter,proto3" json:"data_center,omitempty"` + DataInRemote bool `protobuf:"varint,5,opt,name=data_in_remote,json=dataInRemote,proto3" json:"data_in_remote,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -3495,6 +3496,13 @@ func (x *Location) GetDataCenter() string { return "" } +func (x *Location) GetDataInRemote() bool { + if x != nil { + return x.DataInRemote + } + return false +} + type LookupVolumeResponse struct { state protoimpl.MessageState `protogen:"open.v1"` LocationsMap map[string]*Locations `protobuf:"bytes,1,rep,name=locations_map,json=locationsMap,proto3" json:"locations_map,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` @@ -7287,14 +7295,15 @@ const file_filer_proto_rawDesc = "" + "\n" + "volume_ids\x18\x01 \x03(\tR\tvolumeIds\"=\n" + "\tLocations\x120\n" + - "\tlocations\x18\x01 \x03(\v2\x12.filer_pb.LocationR\tlocations\"y\n" + + "\tlocations\x18\x01 \x03(\v2\x12.filer_pb.LocationR\tlocations\"\x9f\x01\n" + "\bLocation\x12\x10\n" + "\x03url\x18\x01 \x01(\tR\x03url\x12\x1d\n" + "\n" + "public_url\x18\x02 \x01(\tR\tpublicUrl\x12\x1b\n" + "\tgrpc_port\x18\x03 \x01(\rR\bgrpcPort\x12\x1f\n" + "\vdata_center\x18\x04 \x01(\tR\n" + - "dataCenter\"\xc3\x01\n" + + "dataCenter\x12$\n" + + "\x0edata_in_remote\x18\x05 \x01(\bR\fdataInRemote\"\xc3\x01\n" + "\x14LookupVolumeResponse\x12U\n" + "\rlocations_map\x18\x01 \x03(\v20.filer_pb.LookupVolumeResponse.LocationsMapEntryR\flocationsMap\x1aT\n" + "\x11LocationsMapEntry\x12\x10\n" + diff --git a/weed/pb/filer_pb/filer_vtproto.pb.go b/weed/pb/filer_pb/filer_vtproto.pb.go index 8ba5613e0..6d944c6b2 100644 --- a/weed/pb/filer_pb/filer_vtproto.pb.go +++ b/weed/pb/filer_pb/filer_vtproto.pb.go @@ -3136,6 +3136,16 @@ func (m *Location) MarshalToSizedBufferVT(dAtA []byte) (int, error) { i -= len(m.unknownFields) copy(dAtA[i:], m.unknownFields) } + if m.DataInRemote { + i-- + if m.DataInRemote { + dAtA[i] = 1 + } else { + dAtA[i] = 0 + } + i-- + dAtA[i] = 0x28 + } if len(m.DataCenter) > 0 { i -= len(m.DataCenter) copy(dAtA[i:], m.DataCenter) @@ -7494,6 +7504,9 @@ func (m *Location) SizeVT() (n int) { if l > 0 { n += 1 + l + protohelpers.SizeOfVarint(uint64(l)) } + if m.DataInRemote { + n += 2 + } n += len(m.unknownFields) return n } @@ -17370,6 +17383,26 @@ func (m *Location) UnmarshalVT(dAtA []byte) error { } m.DataCenter = string(dAtA[iNdEx:postIndex]) iNdEx = postIndex + case 5: + if wireType != 0 { + return fmt.Errorf("proto: wrong wireType = %d for field DataInRemote", wireType) + } + var v int + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return protohelpers.ErrIntOverflow + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + v |= int(b&0x7F) << shift + if b < 0x80 { + break + } + } + m.DataInRemote = bool(v != 0) default: iNdEx = preIndex skippy, err := protohelpers.Skip(dAtA[iNdEx:]) diff --git a/weed/pb/master.proto b/weed/pb/master.proto index 2dfd9caa0..eec385e4b 100644 --- a/weed/pb/master.proto +++ b/weed/pb/master.proto @@ -222,6 +222,7 @@ message VolumeLocation { uint32 grpc_port = 7; repeated uint32 new_ec_vids = 8; repeated uint32 deleted_ec_vids = 9; + repeated uint32 remote_vids = 10; } message ClusterNodeUpdate { @@ -267,6 +268,7 @@ message Location { string public_url = 2; uint32 grpc_port = 3; string data_center = 4; + bool data_in_remote = 5; } message AssignRequest { diff --git a/weed/pb/master_pb/master.pb.go b/weed/pb/master_pb/master.pb.go index c64f8931a..95e34060f 100644 --- a/weed/pb/master_pb/master.pb.go +++ b/weed/pb/master_pb/master.pb.go @@ -1061,6 +1061,7 @@ type VolumeLocation struct { GrpcPort uint32 `protobuf:"varint,7,opt,name=grpc_port,json=grpcPort,proto3" json:"grpc_port,omitempty"` NewEcVids []uint32 `protobuf:"varint,8,rep,packed,name=new_ec_vids,json=newEcVids,proto3" json:"new_ec_vids,omitempty"` DeletedEcVids []uint32 `protobuf:"varint,9,rep,packed,name=deleted_ec_vids,json=deletedEcVids,proto3" json:"deleted_ec_vids,omitempty"` + RemoteVids []uint32 `protobuf:"varint,10,rep,packed,name=remote_vids,json=remoteVids,proto3" json:"remote_vids,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -1158,6 +1159,13 @@ func (x *VolumeLocation) GetDeletedEcVids() []uint32 { return nil } +func (x *VolumeLocation) GetRemoteVids() []uint32 { + if x != nil { + return x.RemoteVids + } + return nil +} + type ClusterNodeUpdate struct { state protoimpl.MessageState `protogen:"open.v1"` NodeType string `protobuf:"bytes,1,opt,name=node_type,json=nodeType,proto3" json:"node_type,omitempty"` @@ -1460,6 +1468,7 @@ type Location struct { PublicUrl string `protobuf:"bytes,2,opt,name=public_url,json=publicUrl,proto3" json:"public_url,omitempty"` GrpcPort uint32 `protobuf:"varint,3,opt,name=grpc_port,json=grpcPort,proto3" json:"grpc_port,omitempty"` DataCenter string `protobuf:"bytes,4,opt,name=data_center,json=dataCenter,proto3" json:"data_center,omitempty"` + DataInRemote bool `protobuf:"varint,5,opt,name=data_in_remote,json=dataInRemote,proto3" json:"data_in_remote,omitempty"` unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -1522,6 +1531,13 @@ func (x *Location) GetDataCenter() string { return "" } +func (x *Location) GetDataInRemote() bool { + if x != nil { + return x.DataInRemote + } + return false +} + type AssignRequest struct { state protoimpl.MessageState `protogen:"open.v1"` Count uint64 `protobuf:"varint,1,opt,name=count,proto3" json:"count,omitempty"` @@ -5072,7 +5088,7 @@ const file_master_proto_rawDesc = "" + "filerGroup\x12\x1f\n" + "\vdata_center\x18\x06 \x01(\tR\n" + "dataCenter\x12\x12\n" + - "\x04rack\x18\a \x01(\tR\x04rack\"\x9d\x02\n" + + "\x04rack\x18\a \x01(\tR\x04rack\"\xbe\x02\n" + "\x0eVolumeLocation\x12\x10\n" + "\x03url\x18\x01 \x01(\tR\x03url\x12\x1d\n" + "\n" + @@ -5084,7 +5100,10 @@ const file_master_proto_rawDesc = "" + "dataCenter\x12\x1b\n" + "\tgrpc_port\x18\a \x01(\rR\bgrpcPort\x12\x1e\n" + "\vnew_ec_vids\x18\b \x03(\rR\tnewEcVids\x12&\n" + - "\x0fdeleted_ec_vids\x18\t \x03(\rR\rdeletedEcVids\"\xa6\x01\n" + + "\x0fdeleted_ec_vids\x18\t \x03(\rR\rdeletedEcVids\x12\x1f\n" + + "\vremote_vids\x18\n" + + " \x03(\rR\n" + + "remoteVids\"\xa6\x01\n" + "\x11ClusterNodeUpdate\x12\x1b\n" + "\tnode_type\x18\x01 \x01(\tR\bnodeType\x12\x18\n" + "\aaddress\x18\x02 \x01(\tR\aaddress\x12\x15\n" + @@ -5112,14 +5131,15 @@ const file_master_proto_rawDesc = "" + "\x11volume_or_file_id\x18\x01 \x01(\tR\x0evolumeOrFileId\x121\n" + "\tlocations\x18\x02 \x03(\v2\x13.master_pb.LocationR\tlocations\x12\x14\n" + "\x05error\x18\x03 \x01(\tR\x05error\x12\x12\n" + - "\x04auth\x18\x04 \x01(\tR\x04auth\"y\n" + + "\x04auth\x18\x04 \x01(\tR\x04auth\"\x9f\x01\n" + "\bLocation\x12\x10\n" + "\x03url\x18\x01 \x01(\tR\x03url\x12\x1d\n" + "\n" + "public_url\x18\x02 \x01(\tR\tpublicUrl\x12\x1b\n" + "\tgrpc_port\x18\x03 \x01(\rR\bgrpcPort\x12\x1f\n" + "\vdata_center\x18\x04 \x01(\tR\n" + - "dataCenter\"\xfe\x02\n" + + "dataCenter\x12$\n" + + "\x0edata_in_remote\x18\x05 \x01(\bR\fdataInRemote\"\xfe\x02\n" + "\rAssignRequest\x12\x14\n" + "\x05count\x18\x01 \x01(\x04R\x05count\x12 \n" + "\vreplication\x18\x02 \x01(\tR\vreplication\x12\x1e\n" + diff --git a/weed/server/filer_grpc_server.go b/weed/server/filer_grpc_server.go index 1679b1b4e..c30ed2c98 100644 --- a/weed/server/filer_grpc_server.go +++ b/weed/server/filer_grpc_server.go @@ -162,10 +162,11 @@ func wdclientLocationsToPb(locations []wdclient.Location) []*filer_pb.Location { locs := make([]*filer_pb.Location, 0, len(locations)) for _, loc := range locations { locs = append(locs, &filer_pb.Location{ - Url: loc.Url, - PublicUrl: loc.PublicUrl, - GrpcPort: uint32(loc.GrpcPort), - DataCenter: loc.DataCenter, + Url: loc.Url, + PublicUrl: loc.PublicUrl, + GrpcPort: uint32(loc.GrpcPort), + DataCenter: loc.DataCenter, + DataInRemote: loc.DataInRemote, }) } return locs diff --git a/weed/server/master_grpc_server.go b/weed/server/master_grpc_server.go index 15cac59f6..db6106025 100644 --- a/weed/server/master_grpc_server.go +++ b/weed/server/master_grpc_server.go @@ -12,6 +12,7 @@ import ( "github.com/seaweedfs/seaweedfs/weed/cluster" "github.com/seaweedfs/seaweedfs/weed/cluster/maintenance" + "github.com/seaweedfs/seaweedfs/weed/pb" "github.com/seaweedfs/seaweedfs/weed/stats" "github.com/seaweedfs/seaweedfs/weed/storage/backend" @@ -262,7 +263,11 @@ func (ms *MasterServer) SendHeartbeat(stream master_pb.Seaweed_SendHeartbeatServ for _, v := range newVolumes { glog.V(1).Infof("master see new volume %d from %s", uint32(v.Id), dn.Url()) - message.NewVids = append(message.NewVids, uint32(v.Id)) + if v.IsRemote() { + message.RemoteVids = append(message.RemoteVids, uint32(v.Id)) + } else { + message.NewVids = append(message.NewVids, uint32(v.Id)) + } } for _, v := range deletedVolumes { glog.V(1).Infof("master see deleted volume %d from %s", uint32(v.Id), dn.Url()) @@ -307,7 +312,7 @@ func (ms *MasterServer) SendHeartbeat(stream master_pb.Seaweed_SendHeartbeatServ } } - if len(message.NewVids) > 0 || len(message.DeletedVids) > 0 || len(message.NewEcVids) > 0 || len(message.DeletedEcVids) > 0 { + if len(message.NewVids) > 0 || len(message.DeletedVids) > 0 || len(message.NewEcVids) > 0 || len(message.DeletedEcVids) > 0 || len(message.RemoteVids) > 0 { ms.broadcastToClients(&master_pb.KeepConnectedResponse{VolumeLocation: message}) } diff --git a/weed/server/master_grpc_server_volume.go b/weed/server/master_grpc_server_volume.go index ef09c8f05..5febc4f7d 100644 --- a/weed/server/master_grpc_server_volume.go +++ b/weed/server/master_grpc_server_volume.go @@ -167,10 +167,11 @@ func (ms *MasterServer) LookupVolume(ctx context.Context, req *master_pb.LookupV var locations []*master_pb.Location for _, loc := range result.Locations { locations = append(locations, &master_pb.Location{ - Url: loc.Url, - PublicUrl: loc.PublicUrl, - DataCenter: loc.DataCenter, - GrpcPort: uint32(loc.GrpcPort), + Url: loc.Url, + PublicUrl: loc.PublicUrl, + DataCenter: loc.DataCenter, + GrpcPort: uint32(loc.GrpcPort), + DataInRemote: loc.DataInRemote, }) } var auth string diff --git a/weed/server/master_server_handlers.go b/weed/server/master_server_handlers.go index f028a5a04..10e84fe01 100644 --- a/weed/server/master_server_handlers.go +++ b/weed/server/master_server_handlers.go @@ -81,15 +81,21 @@ func (ms *MasterServer) findVolumeLocation(collection, vid string) operation.Loo if ms.Topo.IsLeader() { volumeId, newVolumeIdErr := needle.NewVolumeId(vid) if newVolumeIdErr != nil { - err = fmt.Errorf("Unknown volume id %s", vid) + err = fmt.Errorf("unknown volume id %s", vid) } else { machines := ms.Topo.Lookup(collection, volumeId) for _, loc := range machines { + volInfo, err := loc.GetVolumesById(volumeId) + if err != nil { + glog.V(0).Infof("failed to get volume info from %s: %v", loc.Url(), err) + continue + } locations = append(locations, operation.Location{ - Url: loc.Url(), - PublicUrl: loc.PublicUrl, - DataCenter: loc.GetDataCenterId(), - GrpcPort: loc.GrpcPort, + Url: loc.Url(), + PublicUrl: loc.PublicUrl, + DataCenter: loc.GetDataCenterId(), + GrpcPort: loc.GrpcPort, + DataInRemote: volInfo.IsRemote(), }) } } @@ -97,10 +103,11 @@ func (ms *MasterServer) findVolumeLocation(collection, vid string) operation.Loo machines, getVidLocationsErr := ms.MasterClient.GetVidLocations(vid) for _, loc := range machines { locations = append(locations, operation.Location{ - Url: loc.Url, - PublicUrl: loc.PublicUrl, - DataCenter: loc.DataCenter, - GrpcPort: loc.GrpcPort, + Url: loc.Url, + PublicUrl: loc.PublicUrl, + DataCenter: loc.DataCenter, + GrpcPort: loc.GrpcPort, + DataInRemote: loc.DataInRemote, }) } err = getVidLocationsErr diff --git a/weed/topology/topology_info.go b/weed/topology/topology_info.go index f8df207b0..b0b157a11 100644 --- a/weed/topology/topology_info.go +++ b/weed/topology/topology_info.go @@ -100,7 +100,15 @@ func (t *Topology) ToVolumeLocations() (volumeLocations []*master_pb.VolumeLocat DataCenter: dn.GetDataCenterId(), GrpcPort: uint32(dn.GrpcPort), } - volumeLocation.NewVids = dn.AppendVolumeIds(nil) + + for _, v := range dn.GetVolumes() { + if v.IsRemote() { + volumeLocation.RemoteVids = append(volumeLocation.RemoteVids, uint32(v.Id)) + } else { + volumeLocation.NewVids = append(volumeLocation.NewVids, uint32(v.Id)) + } + } + // A single EC volume's shards can live on multiple disks of // one DataNode, so GetEcShards returns per-(vid,disk) entries. // Dedupe so the snapshot carries each vid once. diff --git a/weed/util/slice.go b/weed/util/slice.go index 6b60beec5..b5d3cad81 100644 --- a/weed/util/slice.go +++ b/weed/util/slice.go @@ -14,3 +14,18 @@ func DrainChannel[T any](ch chan T, first T) []T { } } } + +func ReorderToFront[T comparable](frontMap map[T]bool, inputSlice []T) []T { + var prioritized []T + var remaining []T + + for _, item := range inputSlice { + if frontMap[item] { + prioritized = append(prioritized, item) + } else { + remaining = append(remaining, item) + } + } + + return append(prioritized, remaining...) +} diff --git a/weed/util/slice_test.go b/weed/util/slice_test.go new file mode 100644 index 000000000..52fd29ac5 --- /dev/null +++ b/weed/util/slice_test.go @@ -0,0 +1,40 @@ +package util + +import ( + "reflect" + "testing" +) + +func TestReorderToFront_StringSlice(t *testing.T) { + localUrls := map[string]bool{ + "http://local1": true, + "http://local2": true, + } + + sameDcTargetUrls := []string{ + "http://remote1", + "http://local2", + "http://remote2", + "http://local1", + } + + expected1 := []string{ + "http://local1", + "http://local2", + "http://remote1", + "http://remote2", + } + + expected2 := []string{ + "http://local2", + "http://local1", + "http://remote1", + "http://remote2", + } + + result := ReorderToFront(localUrls, sameDcTargetUrls) + + if !reflect.DeepEqual(result, expected1) && !reflect.DeepEqual(result, expected2) { + t.Errorf("ReorderToFront failed for strings. Got: %v, Expected1: %v, Expected2: %v", result, expected1, expected2) + } +} diff --git a/weed/wdclient/masterclient.go b/weed/wdclient/masterclient.go index 4fe8a1e2a..86bc801f2 100644 --- a/weed/wdclient/masterclient.go +++ b/weed/wdclient/masterclient.go @@ -109,6 +109,7 @@ func (p *masterVolumeProvider) LookupVolumeIds(ctx context.Context, volumeIds [] PublicUrl: masterLoc.PublicUrl, GrpcPort: int(masterLoc.GrpcPort), DataCenter: masterLoc.DataCenter, + DataInRemote: masterLoc.DataInRemote, } // Update cache with the location p.masterClient.addLocation(uint32(vid), loc) @@ -384,6 +385,12 @@ func (mc *MasterClient) updateVidMap(resp *master_pb.KeepConnectedResponse) { glog.V(2).Infof("%s.%s: %s masterClient adds volume %d", mc.FilerGroup, mc.clientType, loc.Url, newVid) mc.addLocation(newVid, loc) } + for _, remoteVid := range resp.VolumeLocation.RemoteVids { + remoteLoc := loc + remoteLoc.DataInRemote = true + glog.V(2).Infof("%s.%s: %s masterClient adds remote volume %d", mc.FilerGroup, mc.clientType, remoteLoc.Url, remoteVid) + mc.addLocation(remoteVid, remoteLoc) + } for _, deletedVid := range resp.VolumeLocation.DeletedVids { if _, moved := stillOnServer[deletedVid]; moved { continue @@ -403,10 +410,11 @@ func (mc *MasterClient) updateVidMap(resp *master_pb.KeepConnectedResponse) { glog.V(2).Infof("%s.%s: %s masterClient removes ec volume %d", mc.FilerGroup, mc.clientType, loc.Url, deletedEcVid) mc.deleteEcLocation(deletedEcVid, loc) } - glog.V(1).Infof("updateVidMap(%s) %s.%s: %s volume add: %d, del: %d, add ec: %d del ec: %d", + glog.V(1).Infof("updateVidMap(%s) %s.%s: %s volume add local: %d, remote: %d, del: %d, add ec: %d del ec: %d", resp.VolumeLocation.DataCenter, mc.FilerGroup, mc.clientType, loc.Url, - len(resp.VolumeLocation.NewVids), len(resp.VolumeLocation.DeletedVids), - len(resp.VolumeLocation.NewEcVids), len(resp.VolumeLocation.DeletedEcVids)) + len(resp.VolumeLocation.NewVids), len(resp.VolumeLocation.RemoteVids), + len(resp.VolumeLocation.DeletedVids), len(resp.VolumeLocation.NewEcVids), + len(resp.VolumeLocation.DeletedEcVids)) } func (mc *MasterClient) WithClient(ctx context.Context, streamingMode bool, fn func(client master_pb.SeaweedClient) error) error { diff --git a/weed/wdclient/vid_map.go b/weed/wdclient/vid_map.go index fdae7dfca..50ace9dc5 100644 --- a/weed/wdclient/vid_map.go +++ b/weed/wdclient/vid_map.go @@ -9,9 +9,9 @@ import ( "strings" "sync" - "github.com/seaweedfs/seaweedfs/weed/pb" - "github.com/seaweedfs/seaweedfs/weed/glog" + "github.com/seaweedfs/seaweedfs/weed/pb" + "github.com/seaweedfs/seaweedfs/weed/util" ) type HasLookupFileIdFunction interface { @@ -21,10 +21,11 @@ type HasLookupFileIdFunction interface { type LookupFileIdFunctionType func(ctx context.Context, fileId string) (targetUrls []string, err error) type Location struct { - Url string `json:"url,omitempty"` - PublicUrl string `json:"publicUrl,omitempty"` - DataCenter string `json:"dataCenter,omitempty"` - GrpcPort int `json:"grpcPort,omitempty"` + Url string `json:"url,omitempty"` + PublicUrl string `json:"publicUrl,omitempty"` + DataCenter string `json:"dataCenter,omitempty"` + GrpcPort int `json:"grpcPort,omitempty"` + DataInRemote bool `json:"dataInRemote,omitempty"` } func (l Location) ServerAddress() pb.ServerAddress { @@ -99,7 +100,15 @@ func (vc *vidMap) LookupVolumeServerUrl(vid string) (serverUrls []string, err er return nil, fmt.Errorf("volume %d not found", id) } var sameDcServers, otherDcServers []string + localUrls := make(map[string]bool) + for _, loc := range locations { + glog.V(4).Infof("lookup %s => %s, data in remote storage tier: %v", vid, loc.Url, loc.DataInRemote) + + if !loc.DataInRemote { + localUrls[loc.Url] = true + } + if vc.isSameDataCenter(&loc) { sameDcServers = append(sameDcServers, loc.Url) } else { @@ -112,6 +121,10 @@ func (vc *vidMap) LookupVolumeServerUrl(vid string) (serverUrls []string, err er rand.Shuffle(len(otherDcServers), func(i, j int) { otherDcServers[i], otherDcServers[j] = otherDcServers[j], otherDcServers[i] }) + if len(localUrls) > 0 { + sameDcServers = util.ReorderToFront(localUrls, sameDcServers) + otherDcServers = util.ReorderToFront(localUrls, otherDcServers) + } // Prefer same data center serverUrls = append(sameDcServers, otherDcServers...) return