mq_broker.go 3.4 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697
  1. package command
  2. import (
  3. "google.golang.org/grpc/reflection"
  4. "github.com/seaweedfs/seaweedfs/weed/util/grace"
  5. "github.com/seaweedfs/seaweedfs/weed/glog"
  6. "github.com/seaweedfs/seaweedfs/weed/mq/broker"
  7. "github.com/seaweedfs/seaweedfs/weed/pb"
  8. "github.com/seaweedfs/seaweedfs/weed/pb/mq_pb"
  9. "github.com/seaweedfs/seaweedfs/weed/security"
  10. "github.com/seaweedfs/seaweedfs/weed/util"
  11. )
  12. var (
  13. mqBrokerStandaloneOptions MessageQueueBrokerOptions
  14. )
  15. type MessageQueueBrokerOptions struct {
  16. masters map[string]pb.ServerAddress
  17. mastersString *string
  18. filerGroup *string
  19. ip *string
  20. port *int
  21. dataCenter *string
  22. rack *string
  23. cpuprofile *string
  24. memprofile *string
  25. }
  26. func init() {
  27. cmdMqBroker.Run = runMqBroker // break init cycle
  28. mqBrokerStandaloneOptions.mastersString = cmdMqBroker.Flag.String("master", "localhost:9333", "comma-separated master servers")
  29. mqBrokerStandaloneOptions.filerGroup = cmdMqBroker.Flag.String("filerGroup", "", "share metadata with other filers in the same filerGroup")
  30. mqBrokerStandaloneOptions.ip = cmdMqBroker.Flag.String("ip", util.DetectedHostAddress(), "broker host address")
  31. mqBrokerStandaloneOptions.port = cmdMqBroker.Flag.Int("port", 17777, "broker gRPC listen port")
  32. mqBrokerStandaloneOptions.dataCenter = cmdMqBroker.Flag.String("dataCenter", "", "prefer to read and write to volumes in this data center")
  33. mqBrokerStandaloneOptions.rack = cmdMqBroker.Flag.String("rack", "", "prefer to write to volumes in this rack")
  34. mqBrokerStandaloneOptions.cpuprofile = cmdMqBroker.Flag.String("cpuprofile", "", "cpu profile output file")
  35. mqBrokerStandaloneOptions.memprofile = cmdMqBroker.Flag.String("memprofile", "", "memory profile output file")
  36. }
  37. var cmdMqBroker = &Command{
  38. UsageLine: "mq.broker [-port=17777] [-master=<ip:port>]",
  39. Short: "<WIP> start a message queue broker",
  40. Long: `start a message queue broker
  41. The broker can accept gRPC calls to write or read messages. The messages are stored via filer.
  42. The brokers are stateless. To scale up, just add more brokers.
  43. `,
  44. }
  45. func runMqBroker(cmd *Command, args []string) bool {
  46. util.LoadSecurityConfiguration()
  47. mqBrokerStandaloneOptions.masters = pb.ServerAddresses(*mqBrokerStandaloneOptions.mastersString).ToAddressMap()
  48. return mqBrokerStandaloneOptions.startQueueServer()
  49. }
  50. func (mqBrokerOpt *MessageQueueBrokerOptions) startQueueServer() bool {
  51. grace.SetupProfiling(*mqBrokerStandaloneOptions.cpuprofile, *mqBrokerStandaloneOptions.memprofile)
  52. grpcDialOption := security.LoadClientTLS(util.GetViper(), "grpc.msg_broker")
  53. qs, err := broker.NewMessageBroker(&broker.MessageQueueBrokerOption{
  54. Masters: mqBrokerOpt.masters,
  55. FilerGroup: *mqBrokerOpt.filerGroup,
  56. DataCenter: *mqBrokerOpt.dataCenter,
  57. Rack: *mqBrokerOpt.rack,
  58. DefaultReplication: "",
  59. MaxMB: 0,
  60. Ip: *mqBrokerOpt.ip,
  61. Port: *mqBrokerOpt.port,
  62. }, grpcDialOption)
  63. if err != nil {
  64. glog.Fatalf("failed to create new message broker for queue server: %v", err)
  65. }
  66. // start grpc listener
  67. grpcL, _, err := util.NewIpAndLocalListeners("", *mqBrokerOpt.port, 0)
  68. if err != nil {
  69. glog.Fatalf("failed to listen on grpc port %d: %v", *mqBrokerOpt.port, err)
  70. }
  71. grpcS := pb.NewGrpcServer(security.LoadServerTLS(util.GetViper(), "grpc.msg_broker"))
  72. mq_pb.RegisterSeaweedMessagingServer(grpcS, qs)
  73. reflection.Register(grpcS)
  74. grpcS.Serve(grpcL)
  75. return true
  76. }