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
 }

Reply via email to