aboutsummaryrefslogtreecommitdiff
path: root/weed/shell/command_mq_topic_create.go
diff options
context:
space:
mode:
authorchrislu <chris.lu@gmail.com>2023-09-24 14:22:11 -0700
committerchrislu <chris.lu@gmail.com>2023-09-24 22:00:43 -0700
commit0361c321b40e5b0ff13edf01c8bb1b095f612caf (patch)
tree9117922ea1663c9a81ada68dad10f35528ff6171 /weed/shell/command_mq_topic_create.go
parent0f8168c0c928bba3d2f48b0680d3bdce9c617559 (diff)
downloadseaweedfs-0361c321b40e5b0ff13edf01c8bb1b095f612caf.tar.xz
seaweedfs-0361c321b40e5b0ff13edf01c8bb1b095f612caf.zip
add CreateTopic API
Diffstat (limited to 'weed/shell/command_mq_topic_create.go')
-rw-r--r--weed/shell/command_mq_topic_create.go65
1 files changed, 65 insertions, 0 deletions
diff --git a/weed/shell/command_mq_topic_create.go b/weed/shell/command_mq_topic_create.go
new file mode 100644
index 000000000..b09f94451
--- /dev/null
+++ b/weed/shell/command_mq_topic_create.go
@@ -0,0 +1,65 @@
+package shell
+
+import (
+ "context"
+ "flag"
+ "fmt"
+ "github.com/seaweedfs/seaweedfs/weed/pb"
+ "github.com/seaweedfs/seaweedfs/weed/pb/mq_pb"
+ "io"
+)
+
+func init() {
+ Commands = append(Commands, &commandMqTopicCreate{})
+}
+
+type commandMqTopicCreate struct {
+}
+
+func (c *commandMqTopicCreate) Name() string {
+ return "mq.topic.create"
+}
+
+func (c *commandMqTopicCreate) Help() string {
+ return `create a topic with a given name
+
+ Example:
+ mq.topic.create -namespace <namespace> -topic <topic_name> -partition_count <partition_count>
+`
+}
+
+func (c *commandMqTopicCreate) Do(args []string, commandEnv *CommandEnv, writer io.Writer) error {
+
+ // parse parameters
+ mqCommand := flag.NewFlagSet(c.Name(), flag.ContinueOnError)
+ namespace := mqCommand.String("namespace", "", "namespace name")
+ topicName := mqCommand.String("topic", "", "topic name")
+ partitionCount := mqCommand.Int("partitionCount", 6, "partition count")
+ if err := mqCommand.Parse(args); err != nil {
+ return err
+ }
+
+ // find the broker balancer
+ brokerBalancer, err := findBrokerBalancer(commandEnv)
+ if err != nil {
+ return err
+ }
+ fmt.Fprintf(writer, "current balancer: %s\n", brokerBalancer)
+
+ // create topic
+ return pb.WithBrokerGrpcClient(false, brokerBalancer, commandEnv.option.GrpcDialOption, func(client mq_pb.SeaweedMessagingClient) error {
+ resp, err := client.CreateTopic(context.Background(), &mq_pb.CreateTopicRequest{
+ Topic: &mq_pb.Topic{
+ Namespace: *namespace,
+ Name: *topicName,
+ },
+ PartitionCount: int32(*partitionCount),
+ })
+ if err != nil {
+ return err
+ }
+ fmt.Fprintf(writer, "response: %+v\n", resp)
+ return nil
+ })
+
+}