From 085a86e47549a97a3a6e493058f27d650753cbf7 Mon Sep 17 00:00:00 2001 From: Brian Goff Date: Fri, 12 Feb 2016 10:20:16 -0500 Subject: [PATCH 1/2] Fix some issues with concurrency in aufs. Adds a benchmark to measure performance under concurrent actions. Signed-off-by: Brian Goff Upstream-commit: 55c91f2ab9bcd48cfa248a4e842bb78257c14134 Component: engine --- .../engine/daemon/graphdriver/aufs/aufs.go | 31 ++++---- .../daemon/graphdriver/aufs/aufs_test.go | 71 ++++++++++++++++++- 2 files changed, 84 insertions(+), 18 deletions(-) diff --git a/components/engine/daemon/graphdriver/aufs/aufs.go b/components/engine/daemon/graphdriver/aufs/aufs.go index 51054fa6ef..e03576aa8a 100644 --- a/components/engine/daemon/graphdriver/aufs/aufs.go +++ b/components/engine/daemon/graphdriver/aufs/aufs.go @@ -227,7 +227,9 @@ func (a *Driver) Create(id, parent, mountLabel string) error { } } } + a.Lock() a.active[id] = &data{} + a.Unlock() return nil } @@ -285,20 +287,17 @@ func (a *Driver) Remove(id string) error { if err := os.Remove(path.Join(a.rootPath(), "layers", id)); err != nil && !os.IsNotExist(err) { return err } + if m != nil { + a.Lock() + delete(a.active, id) + a.Unlock() + } return nil } // Get returns the rootfs path for the id. // This will mount the dir at it's given path func (a *Driver) Get(id, mountLabel string) (string, error) { - ids, err := getParentIds(a.rootPath(), id) - if err != nil { - if !os.IsNotExist(err) { - return "", err - } - ids = []string{} - } - // Protect the a.active from concurrent access a.Lock() defer a.Unlock() @@ -309,13 +308,18 @@ func (a *Driver) Get(id, mountLabel string) (string, error) { a.active[id] = m } + parents, err := a.getParentLayerPaths(id) + if err != nil && !os.IsNotExist(err) { + return "", err + } + // If a dir does not have a parent ( no layers )do not try to mount // just return the diff path to the data m.path = path.Join(a.rootPath(), "diff", id) - if len(ids) > 0 { + if len(parents) > 0 { m.path = path.Join(a.rootPath(), "mnt", id) if m.referenceCount == 0 { - if err := a.mount(id, m, mountLabel); err != nil { + if err := a.mount(id, m, mountLabel, parents); err != nil { return "", err } } @@ -426,7 +430,7 @@ func (a *Driver) getParentLayerPaths(id string) ([]string, error) { return layers, nil } -func (a *Driver) mount(id string, m *data, mountLabel string) error { +func (a *Driver) mount(id string, m *data, mountLabel string, layers []string) error { // If the id is mounted or we get an error return if mounted, err := a.mounted(m); err != nil || mounted { return err @@ -437,11 +441,6 @@ func (a *Driver) mount(id string, m *data, mountLabel string) error { rw = path.Join(a.rootPath(), "diff", id) ) - layers, err := a.getParentLayerPaths(id) - if err != nil { - return err - } - if err := a.aufsMount(layers, rw, target, mountLabel); err != nil { return fmt.Errorf("error creating aufs mount to %s: %v", target, err) } diff --git a/components/engine/daemon/graphdriver/aufs/aufs_test.go b/components/engine/daemon/graphdriver/aufs/aufs_test.go index 761b5b6872..0f6d59d054 100644 --- a/components/engine/daemon/graphdriver/aufs/aufs_test.go +++ b/components/engine/daemon/graphdriver/aufs/aufs_test.go @@ -9,11 +9,13 @@ import ( "io/ioutil" "os" "path" + "sync" "testing" "github.com/docker/docker/daemon/graphdriver" "github.com/docker/docker/pkg/archive" "github.com/docker/docker/pkg/reexec" + "github.com/docker/docker/pkg/stringid" ) var ( @@ -25,7 +27,7 @@ func init() { reexec.Init() } -func testInit(dir string, t *testing.T) graphdriver.Driver { +func testInit(dir string, t testing.TB) graphdriver.Driver { d, err := Init(dir, nil, nil, nil) if err != nil { if err == graphdriver.ErrNotSupported { @@ -37,7 +39,7 @@ func testInit(dir string, t *testing.T) graphdriver.Driver { return d } -func newDriver(t *testing.T) *Driver { +func newDriver(t testing.TB) *Driver { if err := os.MkdirAll(tmp, 0755); err != nil { t.Fatal(err) } @@ -732,3 +734,68 @@ func TestMountMoreThan42LayersMatchingPathLength(t *testing.T) { zeroes += "0" } } + +func BenchmarkConcurrentAccess(b *testing.B) { + b.StopTimer() + b.ResetTimer() + + d := newDriver(b) + defer os.RemoveAll(tmp) + defer d.Cleanup() + + numConcurent := 256 + // create a bunch of ids + var ids []string + for i := 0; i < numConcurent; i++ { + ids = append(ids, stringid.GenerateNonCryptoID()) + } + + if err := d.Create(ids[0], "", ""); err != nil { + b.Fatal(err) + } + + if err := d.Create(ids[1], ids[0], ""); err != nil { + b.Fatal(err) + } + + parent := ids[1] + ids = append(ids[2:]) + + chErr := make(chan error, numConcurent) + var outerGroup sync.WaitGroup + outerGroup.Add(len(ids)) + b.StartTimer() + + // here's the actual bench + for _, id := range ids { + go func(id string) { + defer outerGroup.Done() + if err := d.Create(id, parent, ""); err != nil { + b.Logf("Create %s failed", id) + chErr <- err + return + } + var innerGroup sync.WaitGroup + for i := 0; i < b.N; i++ { + innerGroup.Add(1) + go func() { + d.Get(id, "") + d.Put(id) + innerGroup.Done() + }() + } + innerGroup.Wait() + d.Remove(id) + }(id) + } + + outerGroup.Wait() + b.StopTimer() + close(chErr) + for err := range chErr { + if err != nil { + b.Log(err) + b.Fail() + } + } +} From ac8b4b9a6a663f211348420825ba4bf8b096b53b Mon Sep 17 00:00:00 2001 From: Brian Goff Date: Fri, 12 Feb 2016 11:01:45 -0500 Subject: [PATCH 2/2] Add finer-grained locking for aufs ``` benchmark old ns/op new ns/op delta BenchmarkConcurrentAccess-8 10269529748 26834747 -99.74% benchmark old allocs new allocs delta BenchmarkConcurrentAccess-8 309948 7232 -97.67% benchmark old bytes new bytes delta BenchmarkConcurrentAccess-8 23943576 1578441 -93.41% ``` Signed-off-by: Brian Goff Upstream-commit: f31014197cbe9438cc956ed12c47093a0324c82d Component: engine --- .../engine/daemon/graphdriver/aufs/aufs.go | 94 +++++++++++++------ .../engine/daemon/graphdriver/aufs/mount.go | 2 +- 2 files changed, 67 insertions(+), 29 deletions(-) diff --git a/components/engine/daemon/graphdriver/aufs/aufs.go b/components/engine/daemon/graphdriver/aufs/aufs.go index e03576aa8a..529d44c265 100644 --- a/components/engine/daemon/graphdriver/aufs/aufs.go +++ b/components/engine/daemon/graphdriver/aufs/aufs.go @@ -66,6 +66,7 @@ func init() { type data struct { referenceCount int path string + sync.Mutex } // Driver contains information about the filesystem mounted. @@ -76,7 +77,7 @@ type Driver struct { root string uidMaps []idtools.IDMap gidMaps []idtools.IDMap - sync.Mutex // Protects concurrent modification to active + globalLock sync.Mutex // Protects concurrent modification to active active map[string]*data } @@ -202,7 +203,20 @@ func (a *Driver) Exists(id string) bool { // Create three folders for each id // mnt, layers, and diff func (a *Driver) Create(id, parent, mountLabel string) error { - if err := a.createDirsFor(id); err != nil { + m := a.getActive(id) + m.Lock() + + var err error + defer func() { + a.globalLock.Lock() + if err != nil { + delete(a.active, id) + } + a.globalLock.Unlock() + m.Unlock() + }() + + if err = a.createDirsFor(id); err != nil { return err } // Write the layers metadata @@ -213,23 +227,22 @@ func (a *Driver) Create(id, parent, mountLabel string) error { defer f.Close() if parent != "" { - ids, err := getParentIds(a.rootPath(), parent) + var ids []string + ids, err = getParentIds(a.rootPath(), parent) if err != nil { return err } - if _, err := fmt.Fprintln(f, parent); err != nil { + if _, err = fmt.Fprintln(f, parent); err != nil { return err } for _, i := range ids { - if _, err := fmt.Fprintln(f, i); err != nil { + if _, err = fmt.Fprintln(f, i); err != nil { return err } } } - a.Lock() - a.active[id] = &data{} - a.Unlock() + return nil } @@ -253,11 +266,10 @@ func (a *Driver) createDirsFor(id string) error { // Remove will unmount and remove the given id. func (a *Driver) Remove(id string) error { - // Protect the a.active from concurrent access - a.Lock() - defer a.Unlock() + m := a.getActive(id) + m.Lock() + defer m.Unlock() - m := a.active[id] if m != nil { if m.referenceCount > 0 { return nil @@ -288,9 +300,9 @@ func (a *Driver) Remove(id string) error { return err } if m != nil { - a.Lock() + a.globalLock.Lock() delete(a.active, id) - a.Unlock() + a.globalLock.Unlock() } return nil } @@ -298,21 +310,36 @@ func (a *Driver) Remove(id string) error { // Get returns the rootfs path for the id. // This will mount the dir at it's given path func (a *Driver) Get(id, mountLabel string) (string, error) { - // Protect the a.active from concurrent access - a.Lock() - defer a.Unlock() - - m := a.active[id] - if m == nil { - m = &data{} - a.active[id] = m - } + m := a.getActive(id) + m.Lock() + defer m.Unlock() parents, err := a.getParentLayerPaths(id) if err != nil && !os.IsNotExist(err) { return "", err } + var parentLocks []*data + a.globalLock.Lock() + for _, p := range parents { + parentM, exists := a.active[p] + if !exists { + parentM = &data{} + a.active[p] = parentM + } + parentLocks = append(parentLocks, parentM) + } + a.globalLock.Unlock() + + for _, l := range parentLocks { + l.Lock() + } + defer func() { + for _, l := range parentLocks { + l.Unlock() + } + }() + // If a dir does not have a parent ( no layers )do not try to mount // just return the diff path to the data m.path = path.Join(a.rootPath(), "diff", id) @@ -328,13 +355,24 @@ func (a *Driver) Get(id, mountLabel string) (string, error) { return m.path, nil } +func (a *Driver) getActive(id string) *data { + // Protect the a.active from concurrent access + a.globalLock.Lock() + m, exists := a.active[id] + if !exists { + m = &data{} + a.active[id] = m + } + a.globalLock.Unlock() + return m +} + // Put unmounts and updates list of active mounts. func (a *Driver) Put(id string) error { - // Protect the a.active from concurrent access - a.Lock() - defer a.Unlock() + m := a.getActive(id) + m.Lock() + defer m.Unlock() - m := a.active[id] if m == nil { // but it might be still here if a.Exists(id) { @@ -346,6 +384,7 @@ func (a *Driver) Put(id string) error { } return nil } + if count := m.referenceCount; count > 1 { m.referenceCount = count - 1 } else { @@ -354,7 +393,6 @@ func (a *Driver) Put(id string) error { if ids != nil && len(ids) > 0 { a.unmount(m) } - delete(a.active, id) } return nil } diff --git a/components/engine/daemon/graphdriver/aufs/mount.go b/components/engine/daemon/graphdriver/aufs/mount.go index d7e9bf9fd7..36fa62e41b 100644 --- a/components/engine/daemon/graphdriver/aufs/mount.go +++ b/components/engine/daemon/graphdriver/aufs/mount.go @@ -12,7 +12,7 @@ import ( // Unmount the target specified. func Unmount(target string) error { if err := exec.Command("auplink", target, "flush").Run(); err != nil { - logrus.Errorf("Couldn't run auplink before unmount: %s", err) + logrus.Errorf("Couldn't run auplink before unmount %s: %s", target, err) } if err := syscall.Unmount(target, 0); err != nil { return err