seaweedfs/weed/command/filer_replication.go

82 lines
2.5 KiB
Go
Raw Normal View History

2018-09-17 07:27:56 +00:00
package command
import (
"github.com/chrislusf/seaweedfs/weed/glog"
2018-09-21 08:56:43 +00:00
"github.com/chrislusf/seaweedfs/weed/replication"
2018-09-17 07:27:56 +00:00
"github.com/chrislusf/seaweedfs/weed/server"
"github.com/spf13/viper"
2018-09-23 07:40:36 +00:00
"strings"
2018-09-17 07:27:56 +00:00
)
func init() {
cmdFilerReplicate.Run = runFilerReplicate // break init cycle
}
var cmdFilerReplicate = &Command{
UsageLine: "filer.replicate",
Short: "replicate file changes to another destination",
Long: `replicate file changes to another destination
filer.replicate listens on filer notifications. If any file is updated, it will fetch the updated content,
and write to the other destination.
Run "weed scaffold -config replication" to generate a replication.toml file and customize the parameters.
`,
}
func runFilerReplicate(cmd *Command, args []string) bool {
weed_server.LoadConfiguration("replication", true)
config := viper.GetViper()
var notificationInput replication.NotificationInput
for _, input := range replication.NotificationInputs {
if config.GetBool("notification." + input.GetName() + ".enabled") {
viperSub := config.Sub("notification." + input.GetName())
if err := input.Initialize(viperSub); err != nil {
glog.Fatalf("Failed to initialize notification input for %s: %+v",
input.GetName(), err)
}
glog.V(0).Infof("Configure notification input to %s", input.GetName())
notificationInput = input
break
}
}
2018-09-23 07:40:36 +00:00
// avoid recursive replication
if config.GetBool("notification.source.filer.enabled") && config.GetBool("notification.sink.filer.enabled") {
sourceConfig, sinkConfig := config.Sub("source.filer"), config.Sub("sink.filer")
if sourceConfig.GetString("grpcAddress") == sinkConfig.GetString("grpcAddress") {
fromDir := sourceConfig.GetString("directory")
toDir := sinkConfig.GetString("directory")
if strings.HasPrefix(toDir, fromDir) {
glog.Fatalf("recursive replication! source directory %s includes the sink directory %s", fromDir, toDir)
}
}
}
2018-09-17 08:37:24 +00:00
replicator := replication.NewReplicator(config.Sub("source.filer"), config.Sub("sink.filer"))
2018-09-17 07:27:56 +00:00
for {
key, m, err := notificationInput.ReceiveMessage()
if err != nil {
2018-09-17 09:23:21 +00:00
glog.Errorf("receive %s: %+v", key, err)
2018-09-17 07:27:56 +00:00
continue
}
2018-09-22 07:12:10 +00:00
if m.OldEntry != nil && m.NewEntry == nil {
glog.V(1).Infof("delete: %s", key)
2018-09-22 07:12:10 +00:00
} else if m.OldEntry == nil && m.NewEntry != nil {
glog.V(1).Infof(" add: %s", key)
2018-09-22 07:12:10 +00:00
} else {
glog.V(1).Infof("modify: %s", key)
}
2018-09-17 07:27:56 +00:00
if err = replicator.Replicate(key, m); err != nil {
2018-09-17 09:23:21 +00:00
glog.Errorf("replicate %s: %+v", key, err)
2018-09-17 07:27:56 +00:00
}
}
return true
}