253 lines
8.2 KiB

3 years ago
3 years ago
3 years ago
3 years ago
3 years ago
2 years ago
  1. package command
  2. import (
  3. "context"
  4. "fmt"
  5. "github.com/seaweedfs/seaweedfs/weed/filer"
  6. "github.com/seaweedfs/seaweedfs/weed/glog"
  7. "github.com/spf13/viper"
  8. "google.golang.org/grpc"
  9. "reflect"
  10. "time"
  11. "github.com/seaweedfs/seaweedfs/weed/pb"
  12. "github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
  13. "github.com/seaweedfs/seaweedfs/weed/security"
  14. "github.com/seaweedfs/seaweedfs/weed/util"
  15. )
  16. var (
  17. metaBackup FilerMetaBackupOptions
  18. )
  19. type FilerMetaBackupOptions struct {
  20. grpcDialOption grpc.DialOption
  21. filerAddress *string
  22. filerDirectory *string
  23. restart *bool
  24. backupFilerConfig *string
  25. store filer.FilerStore
  26. clientId int32
  27. clientEpoch int32
  28. }
  29. func init() {
  30. cmdFilerMetaBackup.Run = runFilerMetaBackup // break init cycle
  31. metaBackup.filerAddress = cmdFilerMetaBackup.Flag.String("filer", "localhost:8888", "filer hostname:port")
  32. metaBackup.filerDirectory = cmdFilerMetaBackup.Flag.String("filerDir", "/", "a folder on the filer")
  33. metaBackup.restart = cmdFilerMetaBackup.Flag.Bool("restart", false, "copy the full metadata before async incremental backup")
  34. metaBackup.backupFilerConfig = cmdFilerMetaBackup.Flag.String("config", "", "path to filer.toml specifying backup filer store")
  35. metaBackup.clientId = util.RandomInt32()
  36. }
  37. var cmdFilerMetaBackup = &Command{
  38. UsageLine: "filer.meta.backup [-filer=localhost:8888] [-filerDir=/] [-restart] -config=/path/to/backup_filer.toml",
  39. Short: "continuously backup filer meta data changes to anther filer store specified in a backup_filer.toml",
  40. Long: `continuously backup filer meta data changes.
  41. The backup writes to another filer store specified in a backup_filer.toml.
  42. weed filer.meta.backup -config=/path/to/backup_filer.toml -filer="localhost:8888"
  43. weed filer.meta.backup -config=/path/to/backup_filer.toml -filer="localhost:8888" -restart
  44. `,
  45. }
  46. func runFilerMetaBackup(cmd *Command, args []string) bool {
  47. util.LoadConfiguration("security", false)
  48. metaBackup.grpcDialOption = security.LoadClientTLS(util.GetViper(), "grpc.client")
  49. // load backup_filer.toml
  50. v := viper.New()
  51. v.SetConfigFile(*metaBackup.backupFilerConfig)
  52. if err := v.ReadInConfig(); err != nil { // Handle errors reading the config file
  53. glog.Fatalf("Failed to load %s file.\nPlease use this command to generate the a %s.toml file\n"+
  54. " weed scaffold -config=%s -output=.\n\n\n",
  55. *metaBackup.backupFilerConfig, "backup_filer", "filer")
  56. }
  57. if err := metaBackup.initStore(v); err != nil {
  58. glog.V(0).Infof("init backup filer store: %v", err)
  59. return true
  60. }
  61. missingPreviousBackup := false
  62. _, err := metaBackup.getOffset()
  63. if err != nil {
  64. missingPreviousBackup = true
  65. }
  66. if *metaBackup.restart || missingPreviousBackup {
  67. glog.V(0).Infof("traversing metadata tree...")
  68. startTime := time.Now()
  69. if err := metaBackup.traverseMetadata(); err != nil {
  70. glog.Errorf("traverse meta data: %v", err)
  71. return true
  72. }
  73. glog.V(0).Infof("metadata copied up to %v", startTime)
  74. if err := metaBackup.setOffset(startTime); err != nil {
  75. startTime = time.Now()
  76. }
  77. }
  78. for {
  79. err := metaBackup.streamMetadataBackup()
  80. if err != nil {
  81. glog.Errorf("filer meta backup from %s: %v", *metaBackup.filerAddress, err)
  82. time.Sleep(1747 * time.Millisecond)
  83. }
  84. }
  85. return true
  86. }
  87. func (metaBackup *FilerMetaBackupOptions) initStore(v *viper.Viper) error {
  88. // load configuration for default filer store
  89. hasDefaultStoreConfigured := false
  90. for _, store := range filer.Stores {
  91. if v.GetBool(store.GetName() + ".enabled") {
  92. store = reflect.New(reflect.ValueOf(store).Elem().Type()).Interface().(filer.FilerStore)
  93. if err := store.Initialize(v, store.GetName()+"."); err != nil {
  94. glog.Fatalf("failed to initialize store for %s: %+v", store.GetName(), err)
  95. }
  96. glog.V(0).Infof("configured filer store to %s", store.GetName())
  97. hasDefaultStoreConfigured = true
  98. metaBackup.store = filer.NewFilerStoreWrapper(store)
  99. break
  100. }
  101. }
  102. if !hasDefaultStoreConfigured {
  103. return fmt.Errorf("no filer store enabled in %s", v.ConfigFileUsed())
  104. }
  105. return nil
  106. }
  107. func (metaBackup *FilerMetaBackupOptions) traverseMetadata() (err error) {
  108. var saveErr error
  109. traverseErr := filer_pb.TraverseBfs(metaBackup, util.FullPath(*metaBackup.filerDirectory), func(parentPath util.FullPath, entry *filer_pb.Entry) {
  110. println("+", parentPath.Child(entry.Name))
  111. if err := metaBackup.store.InsertEntry(context.Background(), filer.FromPbEntry(string(parentPath), entry)); err != nil {
  112. saveErr = fmt.Errorf("insert entry error: %v\n", err)
  113. return
  114. }
  115. })
  116. if traverseErr != nil {
  117. return fmt.Errorf("traverse: %v", traverseErr)
  118. }
  119. return saveErr
  120. }
  121. var (
  122. MetaBackupKey = []byte("metaBackup")
  123. )
  124. func (metaBackup *FilerMetaBackupOptions) streamMetadataBackup() error {
  125. startTime, err := metaBackup.getOffset()
  126. if err != nil {
  127. startTime = time.Now()
  128. }
  129. glog.V(0).Infof("streaming from %v", startTime)
  130. store := metaBackup.store
  131. eachEntryFunc := func(resp *filer_pb.SubscribeMetadataResponse) error {
  132. ctx := context.Background()
  133. message := resp.EventNotification
  134. if filer_pb.IsEmpty(resp) {
  135. return nil
  136. } else if filer_pb.IsCreate(resp) {
  137. println("+", util.FullPath(message.NewParentPath).Child(message.NewEntry.Name))
  138. entry := filer.FromPbEntry(message.NewParentPath, message.NewEntry)
  139. return store.InsertEntry(ctx, entry)
  140. } else if filer_pb.IsDelete(resp) {
  141. println("-", util.FullPath(resp.Directory).Child(message.OldEntry.Name))
  142. return store.DeleteEntry(ctx, util.FullPath(resp.Directory).Child(message.OldEntry.Name))
  143. } else if filer_pb.IsUpdate(resp) {
  144. println("~", util.FullPath(message.NewParentPath).Child(message.NewEntry.Name))
  145. entry := filer.FromPbEntry(message.NewParentPath, message.NewEntry)
  146. return store.UpdateEntry(ctx, entry)
  147. } else {
  148. // renaming
  149. println("-", util.FullPath(resp.Directory).Child(message.OldEntry.Name))
  150. if err := store.DeleteEntry(ctx, util.FullPath(resp.Directory).Child(message.OldEntry.Name)); err != nil {
  151. return err
  152. }
  153. println("+", util.FullPath(message.NewParentPath).Child(message.NewEntry.Name))
  154. return store.InsertEntry(ctx, filer.FromPbEntry(message.NewParentPath, message.NewEntry))
  155. }
  156. return nil
  157. }
  158. processEventFnWithOffset := pb.AddOffsetFunc(eachEntryFunc, 3*time.Second, func(counter int64, lastTsNs int64) error {
  159. lastTime := time.Unix(0, lastTsNs)
  160. glog.V(0).Infof("meta backup %s progressed to %v %0.2f/sec", *metaBackup.filerAddress, lastTime, float64(counter)/float64(3))
  161. return metaBackup.setOffset(lastTime)
  162. })
  163. metaBackup.clientEpoch++
  164. metadataFollowOption := &pb.MetadataFollowOption{
  165. ClientName: "meta_backup",
  166. ClientId: metaBackup.clientId,
  167. ClientEpoch: metaBackup.clientEpoch,
  168. SelfSignature: 0,
  169. PathPrefix: *metaBackup.filerDirectory,
  170. AdditionalPathPrefixes: nil,
  171. DirectoriesToWatch: nil,
  172. StartTsNs: startTime.UnixNano(),
  173. StopTsNs: 0,
  174. EventErrorType: pb.TrivialOnError,
  175. }
  176. return pb.FollowMetadata(pb.ServerAddress(*metaBackup.filerAddress), metaBackup.grpcDialOption, metadataFollowOption, processEventFnWithOffset)
  177. }
  178. func (metaBackup *FilerMetaBackupOptions) getOffset() (lastWriteTime time.Time, err error) {
  179. value, err := metaBackup.store.KvGet(context.Background(), MetaBackupKey)
  180. if err != nil {
  181. return
  182. }
  183. tsNs := util.BytesToUint64(value)
  184. return time.Unix(0, int64(tsNs)), nil
  185. }
  186. func (metaBackup *FilerMetaBackupOptions) setOffset(lastWriteTime time.Time) error {
  187. valueBuf := make([]byte, 8)
  188. util.Uint64toBytes(valueBuf, uint64(lastWriteTime.UnixNano()))
  189. if err := metaBackup.store.KvPut(context.Background(), MetaBackupKey, valueBuf); err != nil {
  190. return err
  191. }
  192. return nil
  193. }
  194. var _ = filer_pb.FilerClient(&FilerMetaBackupOptions{})
  195. func (metaBackup *FilerMetaBackupOptions) WithFilerClient(streamingMode bool, fn func(filer_pb.SeaweedFilerClient) error) error {
  196. return pb.WithFilerClient(streamingMode, metaBackup.clientId, pb.ServerAddress(*metaBackup.filerAddress), metaBackup.grpcDialOption, func(client filer_pb.SeaweedFilerClient) error {
  197. return fn(client)
  198. })
  199. }
  200. func (metaBackup *FilerMetaBackupOptions) AdjustedUrl(location *filer_pb.Location) string {
  201. return location.Url
  202. }
  203. func (metaBackup *FilerMetaBackupOptions) GetDataCenter() string {
  204. return ""
  205. }