This is an automated email from the ASF dual-hosted git repository.
AlexStocks pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/dubbo-go.git
The following commit(s) were added to refs/heads/main by this push:
new 7125a09bd perf(metadata): pipeline ZooKeeper revision reads (#3602)
7125a09bd is described below
commit 7125a09bdbfdadd1c5922d669730e29d7416c395
Author: xiaobaicai66695 <[email protected]>
AuthorDate: Wed Aug 12 15:43:33 2026 +0800
perf(metadata): pipeline ZooKeeper revision reads (#3602)
* perf(metadata): pipeline ZooKeeper revision reads
---
common/constant/default.go | 3 ++
metadata/report/zookeeper/report.go | 44 +++++++++++----
metadata/report/zookeeper/report_test.go | 92 ++++++++++++++++++++++++++++++++
3 files changed, 130 insertions(+), 9 deletions(-)
diff --git a/common/constant/default.go b/common/constant/default.go
index 4df46d956..a72f6d94d 100644
--- a/common/constant/default.go
+++ b/common/constant/default.go
@@ -95,6 +95,9 @@ const (
const (
SimpleMetadataServiceName = "MetadataService"
DefaultRevision = "N/A"
+
+ // ZookeeperListAppRevisionsMaxConcurrency bounds in-flight reads when
listing application revisions.
+ ZookeeperListAppRevisionsMaxConcurrency = 16
)
const (
diff --git a/metadata/report/zookeeper/report.go
b/metadata/report/zookeeper/report.go
index 8e594e365..39eb606f8 100644
--- a/metadata/report/zookeeper/report.go
+++ b/metadata/report/zookeeper/report.go
@@ -21,6 +21,7 @@ import (
"encoding/json"
"fmt"
"strings"
+ "sync"
)
import (
@@ -144,17 +145,42 @@ func (m *zookeeperMetadataReport)
ListAppRevisions(application string) ([]report
}
return nil, err
}
+ revisions := make([]report.AppRevision, len(children))
+ found := make([]bool, len(children))
+ jobs := make(chan int, len(children))
+ for i := range children {
+ jobs <- i
+ }
+ close(jobs)
+
+ var workers sync.WaitGroup
+ workerCount := min(len(children),
constant.ZookeeperListAppRevisionsMaxConcurrency)
+ workers.Add(workerCount)
+ for range workerCount {
+ go func() {
+ defer workers.Done()
+ for i := range jobs {
+ revision := children[i]
+ path := parent + constant.PathSeparator +
revision
+ data, _, getErr := m.client.Get(path)
+ if getErr != nil {
+ continue // skip if node disappeared
between listing and reading
+ }
+ revisions[i] = report.AppRevision{
+ Revision: revision,
+ ModifyTime:
report.ParseMetadataLastUpdatedTime(data),
+ }
+ found[i] = true
+ }
+ }()
+ }
+ workers.Wait()
+
result := make([]report.AppRevision, 0, len(children))
- for _, rev := range children {
- path := parent + constant.PathSeparator + rev
- data, _, err := m.client.Get(path)
- if err != nil {
- continue // skip if node disappeared between listing
and reading
+ for i := range revisions {
+ if found[i] {
+ result = append(result, revisions[i])
}
- result = append(result, report.AppRevision{
- Revision: rev,
- ModifyTime: report.ParseMetadataLastUpdatedTime(data),
- })
}
return result, nil
}
diff --git a/metadata/report/zookeeper/report_test.go
b/metadata/report/zookeeper/report_test.go
index 32fd4de50..a4a56d820 100644
--- a/metadata/report/zookeeper/report_test.go
+++ b/metadata/report/zookeeper/report_test.go
@@ -19,8 +19,11 @@ package zookeeper
import (
"encoding/json"
+ "fmt"
"strings"
+ "sync"
"testing"
+ "time"
)
import (
@@ -134,6 +137,9 @@ func (m *mockZkClient) Children(path string) ([]string,
*zk.Stat, error) {
}
func (m *mockZkClient) Get(path string) ([]byte, *zk.Stat, error) {
+ if err, ok := m.errors["Get:"+path]; ok {
+ return nil, nil, err
+ }
v, ok := m.data[path]
if !ok {
return nil, nil, zk.ErrNoNode
@@ -141,6 +147,33 @@ func (m *mockZkClient) Get(path string) ([]byte, *zk.Stat,
error) {
return v, m.stats[path], nil
}
+type concurrencyTrackingZkClient struct {
+ *mockZkClient
+ entered chan struct{}
+ release chan struct{}
+
+ mu sync.Mutex
+ active int
+ maxActive int
+}
+
+func (m *concurrencyTrackingZkClient) Get(path string) ([]byte, *zk.Stat,
error) {
+ m.mu.Lock()
+ m.active++
+ if m.active > m.maxActive {
+ m.maxActive = m.active
+ }
+ m.mu.Unlock()
+
+ m.entered <- struct{}{}
+ <-m.release
+
+ m.mu.Lock()
+ m.active--
+ m.mu.Unlock()
+ return m.mockZkClient.Get(path)
+}
+
// --- Helper ---
func newTestReportWithMock() (*zookeeperMetadataReport, *mockZkClient) {
@@ -255,6 +288,65 @@ func TestListAppRevisions(t *testing.T) {
assert.Equal(t, int64(2000), names["r3"])
}
+func TestListAppRevisionsPipelinesReadsWithBoundedConcurrency(t *testing.T) {
+ const extraRevisions = 5
+ revisionCount := constant.ZookeeperListAppRevisionsMaxConcurrency +
extraRevisions
+ mc := newMockZkClient()
+ client := &concurrencyTrackingZkClient{
+ mockZkClient: mc,
+ entered: make(chan struct{}, revisionCount),
+ release: make(chan struct{}),
+ }
+ for i := range revisionCount {
+ path := fmt.Sprintf("/dubbo/my-app/r%d", i)
+ mc.data[path] = fmt.Appendf(nil, `{"lastUpdatedTime":%d}`, i)
+ }
+
+ r := &zookeeperMetadataReport{client: client, rootDir: "/dubbo/"}
+ type listResult struct {
+ revisions []report.AppRevision
+ err error
+ }
+ done := make(chan listResult, 1)
+ go func() {
+ revisions, err := r.ListAppRevisions("my-app")
+ done <- listResult{revisions: revisions, err: err}
+ }()
+
+ for range constant.ZookeeperListAppRevisionsMaxConcurrency {
+ select {
+ case <-client.entered:
+ case <-time.After(time.Second):
+ t.Fatal("timed out waiting for concurrent ZooKeeper
reads")
+ }
+ }
+ select {
+ case <-client.entered:
+ t.Fatal("ListAppRevisions exceeded its concurrency limit")
+ default:
+ }
+ close(client.release)
+
+ result := <-done
+ require.NoError(t, result.err)
+ require.Len(t, result.revisions, revisionCount)
+ client.mu.Lock()
+ assert.Equal(t, constant.ZookeeperListAppRevisionsMaxConcurrency,
client.maxActive)
+ client.mu.Unlock()
+}
+
+func TestListAppRevisionsSkipsRevisionRemovedDuringRead(t *testing.T) {
+ r, mc := newTestReportWithMock()
+ mc.data["/dubbo/my-app/r1"] = []byte(`{"lastUpdatedTime":1000}`)
+ mc.data["/dubbo/my-app/r2"] = []byte(`{"lastUpdatedTime":2000}`)
+ mc.errors["Get:/dubbo/my-app/r2"] = zk.ErrNoNode
+
+ revisions, err := r.ListAppRevisions("my-app")
+ require.NoError(t, err)
+ require.Len(t, revisions, 1)
+ assert.Equal(t, "r1", revisions[0].Revision)
+}
+
func TestRegisterServiceAppMapping_NewKey(t *testing.T) {
r, mc := newTestReportWithMock()