broker

package
v0.0.0-...-ee25ada Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Apr 30, 2024 License: Apache-2.0 Imports: 31 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type MessageQueueBroker

type MessageQueueBroker struct {
	mq_pb.UnimplementedSeaweedMessagingServer

	MasterClient *wdclient.MasterClient

	Balancer *pub_balancer.Balancer

	Coordinator *sub_coordinator.Coordinator
	// contains filtered or unexported fields
}

func NewMessageBroker

func NewMessageBroker(option *MessageQueueBrokerOption, grpcDialOption grpc.DialOption) (mqBroker *MessageQueueBroker, err error)

func (*MessageQueueBroker) AdjustedUrl

func (b *MessageQueueBroker) AdjustedUrl(location *filer_pb.Location) string

func (*MessageQueueBroker) AssignTopicPartitions

AssignTopicPartitions Runs on the assigned broker, to execute the topic partition assignment

func (*MessageQueueBroker) BalanceTopics

func (b *MessageQueueBroker) BalanceTopics(ctx context.Context, request *mq_pb.BalanceTopicsRequest) (resp *mq_pb.BalanceTopicsResponse, err error)

func (*MessageQueueBroker) BrokerConnectToBalancer

func (b *MessageQueueBroker) BrokerConnectToBalancer(brokerBalancer string, stopCh chan struct{}) error

BrokerConnectToBalancer connects to the broker balancer and sends stats

func (*MessageQueueBroker) ClosePublishers

func (b *MessageQueueBroker) ClosePublishers(ctx context.Context, request *mq_pb.ClosePublishersRequest) (resp *mq_pb.ClosePublishersResponse, err error)

func (*MessageQueueBroker) CloseSubscribers

func (b *MessageQueueBroker) CloseSubscribers(ctx context.Context, request *mq_pb.CloseSubscribersRequest) (resp *mq_pb.CloseSubscribersResponse, err error)

func (*MessageQueueBroker) ConfigureTopic

func (b *MessageQueueBroker) ConfigureTopic(ctx context.Context, request *mq_pb.ConfigureTopicRequest) (resp *mq_pb.ConfigureTopicResponse, err error)

ConfigureTopic Runs on any broker, but proxied to the balancer if not the balancer It generates an assignments based on existing allocations, and then assign the partitions to the brokers.

func (*MessageQueueBroker) FindBrokerLeader

func (*MessageQueueBroker) FollowInMemoryMessages

func (*MessageQueueBroker) GetDataCenter

func (b *MessageQueueBroker) GetDataCenter() string

func (*MessageQueueBroker) GetFiler

func (b *MessageQueueBroker) GetFiler() pb.ServerAddress

func (*MessageQueueBroker) GetOrGenLocalPartition

func (b *MessageQueueBroker) GetOrGenLocalPartition(t topic.Topic, partition topic.Partition) (localPartition *topic.LocalPartition, isGenerated bool, err error)

func (*MessageQueueBroker) KeepConnectedToBrokerBalancer

func (b *MessageQueueBroker) KeepConnectedToBrokerBalancer(newBrokerBalancerCh chan string)

func (*MessageQueueBroker) ListTopics

func (b *MessageQueueBroker) ListTopics(ctx context.Context, request *mq_pb.ListTopicsRequest) (resp *mq_pb.ListTopicsResponse, err error)

func (*MessageQueueBroker) LookupTopicBrokers

func (b *MessageQueueBroker) LookupTopicBrokers(ctx context.Context, request *mq_pb.LookupTopicBrokersRequest) (resp *mq_pb.LookupTopicBrokersResponse, err error)

LookupTopicBrokers returns the brokers that are serving the topic

func (*MessageQueueBroker) OnBrokerUpdate

func (b *MessageQueueBroker) OnBrokerUpdate(update *master_pb.ClusterNodeUpdate, startFrom time.Time)

func (*MessageQueueBroker) PublishFollowMe

func (*MessageQueueBroker) PublishMessage

func (*MessageQueueBroker) PublisherToPubBalancer

PublisherToPubBalancer receives connections from brokers and collects stats

func (*MessageQueueBroker) SubscribeMessage

func (*MessageQueueBroker) SubscriberToSubCoordinator

SubscriberToSubCoordinator coordinates the subscribers

func (*MessageQueueBroker) WithFilerClient

func (b *MessageQueueBroker) WithFilerClient(streamingMode bool, fn func(filer_pb.SeaweedFilerClient) error) error

type MessageQueueBrokerOption

type MessageQueueBrokerOption struct {
	Masters            map[string]pb.ServerAddress
	FilerGroup         string
	DataCenter         string
	Rack               string
	DefaultReplication string
	MaxMB              int
	Ip                 string
	Port               int
	Cipher             bool
	VolumeServerAccess string // how to access volume servers
}

func (*MessageQueueBrokerOption) BrokerAddress

func (option *MessageQueueBrokerOption) BrokerAddress() pb.ServerAddress

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL