diff options
Diffstat (limited to 'weed/mq/client/pub_client')
| -rw-r--r-- | weed/mq/client/pub_client/publish.go | 1 | ||||
| -rw-r--r-- | weed/mq/client/pub_client/scheduler.go | 10 |
2 files changed, 6 insertions, 5 deletions
diff --git a/weed/mq/client/pub_client/publish.go b/weed/mq/client/pub_client/publish.go index a25620de1..a85eec31f 100644 --- a/weed/mq/client/pub_client/publish.go +++ b/weed/mq/client/pub_client/publish.go @@ -51,6 +51,7 @@ func (p *TopicPublisher) FinishPublish() error { TsNs: time.Now().UnixNano(), Ctrl: &mq_pb.ControlMessage{ IsClose: true, + PublisherName: p.config.PublisherName, }, }) } diff --git a/weed/mq/client/pub_client/scheduler.go b/weed/mq/client/pub_client/scheduler.go index 03377d653..df2270b2c 100644 --- a/weed/mq/client/pub_client/scheduler.go +++ b/weed/mq/client/pub_client/scheduler.go @@ -142,11 +142,11 @@ func (p *TopicPublisher) doPublishToPartition(job *EachPartitionPublishJob) erro if err = publishClient.Send(&mq_pb.PublishMessageRequest{ Message: &mq_pb.PublishMessageRequest_Init{ Init: &mq_pb.PublishMessageRequest_InitMessage{ - Topic: p.config.Topic.ToPbTopic(), - Partition: job.Partition, - AckInterval: 128, - FollowerBrokers: job.FollowerBrokers, - PublisherName: p.config.PublisherName, + Topic: p.config.Topic.ToPbTopic(), + Partition: job.Partition, + AckInterval: 128, + FollowerBroker: job.FollowerBroker, + PublisherName: p.config.PublisherName, }, }, }); err != nil { |
