aboutsummaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorchrislu <chris.lu@gmail.com>2023-09-25 20:46:00 -0700
committerchrislu <chris.lu@gmail.com>2023-09-25 20:46:00 -0700
commit19505c1cf44ce7c95032b3b007a200862a5c5c47 (patch)
tree5b60f627814950e868c582d766c76201854e839b
parent78dbac77021f72cc005c45ef68cf0793d4607232 (diff)
downloadseaweedfs-19505c1cf44ce7c95032b3b007a200862a5c5c47.tar.xz
seaweedfs-19505c1cf44ce7c95032b3b007a200862a5c5c47.zip
describe a topic
-rw-r--r--weed/shell/command_mq_topic_desc.go59
1 files changed, 59 insertions, 0 deletions
diff --git a/weed/shell/command_mq_topic_desc.go b/weed/shell/command_mq_topic_desc.go
new file mode 100644
index 000000000..a4bf805f9
--- /dev/null
+++ b/weed/shell/command_mq_topic_desc.go
@@ -0,0 +1,59 @@
+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, &commandMqTopicDescribe{})
+}
+
+type commandMqTopicDescribe struct {
+}
+
+func (c *commandMqTopicDescribe) Name() string {
+ return "mq.topic.describe"
+}
+
+func (c *commandMqTopicDescribe) Help() string {
+ return `describe a topic`
+}
+
+func (c *commandMqTopicDescribe) 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")
+ 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)
+
+ return pb.WithBrokerGrpcClient(false, brokerBalancer, commandEnv.option.GrpcDialOption, func(client mq_pb.SeaweedMessagingClient) error {
+ resp, err := client.LookupTopicBrokers(context.Background(), &mq_pb.LookupTopicBrokersRequest{
+ Topic: &mq_pb.Topic{
+ Namespace: *namespace,
+ Name: *topicName,
+ },
+ IsForPublish: false,
+ })
+ if err != nil {
+ return err
+ }
+ for _, assignment := range resp.BrokerPartitionAssignments {
+ fmt.Fprintf(writer, " %+v\n", assignment)
+ }
+ return nil
+ })
+}