command_ec_encode.go 9.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297
  1. package shell
  2. import (
  3. "context"
  4. "flag"
  5. "fmt"
  6. "io"
  7. "sync"
  8. "time"
  9. "google.golang.org/grpc"
  10. "github.com/chrislusf/seaweedfs/weed/operation"
  11. "github.com/chrislusf/seaweedfs/weed/pb/master_pb"
  12. "github.com/chrislusf/seaweedfs/weed/pb/volume_server_pb"
  13. "github.com/chrislusf/seaweedfs/weed/storage/erasure_coding"
  14. "github.com/chrislusf/seaweedfs/weed/storage/needle"
  15. "github.com/chrislusf/seaweedfs/weed/wdclient"
  16. )
  17. func init() {
  18. Commands = append(Commands, &commandEcEncode{})
  19. }
  20. type commandEcEncode struct {
  21. }
  22. func (c *commandEcEncode) Name() string {
  23. return "ec.encode"
  24. }
  25. func (c *commandEcEncode) Help() string {
  26. return `apply erasure coding to a volume
  27. ec.encode [-collection=""] [-fullPercent=95] [-quietFor=1h]
  28. ec.encode [-collection=""] [-volumeId=<volume_id>]
  29. This command will:
  30. 1. freeze one volume
  31. 2. apply erasure coding to the volume
  32. 3. move the encoded shards to multiple volume servers
  33. The erasure coding is 10.4. So ideally you have more than 14 volume servers, and you can afford
  34. to lose 4 volume servers.
  35. If the number of volumes are not high, the worst case is that you only have 4 volume servers,
  36. and the shards are spread as 4,4,3,3, respectively. You can afford to lose one volume server.
  37. If you only have less than 4 volume servers, with erasure coding, at least you can afford to
  38. have 4 corrupted shard files.
  39. `
  40. }
  41. func (c *commandEcEncode) Do(args []string, commandEnv *CommandEnv, writer io.Writer) (err error) {
  42. if err = commandEnv.confirmIsLocked(); err != nil {
  43. return
  44. }
  45. encodeCommand := flag.NewFlagSet(c.Name(), flag.ContinueOnError)
  46. volumeId := encodeCommand.Int("volumeId", 0, "the volume id")
  47. collection := encodeCommand.String("collection", "", "the collection name")
  48. fullPercentage := encodeCommand.Float64("fullPercent", 95, "the volume reaches the percentage of max volume size")
  49. quietPeriod := encodeCommand.Duration("quietFor", time.Hour, "select volumes without no writes for this period")
  50. if err = encodeCommand.Parse(args); err != nil {
  51. return nil
  52. }
  53. vid := needle.VolumeId(*volumeId)
  54. // volumeId is provided
  55. if vid != 0 {
  56. return doEcEncode(commandEnv, *collection, vid)
  57. }
  58. // apply to all volumes in the collection
  59. volumeIds, err := collectVolumeIdsForEcEncode(commandEnv, *collection, *fullPercentage, *quietPeriod)
  60. if err != nil {
  61. return err
  62. }
  63. fmt.Printf("ec encode volumes: %v\n", volumeIds)
  64. for _, vid := range volumeIds {
  65. if err = doEcEncode(commandEnv, *collection, vid); err != nil {
  66. return err
  67. }
  68. }
  69. return nil
  70. }
  71. func doEcEncode(commandEnv *CommandEnv, collection string, vid needle.VolumeId) (err error) {
  72. // find volume location
  73. locations, found := commandEnv.MasterClient.GetLocations(uint32(vid))
  74. if !found {
  75. return fmt.Errorf("volume %d not found", vid)
  76. }
  77. // fmt.Printf("found ec %d shards on %v\n", vid, locations)
  78. // mark the volume as readonly
  79. err = markVolumeReadonly(commandEnv.option.GrpcDialOption, vid, locations)
  80. if err != nil {
  81. return fmt.Errorf("mark volume %d as readonly on %s: %v", vid, locations[0].Url, err)
  82. }
  83. // generate ec shards
  84. err = generateEcShards(commandEnv.option.GrpcDialOption, vid, collection, locations[0].Url)
  85. if err != nil {
  86. return fmt.Errorf("generate ec shards for volume %d on %s: %v", vid, locations[0].Url, err)
  87. }
  88. // balance the ec shards to current cluster
  89. err = spreadEcShards(commandEnv, vid, collection, locations)
  90. if err != nil {
  91. return fmt.Errorf("spread ec shards for volume %d from %s: %v", vid, locations[0].Url, err)
  92. }
  93. return nil
  94. }
  95. func markVolumeReadonly(grpcDialOption grpc.DialOption, volumeId needle.VolumeId, locations []wdclient.Location) error {
  96. for _, location := range locations {
  97. fmt.Printf("markVolumeReadonly %d on %s ...\n", volumeId, location.Url)
  98. err := operation.WithVolumeServerClient(location.Url, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
  99. _, markErr := volumeServerClient.VolumeMarkReadonly(context.Background(), &volume_server_pb.VolumeMarkReadonlyRequest{
  100. VolumeId: uint32(volumeId),
  101. })
  102. return markErr
  103. })
  104. if err != nil {
  105. return err
  106. }
  107. }
  108. return nil
  109. }
  110. func generateEcShards(grpcDialOption grpc.DialOption, volumeId needle.VolumeId, collection string, sourceVolumeServer string) error {
  111. fmt.Printf("generateEcShards %s %d on %s ...\n", collection, volumeId, sourceVolumeServer)
  112. err := operation.WithVolumeServerClient(sourceVolumeServer, grpcDialOption, func(volumeServerClient volume_server_pb.VolumeServerClient) error {
  113. _, genErr := volumeServerClient.VolumeEcShardsGenerate(context.Background(), &volume_server_pb.VolumeEcShardsGenerateRequest{
  114. VolumeId: uint32(volumeId),
  115. Collection: collection,
  116. })
  117. return genErr
  118. })
  119. return err
  120. }
  121. func spreadEcShards(commandEnv *CommandEnv, volumeId needle.VolumeId, collection string, existingLocations []wdclient.Location) (err error) {
  122. allEcNodes, totalFreeEcSlots, err := collectEcNodes(commandEnv, "")
  123. if err != nil {
  124. return err
  125. }
  126. if totalFreeEcSlots < erasure_coding.TotalShardsCount {
  127. return fmt.Errorf("not enough free ec shard slots. only %d left", totalFreeEcSlots)
  128. }
  129. allocatedDataNodes := allEcNodes
  130. if len(allocatedDataNodes) > erasure_coding.TotalShardsCount {
  131. allocatedDataNodes = allocatedDataNodes[:erasure_coding.TotalShardsCount]
  132. }
  133. // calculate how many shards to allocate for these servers
  134. allocatedEcIds := balancedEcDistribution(allocatedDataNodes)
  135. // ask the data nodes to copy from the source volume server
  136. copiedShardIds, err := parallelCopyEcShardsFromSource(commandEnv.option.GrpcDialOption, allocatedDataNodes, allocatedEcIds, volumeId, collection, existingLocations[0])
  137. if err != nil {
  138. return err
  139. }
  140. // unmount the to be deleted shards
  141. err = unmountEcShards(commandEnv.option.GrpcDialOption, volumeId, existingLocations[0].Url, copiedShardIds)
  142. if err != nil {
  143. return err
  144. }
  145. // ask the source volume server to clean up copied ec shards
  146. err = sourceServerDeleteEcShards(commandEnv.option.GrpcDialOption, collection, volumeId, existingLocations[0].Url, copiedShardIds)
  147. if err != nil {
  148. return fmt.Errorf("source delete copied ecShards %s %d.%v: %v", existingLocations[0].Url, volumeId, copiedShardIds, err)
  149. }
  150. // ask the source volume server to delete the original volume
  151. for _, location := range existingLocations {
  152. fmt.Printf("delete volume %d from %s\n", volumeId, location.Url)
  153. err = deleteVolume(commandEnv.option.GrpcDialOption, volumeId, location.Url)
  154. if err != nil {
  155. return fmt.Errorf("deleteVolume %s volume %d: %v", location.Url, volumeId, err)
  156. }
  157. }
  158. return err
  159. }
  160. func parallelCopyEcShardsFromSource(grpcDialOption grpc.DialOption, targetServers []*EcNode, allocatedEcIds [][]uint32, volumeId needle.VolumeId, collection string, existingLocation wdclient.Location) (actuallyCopied []uint32, err error) {
  161. fmt.Printf("parallelCopyEcShardsFromSource %d %s\n", volumeId, existingLocation.Url)
  162. // parallelize
  163. shardIdChan := make(chan []uint32, len(targetServers))
  164. var wg sync.WaitGroup
  165. for i, server := range targetServers {
  166. if len(allocatedEcIds[i]) <= 0 {
  167. continue
  168. }
  169. wg.Add(1)
  170. go func(server *EcNode, allocatedEcShardIds []uint32) {
  171. defer wg.Done()
  172. copiedShardIds, copyErr := oneServerCopyAndMountEcShardsFromSource(grpcDialOption, server,
  173. allocatedEcShardIds, volumeId, collection, existingLocation.Url)
  174. if copyErr != nil {
  175. err = copyErr
  176. } else {
  177. shardIdChan <- copiedShardIds
  178. server.addEcVolumeShards(volumeId, collection, copiedShardIds)
  179. }
  180. }(server, allocatedEcIds[i])
  181. }
  182. wg.Wait()
  183. close(shardIdChan)
  184. if err != nil {
  185. return nil, err
  186. }
  187. for shardIds := range shardIdChan {
  188. actuallyCopied = append(actuallyCopied, shardIds...)
  189. }
  190. return
  191. }
  192. func balancedEcDistribution(servers []*EcNode) (allocated [][]uint32) {
  193. allocated = make([][]uint32, len(servers))
  194. allocatedShardIdIndex := uint32(0)
  195. serverIndex := 0
  196. for allocatedShardIdIndex < erasure_coding.TotalShardsCount {
  197. if servers[serverIndex].freeEcSlot > 0 {
  198. allocated[serverIndex] = append(allocated[serverIndex], allocatedShardIdIndex)
  199. allocatedShardIdIndex++
  200. }
  201. serverIndex++
  202. if serverIndex >= len(servers) {
  203. serverIndex = 0
  204. }
  205. }
  206. return allocated
  207. }
  208. func collectVolumeIdsForEcEncode(commandEnv *CommandEnv, selectedCollection string, fullPercentage float64, quietPeriod time.Duration) (vids []needle.VolumeId, err error) {
  209. // collect topology information
  210. topologyInfo, volumeSizeLimitMb, err := collectTopologyInfo(commandEnv)
  211. if err != nil {
  212. return
  213. }
  214. quietSeconds := int64(quietPeriod / time.Second)
  215. nowUnixSeconds := time.Now().Unix()
  216. fmt.Printf("collect volumes quiet for: %d seconds\n", quietSeconds)
  217. vidMap := make(map[uint32]bool)
  218. eachDataNode(topologyInfo, func(dc string, rack RackId, dn *master_pb.DataNodeInfo) {
  219. for _, diskInfo := range dn.DiskInfos {
  220. for _, v := range diskInfo.VolumeInfos {
  221. if v.Collection == selectedCollection && v.ModifiedAtSecond+quietSeconds < nowUnixSeconds {
  222. if float64(v.Size) > fullPercentage/100*float64(volumeSizeLimitMb)*1024*1024 {
  223. vidMap[v.Id] = true
  224. }
  225. }
  226. }
  227. }
  228. })
  229. for vid := range vidMap {
  230. vids = append(vids, needle.VolumeId(vid))
  231. }
  232. return
  233. }