You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

121 lines
4.2 KiB

5 years ago
5 years ago
3 years ago
5 years ago
3 years ago
5 years ago
5 years ago
5 years ago
5 years ago
5 years ago
5 years ago
5 years ago
5 years ago
5 years ago
5 years ago
5 years ago
5 years ago
5 years ago
5 years ago
5 years ago
3 years ago
5 years ago
5 years ago
5 years ago
5 years ago
3 years ago
5 years ago
  1. package command
  2. import (
  3. "context"
  4. "fmt"
  5. "time"
  6. "google.golang.org/grpc/reflection"
  7. "github.com/chrislusf/seaweedfs/weed/util/grace"
  8. "github.com/chrislusf/seaweedfs/weed/glog"
  9. "github.com/chrislusf/seaweedfs/weed/mq/broker"
  10. "github.com/chrislusf/seaweedfs/weed/pb"
  11. "github.com/chrislusf/seaweedfs/weed/pb/filer_pb"
  12. "github.com/chrislusf/seaweedfs/weed/pb/mq_pb"
  13. "github.com/chrislusf/seaweedfs/weed/security"
  14. "github.com/chrislusf/seaweedfs/weed/util"
  15. )
  16. var (
  17. mqBrokerStandaloneOptions MessageQueueBrokerOptions
  18. )
  19. type MessageQueueBrokerOptions struct {
  20. masters *string
  21. filerGroup *string
  22. filer *string
  23. ip *string
  24. port *int
  25. dataCenter *string
  26. rack *string
  27. cpuprofile *string
  28. memprofile *string
  29. }
  30. func init() {
  31. cmdMqBroker.Run = runMqBroker // break init cycle
  32. mqBrokerStandaloneOptions.masters = cmdMqBroker.Flag.String("master", "localhost:9333", "comma-separated master servers")
  33. mqBrokerStandaloneOptions.filer = cmdMqBroker.Flag.String("filer", "localhost:8888", "filer server address")
  34. mqBrokerStandaloneOptions.filerGroup = cmdMqBroker.Flag.String("filerGroup", "", "share metadata with other filers in the same filerGroup")
  35. mqBrokerStandaloneOptions.ip = cmdMqBroker.Flag.String("ip", util.DetectedHostAddress(), "broker host address")
  36. mqBrokerStandaloneOptions.port = cmdMqBroker.Flag.Int("port", 17777, "broker gRPC listen port")
  37. mqBrokerStandaloneOptions.dataCenter = cmdMqBroker.Flag.String("dataCenter", "", "prefer to read and write to volumes in this data center")
  38. mqBrokerStandaloneOptions.rack = cmdMqBroker.Flag.String("rack", "", "prefer to write to volumes in this rack")
  39. mqBrokerStandaloneOptions.cpuprofile = cmdMqBroker.Flag.String("cpuprofile", "", "cpu profile output file")
  40. mqBrokerStandaloneOptions.memprofile = cmdMqBroker.Flag.String("memprofile", "", "memory profile output file")
  41. }
  42. var cmdMqBroker = &Command{
  43. UsageLine: "mq.broker [-port=17777] [-filer=<ip:port>]",
  44. Short: "start a message queue broker",
  45. Long: `start a message queue broker
  46. The broker can accept gRPC calls to write or read messages. The messages are stored via filer.
  47. The brokers are stateless. To scale up, just add more brokers.
  48. `,
  49. }
  50. func runMqBroker(cmd *Command, args []string) bool {
  51. util.LoadConfiguration("security", false)
  52. return mqBrokerStandaloneOptions.startQueueServer()
  53. }
  54. func (mqBrokerOpt *MessageQueueBrokerOptions) startQueueServer() bool {
  55. grace.SetupProfiling(*mqBrokerStandaloneOptions.cpuprofile, *mqBrokerStandaloneOptions.memprofile)
  56. filerAddress := pb.ServerAddress(*mqBrokerOpt.filer)
  57. grpcDialOption := security.LoadClientTLS(util.GetViper(), "grpc.msg_broker")
  58. cipher := false
  59. for {
  60. err := pb.WithGrpcFilerClient(false, filerAddress, grpcDialOption, func(client filer_pb.SeaweedFilerClient) error {
  61. resp, err := client.GetFilerConfiguration(context.Background(), &filer_pb.GetFilerConfigurationRequest{})
  62. if err != nil {
  63. return fmt.Errorf("get filer %s configuration: %v", filerAddress, err)
  64. }
  65. cipher = resp.Cipher
  66. return nil
  67. })
  68. if err != nil {
  69. glog.V(0).Infof("wait to connect to filer %s grpc address %s", *mqBrokerOpt.filer, filerAddress.ToGrpcAddress())
  70. time.Sleep(time.Second)
  71. } else {
  72. glog.V(0).Infof("connected to filer %s grpc address %s", *mqBrokerOpt.filer, filerAddress.ToGrpcAddress())
  73. break
  74. }
  75. }
  76. qs, err := broker.NewMessageBroker(&broker.MessageQueueBrokerOption{
  77. Masters: pb.ServerAddresses(*mqBrokerOpt.masters).ToAddressMap(),
  78. FilerGroup: *mqBrokerOpt.filerGroup,
  79. DataCenter: *mqBrokerOpt.dataCenter,
  80. Rack: *mqBrokerOpt.rack,
  81. Filers: []pb.ServerAddress{filerAddress},
  82. DefaultReplication: "",
  83. MaxMB: 0,
  84. Ip: *mqBrokerOpt.ip,
  85. Port: *mqBrokerOpt.port,
  86. Cipher: cipher,
  87. }, grpcDialOption)
  88. // start grpc listener
  89. grpcL, _, err := util.NewIpAndLocalListeners("", *mqBrokerOpt.port, 0)
  90. if err != nil {
  91. glog.Fatalf("failed to listen on grpc port %d: %v", *mqBrokerOpt.port, err)
  92. }
  93. grpcS := pb.NewGrpcServer(security.LoadServerTLS(util.GetViper(), "grpc.msg_broker"))
  94. mq_pb.RegisterSeaweedMessagingServer(grpcS, qs)
  95. reflection.Register(grpcS)
  96. grpcS.Serve(grpcL)
  97. return true
  98. }