This is an automated email from the ASF dual-hosted git repository.
RongtongJin pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/rocketmq-clients.git
The following commit(s) were added to refs/heads/master by this push:
new 5939e374 [Go] Fix resource leaks on client close (#1305)
5939e374 is described below
commit 5939e374500b4b1572acd23979a61a5f6426f876
Author: guyinyou <[email protected]>
AuthorDate: Thu Jul 23 13:45:31 2026 +0800
[Go] Fix resource leaks on client close (#1305)
1. ClientManager shutdown: call shutdown() instead of UnRegisterClient()
in GracefulStop to stop background goroutines and close RPC connections
(regression from #235 which made ClientManager per-client but forgot
to add shutdown call).
2. Telemetry session cleanup: release all sessions in GracefulStop before
shutting down the client manager to properly close gRPC streams.
3. release() nil check: guard against nil observer when telemetry
connection failed during startup.
4. Per-client metrics isolation: use view.NewMeter() to give each client
its own Meter, Measures, Views, and Exporter. On client close, the
Meter is stopped and views are unregistered, preventing metrics
cardinality explosion from accumulated client_id tag values.
The client_id tag is preserved in exported metric data.
5. Fix meter lifecycle: call Start() before Register() to avoid deadlock,
and register/unregister the exporter with the per-client Meter so
metrics are actually exported.
Co-authored-by: guyinyou <[email protected]>
---
golang/client.go | 16 ++-
golang/client_manager.go | 1 +
golang/client_manager_mock.go | 12 +++
golang/lite_push_consumer_test.go | 2 +
golang/metric.go | 202 +++++++++++++++++++++++++-------------
5 files changed, 161 insertions(+), 72 deletions(-)
diff --git a/golang/client.go b/golang/client.go
index 471541d4..d58b0545 100644
--- a/golang/client.go
+++ b/golang/client.go
@@ -181,10 +181,12 @@ func (cs *defaultClientSession)
handleTelemetryCommand(response *v2.TelemetryCom
func (cs *defaultClientSession) release() {
cs.observerLock.Lock()
defer cs.observerLock.Unlock()
- if err := cs.observer.CloseSend(); err != nil {
- cs.cli.log.Errorf("release defaultClientSession err=%v", err)
+ if cs.observer != nil {
+ if err := cs.observer.CloseSend(); err != nil {
+ cs.cli.log.Errorf("release defaultClientSession
err=%v", err)
+ }
+ cs.observer = nil
}
- cs.observer = nil
}
func (cs *defaultClientSession) publish(ctx context.Context, common
*v2.TelemetryCommand) error {
var err error
@@ -646,8 +648,14 @@ func (cli *defaultClient) GracefulStop() error {
return fmt.Errorf("client has been closed")
}
cli.notifyClientTermination()
+ cli.endpointsTelemetryClientsLock.Lock()
+ for _, session := range cli.endpointsTelemetryClientTable {
+ session.release()
+ }
+ cli.endpointsTelemetryClientTable =
make(map[string]*defaultClientSession)
+ cli.endpointsTelemetryClientsLock.Unlock()
if cli.clientManager != nil {
- cli.clientManager.UnRegisterClient(cli)
+ cli.clientManager.shutdown()
}
cli.done <- struct{}{}
close(cli.done)
diff --git a/golang/client_manager.go b/golang/client_manager.go
index 0c395c4b..242a4e6a 100644
--- a/golang/client_manager.go
+++ b/golang/client_manager.go
@@ -46,6 +46,7 @@ type ClientManager interface {
ForwardMessageToDeadLetterQueue(ctx context.Context, endpoints
*v2.Endpoints, request *v2.ForwardMessageToDeadLetterQueueRequest, duration
time.Duration) (*v2.ForwardMessageToDeadLetterQueueResponse, error)
SyncLiteSubscription(ctx context.Context, endpoints *v2.Endpoints,
request *v2.SyncLiteSubscriptionRequest, duration time.Duration)
(*v2.SyncLiteSubscriptionResponse, error)
RecallMessage(ctx context.Context, endpoints *v2.Endpoints, request
*v2.RecallMessageRequest, duration time.Duration) (*v2.RecallMessageResponse,
error)
+ shutdown()
}
type clientManagerOptions struct {
diff --git a/golang/client_manager_mock.go b/golang/client_manager_mock.go
index 70daf72b..34c9c836 100644
--- a/golang/client_manager_mock.go
+++ b/golang/client_manager_mock.go
@@ -268,3 +268,15 @@ func (mr *MockClientManagerMockRecorder)
UnRegisterClient(client interface{}) *g
mr.mock.ctrl.T.Helper()
return mr.mock.ctrl.RecordCallWithMethodType(mr.mock,
"UnRegisterClient", reflect.TypeOf((*MockClientManager)(nil).UnRegisterClient),
client)
}
+
+// shutdown mocks base method.
+func (m *MockClientManager) shutdown() {
+ m.ctrl.T.Helper()
+ m.ctrl.Call(m, "shutdown")
+}
+
+// shutdown indicates an expected call of shutdown.
+func (mr *MockClientManagerMockRecorder) shutdown() *gomock.Call {
+ mr.mock.ctrl.T.Helper()
+ return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "shutdown",
reflect.TypeOf((*MockClientManager)(nil).shutdown))
+}
diff --git a/golang/lite_push_consumer_test.go
b/golang/lite_push_consumer_test.go
index a51fc0f2..9a580f28 100644
--- a/golang/lite_push_consumer_test.go
+++ b/golang/lite_push_consumer_test.go
@@ -181,6 +181,8 @@ func (m *mockedClientManager) RecallMessage(ctx
context.Context, endpoints *v2.E
return nil, nil
}
+func (m *mockedClientManager) shutdown() {}
+
func TestLitePushConsumer_SubscribeLite(t *testing.T) {
setupTest(t)
defer teardownTest()
diff --git a/golang/metric.go b/golang/metric.go
index 86755798..3d876f8b 100644
--- a/golang/metric.go
+++ b/golang/metric.go
@@ -19,6 +19,7 @@ package golang
import (
"context"
+ "fmt"
"sync"
"time"
@@ -46,57 +47,113 @@ var (
clientIdTag, _ = tag.NewKey("client_id")
invocationStatusTag, _ = tag.NewKey("invocation_status")
consumerGroupTag, _ = tag.NewKey("consumer_group")
+)
- PublishMLatencyMs = stats.Int64("publish_latency", "Publish
latency in milliseconds", "ms")
- ConsumeDeliveryMLatencyMs = stats.Int64("delivery_latency", "Time spent
delivering messages from servers to clients", "ms")
- ConsumeAwaitMLatencyMs = stats.Int64("await_time", "Client side
queuing time of messages before getting processed", "ms")
- ConsumeProcessMLatencyMs = stats.Int64("process_time", "Process
message time", "ms")
+type meterType int
- PublishLatencyView = view.View{
- Name: "rocketmq_send_cost_time",
- Description: "Publish latency",
- Measure: PublishMLatencyMs,
- Aggregation: view.Distribution(1, 5, 10, 20, 50, 200, 500),
- TagKeys: []tag.Key{topicTag, clientIdTag,
invocationStatusTag},
- }
+const (
+ meterPublishLatency meterType = iota
+ meterDeliveryLatency
+ meterAwaitTime
+ meterProcessTime
+)
- ConsumeDeliveryLatencyView = view.View{
- Name: "rocketmq_delivery_latency",
- Description: "Message delivery latency",
- Measure: ConsumeDeliveryMLatencyMs,
- Aggregation: view.Distribution(1, 5, 10, 20, 50, 200, 500),
- TagKeys: []tag.Key{topicTag, clientIdTag, consumerGroupTag},
- }
+type defaultClientMeter struct {
+ enabled atomic.Bool
+ endpoints *v2.Endpoints
+ ocaExporter view.Exporter
+ mutex sync.Mutex
- ConsumeAwaitTimeView = view.View{
- Name: "rocketmq_await_time",
- Description: "Message await time",
- Measure: ConsumeAwaitMLatencyMs,
- Aggregation: view.Distribution(1, 5, 20, 100, 1000, 5000,
10000),
- TagKeys: []tag.Key{topicTag, clientIdTag, consumerGroupTag},
- }
+ meter view.Meter
+ publishMeasure *stats.Int64Measure
+ deliveryMeasure *stats.Int64Measure
+ awaitMeasure *stats.Int64Measure
+ processMeasure *stats.Int64Measure
+ registeredViews []*view.View
+}
- ConsumeProcessTimeView = view.View{
- Name: "rocketmq_process_time",
- Description: "Message process time",
- Measure: ConsumeProcessMLatencyMs,
- Aggregation: view.Distribution(1, 5, 10, 100, 1000, 10000,
60000),
- TagKeys: []tag.Key{topicTag, clientIdTag, consumerGroupTag,
invocationStatusTag},
+func newDefaultClientMeter(enabled bool, exporter view.Exporter, endpoints
*v2.Endpoints, clientID string) *defaultClientMeter {
+ dcm := &defaultClientMeter{
+ enabled: *atomic.NewBool(enabled),
+ endpoints: endpoints,
+ ocaExporter: exporter,
}
-)
+ if enabled {
+ dcm.initPerClientResources(clientID)
+ }
+ return dcm
+}
-func init() {
- if err := view.Register(&PublishLatencyView,
&ConsumeDeliveryLatencyView, &ConsumeAwaitTimeView, &ConsumeProcessTimeView);
err != nil {
- sugarBaseLogger.Fatalf("failed to register views: %v", err)
+func (dcm *defaultClientMeter) initPerClientResources(clientID string) {
+ dcm.meter = view.NewMeter()
+
+ prefix := fmt.Sprintf("rocketmq_%s", clientID)
+ dcm.publishMeasure = stats.Int64(prefix+"_publish_latency", "Publish
latency in milliseconds", "ms")
+ dcm.deliveryMeasure = stats.Int64(prefix+"_delivery_latency", "Time
spent delivering messages from servers to clients", "ms")
+ dcm.awaitMeasure = stats.Int64(prefix+"_await_time", "Client side
queuing time of messages before getting processed", "ms")
+ dcm.processMeasure = stats.Int64(prefix+"_process_time", "Process
message time", "ms")
+
+ dcm.registeredViews = []*view.View{
+ {
+ Name: "rocketmq_send_cost_time",
+ Description: "Publish latency",
+ Measure: dcm.publishMeasure,
+ Aggregation: view.Distribution(1, 5, 10, 20, 50, 200,
500),
+ TagKeys: []tag.Key{topicTag, clientIdTag,
invocationStatusTag},
+ },
+ {
+ Name: "rocketmq_delivery_latency",
+ Description: "Message delivery latency",
+ Measure: dcm.deliveryMeasure,
+ Aggregation: view.Distribution(1, 5, 10, 20, 50, 200,
500),
+ TagKeys: []tag.Key{topicTag, clientIdTag,
consumerGroupTag},
+ },
+ {
+ Name: "rocketmq_await_time",
+ Description: "Message await time",
+ Measure: dcm.awaitMeasure,
+ Aggregation: view.Distribution(1, 5, 20, 100, 1000,
5000, 10000),
+ TagKeys: []tag.Key{topicTag, clientIdTag,
consumerGroupTag},
+ },
+ {
+ Name: "rocketmq_process_time",
+ Description: "Message process time",
+ Measure: dcm.processMeasure,
+ Aggregation: view.Distribution(1, 5, 10, 100, 1000,
10000, 60000),
+ TagKeys: []tag.Key{topicTag, clientIdTag,
consumerGroupTag, invocationStatusTag},
+ },
+ }
+
+ dcm.meter.Start()
+ if err := dcm.meter.Register(dcm.registeredViews...); err != nil {
+ sugarBaseLogger.Errorf("failed to register per-client views:
%v", err)
}
- view.SetReportingPeriod(time.Minute)
}
-type defaultClientMeter struct {
- enabled atomic.Bool
- endpoints *v2.Endpoints
- ocaExporter view.Exporter
- mutex sync.Mutex
+func (dcm *defaultClientMeter) record(mt meterType, mutators []tag.Mutator,
val int64) {
+ if !dcm.enabled.Load() || dcm.meter == nil {
+ return
+ }
+ var measure *stats.Int64Measure
+ switch mt {
+ case meterPublishLatency:
+ measure = dcm.publishMeasure
+ case meterDeliveryLatency:
+ measure = dcm.deliveryMeasure
+ case meterAwaitTime:
+ measure = dcm.awaitMeasure
+ case meterProcessTime:
+ measure = dcm.processMeasure
+ default:
+ return
+ }
+ ctx, err := tag.New(context.Background(), mutators...)
+ if err != nil {
+ sugarBaseLogger.Errorf("failed to create tag map: %v", err)
+ return
+ }
+ tagMap := tag.FromContext(ctx)
+ dcm.meter.Record(tagMap, []stats.Measurement{measure.M(val)}, nil)
}
func (dcm *defaultClientMeter) shutdown() {
@@ -105,31 +162,38 @@ func (dcm *defaultClientMeter) shutdown() {
}
dcm.mutex.Lock()
defer dcm.mutex.Unlock()
- view.UnregisterExporter(dcm.ocaExporter)
+
+ if dcm.meter != nil {
+ if dcm.ocaExporter != nil {
+ dcm.meter.UnregisterExporter(dcm.ocaExporter)
+ }
+ dcm.meter.Unregister(dcm.registeredViews...)
+ dcm.meter.Stop()
+ dcm.meter = nil
+ }
+
if dcm.ocaExporter != nil {
- exporter, ok := dcm.ocaExporter.(*ocagent.Exporter)
- if ok {
- err := exporter.Stop()
- if err != nil {
+ if exporter, ok := dcm.ocaExporter.(*ocagent.Exporter); ok {
+ if err := exporter.Stop(); err != nil {
sugarBaseLogger.Errorf("ocExporter stop failed,
err=%w", err)
}
}
+ dcm.ocaExporter = nil
}
+ dcm.enabled.Store(false)
}
func (dcm *defaultClientMeter) start() {
if !dcm.enabled.Load() {
return
}
- view.RegisterExporter(dcm.ocaExporter)
+ if dcm.meter != nil && dcm.ocaExporter != nil {
+ dcm.meter.RegisterExporter(dcm.ocaExporter)
+ }
}
var NewDefaultClientMeter = func(exporter view.Exporter, on bool, endpoints
*v2.Endpoints, clientID string) *defaultClientMeter {
- return &defaultClientMeter{
- enabled: *atomic.NewBool(on),
- endpoints: endpoints,
- ocaExporter: exporter,
- }
+ return newDefaultClientMeter(on, exporter, endpoints, clientID)
}
type MessageMeterInterceptor interface {
@@ -145,6 +209,7 @@ type ClientMeterProvider interface {
isEnabled() bool
getClientID() string
getClientImpl() isClient
+ record(mt meterType, tags []tag.Mutator, val int64)
}
var _ = ClientMeterProvider(&defaultClientMeterProvider{})
@@ -155,6 +220,10 @@ type defaultClientMeterProvider struct {
globalMutex sync.Mutex
}
+func (dcmp *defaultClientMeterProvider) record(mt meterType, tags
[]tag.Mutator, val int64) {
+ dcmp.clientMeter.record(mt, tags, val)
+}
+
func (dcmp *defaultClientMeterProvider) getClientImpl() isClient {
if dc, ok := dcmp.client.(*defaultClient); ok {
return dc.clientImpl
@@ -195,10 +264,9 @@ func (dmmi *defaultMessageMeterInterceptor)
doBeforeConsumeMessage(messageCommon
continue
}
duration := time.Since(*messageCommon.decodeStopwatch)
- err := stats.RecordWithTags(context.Background(),
[]tag.Mutator{tag.Insert(topicTag, messageCommon.topic),
tag.Insert(clientIdTag, dmmi.clientMeterProvider.getClientID()),
tag.Insert(consumerGroupTag, consumerGroup)},
ConsumeAwaitMLatencyMs.M(duration.Milliseconds()))
- if err != nil {
- return err
- }
+ dmmi.clientMeterProvider.record(meterAwaitTime,
+ []tag.Mutator{tag.Insert(topicTag,
messageCommon.topic), tag.Insert(clientIdTag, clientId),
tag.Insert(consumerGroupTag, consumerGroup)},
+ duration.Milliseconds())
}
return nil
@@ -230,10 +298,9 @@ func (dmmi *defaultMessageMeterInterceptor)
doAfterConsumeMessage(messageCommons
invocationStatus = InvocationStatus_SUCCESS
}
for _, messageCommon := range messageCommons {
- err := stats.RecordWithTags(context.Background(),
[]tag.Mutator{tag.Insert(topicTag, messageCommon.topic),
tag.Insert(clientIdTag, dmmi.clientMeterProvider.getClientID()),
tag.Insert(consumerGroupTag, consumerGroup), tag.Insert(invocationStatusTag,
string(invocationStatus))}, ConsumeProcessMLatencyMs.M(duration.Milliseconds()))
- if err != nil {
- return err
- }
+ dmmi.clientMeterProvider.record(meterProcessTime,
+ []tag.Mutator{tag.Insert(topicTag,
messageCommon.topic), tag.Insert(clientIdTag, clientId),
tag.Insert(consumerGroupTag, consumerGroup), tag.Insert(invocationStatusTag,
string(invocationStatus))},
+ duration.Milliseconds())
}
return nil
@@ -265,10 +332,9 @@ func (dmmi *defaultMessageMeterInterceptor)
doAfterReceiveMessage(messageCommons
continue
}
latency := time.Since(*messageCommon.deliveryTimestamp)
- err := stats.RecordWithTags(context.Background(),
[]tag.Mutator{tag.Insert(topicTag, messageCommon.topic),
tag.Insert(clientIdTag, dmmi.clientMeterProvider.getClientID()),
tag.Insert(consumerGroupTag, consumerGroup)},
ConsumeDeliveryMLatencyMs.M(latency.Milliseconds()))
- if err != nil {
- return err
- }
+ dmmi.clientMeterProvider.record(meterDeliveryLatency,
+ []tag.Mutator{tag.Insert(topicTag,
messageCommon.topic), tag.Insert(clientIdTag, clientId),
tag.Insert(consumerGroupTag, consumerGroup)},
+ latency.Milliseconds())
}
return nil
@@ -292,11 +358,11 @@ func (dmmi *defaultMessageMeterInterceptor)
doAfterSendMessage(messageCommons []
if status == MessageHookPointsStatus_OK {
invocationStatus = InvocationStatus_SUCCESS
}
+ clientId := dmmi.clientMeterProvider.getClientID()
for _, messageCommon := range messageCommons {
- err := stats.RecordWithTags(context.Background(),
[]tag.Mutator{tag.Insert(topicTag, messageCommon.topic),
tag.Insert(clientIdTag, dmmi.clientMeterProvider.getClientID()),
tag.Insert(invocationStatusTag, string(invocationStatus))},
PublishMLatencyMs.M(duration.Milliseconds()))
- if err != nil {
- return err
- }
+ dmmi.clientMeterProvider.record(meterPublishLatency,
+ []tag.Mutator{tag.Insert(topicTag,
messageCommon.topic), tag.Insert(clientIdTag, clientId),
tag.Insert(invocationStatusTag, string(invocationStatus))},
+ duration.Milliseconds())
}
return nil
}