mirror of https://github.com/docker/compose.git
Implement --scale option on up command, allow scale config in v2.2 format
docker-compose scale modified to reuse code between up and scale Signed-off-by: Joffrey F <joffrey@docker.com>
This commit is contained in:
parent
baf457c78c
commit
78ee612333
|
@ -771,15 +771,7 @@ class TopLevelCommand(object):
|
|||
"""
|
||||
timeout = timeout_from_opts(options)
|
||||
|
||||
for s in options['SERVICE=NUM']:
|
||||
if '=' not in s:
|
||||
raise UserError('Arguments to scale should be in the form service=num')
|
||||
service_name, num = s.split('=', 1)
|
||||
try:
|
||||
num = int(num)
|
||||
except ValueError:
|
||||
raise UserError('Number of containers for service "%s" is not a '
|
||||
'number' % service_name)
|
||||
for service_name, num in parse_scale_args(options['SERVICE=NUM']).items():
|
||||
self.project.get_service(service_name).scale(num, timeout=timeout)
|
||||
|
||||
def start(self, options):
|
||||
|
@ -875,7 +867,7 @@ class TopLevelCommand(object):
|
|||
If you want to force Compose to stop and recreate all containers, use the
|
||||
`--force-recreate` flag.
|
||||
|
||||
Usage: up [options] [SERVICE...]
|
||||
Usage: up [options] [--scale SERVICE=NUM...] [SERVICE...]
|
||||
|
||||
Options:
|
||||
-d Detached mode: Run containers in the background,
|
||||
|
@ -898,7 +890,9 @@ class TopLevelCommand(object):
|
|||
--remove-orphans Remove containers for services not
|
||||
defined in the Compose file
|
||||
--exit-code-from SERVICE Return the exit code of the selected service container.
|
||||
Requires --abort-on-container-exit.
|
||||
Implies --abort-on-container-exit.
|
||||
--scale SERVICE=NUM Scale SERVICE to NUM instances. Overrides the `scale`
|
||||
setting in the Compose file if present.
|
||||
"""
|
||||
start_deps = not options['--no-deps']
|
||||
exit_value_from = exitval_from_opts(options, self.project)
|
||||
|
@ -919,7 +913,9 @@ class TopLevelCommand(object):
|
|||
do_build=build_action_from_opts(options),
|
||||
timeout=timeout,
|
||||
detached=detached,
|
||||
remove_orphans=remove_orphans)
|
||||
remove_orphans=remove_orphans,
|
||||
scale_override=parse_scale_args(options['--scale']),
|
||||
)
|
||||
|
||||
if detached:
|
||||
return
|
||||
|
@ -1238,3 +1234,19 @@ def call_docker(args):
|
|||
log.debug(" ".join(map(pipes.quote, args)))
|
||||
|
||||
return subprocess.call(args)
|
||||
|
||||
|
||||
def parse_scale_args(options):
|
||||
res = {}
|
||||
for s in options:
|
||||
if '=' not in s:
|
||||
raise UserError('Arguments to scale should be in the form service=num')
|
||||
service_name, num = s.split('=', 1)
|
||||
try:
|
||||
num = int(num)
|
||||
except ValueError:
|
||||
raise UserError(
|
||||
'Number of containers for service "%s" is not a number' % service_name
|
||||
)
|
||||
res[service_name] = num
|
||||
return res
|
||||
|
|
|
@ -222,6 +222,7 @@
|
|||
"privileged": {"type": "boolean"},
|
||||
"read_only": {"type": "boolean"},
|
||||
"restart": {"type": "string"},
|
||||
"scale": {"type": "integer"},
|
||||
"security_opt": {"type": "array", "items": {"type": "string"}, "uniqueItems": true},
|
||||
"shm_size": {"type": ["number", "string"]},
|
||||
"sysctls": {"$ref": "#/definitions/list_or_dict"},
|
||||
|
|
|
@ -260,10 +260,6 @@ def parallel_remove(containers, options):
|
|||
parallel_operation(stopped_containers, 'remove', options, 'Removing')
|
||||
|
||||
|
||||
def parallel_start(containers, options):
|
||||
parallel_operation(containers, 'start', options, 'Starting')
|
||||
|
||||
|
||||
def parallel_pause(containers, options):
|
||||
parallel_operation(containers, 'pause', options, 'Pausing')
|
||||
|
||||
|
|
|
@ -380,13 +380,17 @@ class Project(object):
|
|||
do_build=BuildAction.none,
|
||||
timeout=None,
|
||||
detached=False,
|
||||
remove_orphans=False):
|
||||
remove_orphans=False,
|
||||
scale_override=None):
|
||||
|
||||
warn_for_swarm_mode(self.client)
|
||||
|
||||
self.initialize()
|
||||
self.find_orphan_containers(remove_orphans)
|
||||
|
||||
if scale_override is None:
|
||||
scale_override = {}
|
||||
|
||||
services = self.get_services_without_duplicate(
|
||||
service_names,
|
||||
include_deps=start_deps)
|
||||
|
@ -399,7 +403,8 @@ class Project(object):
|
|||
return service.execute_convergence_plan(
|
||||
plans[service.name],
|
||||
timeout=timeout,
|
||||
detached=detached
|
||||
detached=detached,
|
||||
scale_override=scale_override.get(service.name)
|
||||
)
|
||||
|
||||
def get_deps(service):
|
||||
|
@ -589,10 +594,13 @@ def get_secrets(service, service_secrets, secret_defs):
|
|||
continue
|
||||
|
||||
if secret.uid or secret.gid or secret.mode:
|
||||
log.warn("Service \"{service}\" uses secret \"{secret}\" with uid, "
|
||||
"gid, or mode. These fields are not supported by this "
|
||||
"implementation of the Compose file".format(
|
||||
service=service, secret=secret.source))
|
||||
log.warn(
|
||||
"Service \"{service}\" uses secret \"{secret}\" with uid, "
|
||||
"gid, or mode. These fields are not supported by this "
|
||||
"implementation of the Compose file".format(
|
||||
service=service, secret=secret.source
|
||||
)
|
||||
)
|
||||
|
||||
secrets.append({'secret': secret, 'file': secret_def.get('file')})
|
||||
|
||||
|
|
|
@ -38,7 +38,6 @@ from .errors import HealthCheckFailed
|
|||
from .errors import NoHealthCheckConfigured
|
||||
from .errors import OperationFailedError
|
||||
from .parallel import parallel_execute
|
||||
from .parallel import parallel_start
|
||||
from .progress_stream import stream_output
|
||||
from .progress_stream import StreamOutputError
|
||||
from .utils import json_hash
|
||||
|
@ -148,6 +147,7 @@ class Service(object):
|
|||
network_mode=None,
|
||||
networks=None,
|
||||
secrets=None,
|
||||
scale=None,
|
||||
**options
|
||||
):
|
||||
self.name = name
|
||||
|
@ -159,6 +159,7 @@ class Service(object):
|
|||
self.network_mode = network_mode or NetworkMode(None)
|
||||
self.networks = networks or {}
|
||||
self.secrets = secrets or []
|
||||
self.scale_num = scale or 1
|
||||
self.options = options
|
||||
|
||||
def __repr__(self):
|
||||
|
@ -189,16 +190,7 @@ class Service(object):
|
|||
self.start_container_if_stopped(c, **options)
|
||||
return containers
|
||||
|
||||
def scale(self, desired_num, timeout=None):
|
||||
"""
|
||||
Adjusts the number of containers to the specified number and ensures
|
||||
they are running.
|
||||
|
||||
- creates containers until there are at least `desired_num`
|
||||
- stops containers until there are at most `desired_num` running
|
||||
- starts containers until there are at least `desired_num` running
|
||||
- removes all stopped containers
|
||||
"""
|
||||
def show_scale_warnings(self, desired_num):
|
||||
if self.custom_container_name and desired_num > 1:
|
||||
log.warn('The "%s" service is using the custom container name "%s". '
|
||||
'Docker requires each container to have a unique name. '
|
||||
|
@ -210,14 +202,18 @@ class Service(object):
|
|||
'for this service are created on a single host, the port will clash.'
|
||||
% self.name)
|
||||
|
||||
def create_and_start(service, number):
|
||||
container = service.create_container(number=number, quiet=True)
|
||||
service.start_container(container)
|
||||
return container
|
||||
def scale(self, desired_num, timeout=None):
|
||||
"""
|
||||
Adjusts the number of containers to the specified number and ensures
|
||||
they are running.
|
||||
|
||||
def stop_and_remove(container):
|
||||
container.stop(timeout=self.stop_timeout(timeout))
|
||||
container.remove()
|
||||
- creates containers until there are at least `desired_num`
|
||||
- stops containers until there are at most `desired_num` running
|
||||
- starts containers until there are at least `desired_num` running
|
||||
- removes all stopped containers
|
||||
"""
|
||||
|
||||
self.show_scale_warnings(desired_num)
|
||||
|
||||
running_containers = self.containers(stopped=False)
|
||||
num_running = len(running_containers)
|
||||
|
@ -228,11 +224,10 @@ class Service(object):
|
|||
return
|
||||
|
||||
if desired_num > num_running:
|
||||
# we need to start/create until we have desired_num
|
||||
all_containers = self.containers(stopped=True)
|
||||
|
||||
if num_running != len(all_containers):
|
||||
# we have some stopped containers, let's start them up again
|
||||
# we have some stopped containers, check for divergences
|
||||
stopped_containers = [
|
||||
c for c in all_containers if not c.is_running
|
||||
]
|
||||
|
@ -241,38 +236,14 @@ class Service(object):
|
|||
divergent_containers = [
|
||||
c for c in stopped_containers if self._containers_have_diverged([c])
|
||||
]
|
||||
stopped_containers = sorted(
|
||||
set(stopped_containers) - set(divergent_containers),
|
||||
key=attrgetter('number')
|
||||
)
|
||||
for c in divergent_containers:
|
||||
c.remove()
|
||||
|
||||
num_stopped = len(stopped_containers)
|
||||
all_containers = list(set(all_containers) - set(divergent_containers))
|
||||
|
||||
if num_stopped + num_running > desired_num:
|
||||
num_to_start = desired_num - num_running
|
||||
containers_to_start = stopped_containers[:num_to_start]
|
||||
else:
|
||||
containers_to_start = stopped_containers
|
||||
|
||||
parallel_start(containers_to_start, {})
|
||||
|
||||
num_running += len(containers_to_start)
|
||||
|
||||
num_to_create = desired_num - num_running
|
||||
next_number = self._next_container_number()
|
||||
container_numbers = [
|
||||
number for number in range(
|
||||
next_number, next_number + num_to_create
|
||||
)
|
||||
]
|
||||
|
||||
parallel_execute(
|
||||
container_numbers,
|
||||
lambda n: create_and_start(service=self, number=n),
|
||||
lambda n: self.get_container_name(n),
|
||||
"Creating and starting"
|
||||
sorted_containers = sorted(all_containers, key=attrgetter('number'))
|
||||
self._execute_convergence_start(
|
||||
sorted_containers, desired_num, timeout, True, True
|
||||
)
|
||||
|
||||
if desired_num < num_running:
|
||||
|
@ -282,12 +253,7 @@ class Service(object):
|
|||
running_containers,
|
||||
key=attrgetter('number'))
|
||||
|
||||
parallel_execute(
|
||||
sorted_running_containers[-num_to_stop:],
|
||||
stop_and_remove,
|
||||
lambda c: c.name,
|
||||
"Stopping and removing",
|
||||
)
|
||||
self._downscale(sorted_running_containers[-num_to_stop:], timeout)
|
||||
|
||||
def create_container(self,
|
||||
one_off=False,
|
||||
|
@ -400,51 +366,109 @@ class Service(object):
|
|||
|
||||
return has_diverged
|
||||
|
||||
def execute_convergence_plan(self,
|
||||
plan,
|
||||
timeout=None,
|
||||
detached=False,
|
||||
start=True):
|
||||
(action, containers) = plan
|
||||
should_attach_logs = not detached
|
||||
def _execute_convergence_create(self, scale, detached, start):
|
||||
i = self._next_container_number()
|
||||
|
||||
if action == 'create':
|
||||
container = self.create_container()
|
||||
def create_and_start(service, n):
|
||||
container = service.create_container(number=n)
|
||||
if not detached:
|
||||
container.attach_log_stream()
|
||||
if start:
|
||||
self.start_container(container)
|
||||
return container
|
||||
|
||||
if should_attach_logs:
|
||||
container.attach_log_stream()
|
||||
return parallel_execute(
|
||||
range(i, i + scale),
|
||||
lambda n: create_and_start(self, n),
|
||||
lambda n: self.get_container_name(n),
|
||||
"Creating"
|
||||
)[0]
|
||||
|
||||
if start:
|
||||
self.start_container(container)
|
||||
def _execute_convergence_recreate(self, containers, scale, timeout, detached, start):
|
||||
if len(containers) > scale:
|
||||
self._downscale(containers[scale:], timeout)
|
||||
containers = containers[:scale]
|
||||
|
||||
return [container]
|
||||
|
||||
elif action == 'recreate':
|
||||
return [
|
||||
self.recreate_container(
|
||||
container,
|
||||
timeout=timeout,
|
||||
attach_logs=should_attach_logs,
|
||||
def recreate(container):
|
||||
return self.recreate_container(
|
||||
container, timeout=timeout, attach_logs=not detached,
|
||||
start_new_container=start
|
||||
)
|
||||
for container in containers
|
||||
]
|
||||
|
||||
elif action == 'start':
|
||||
if start:
|
||||
for container in containers:
|
||||
self.start_container_if_stopped(container, attach_logs=should_attach_logs)
|
||||
|
||||
containers = parallel_execute(
|
||||
containers,
|
||||
recreate,
|
||||
lambda c: c.name,
|
||||
"Recreating"
|
||||
)[0]
|
||||
if len(containers) < scale:
|
||||
containers.extend(self._execute_convergence_create(
|
||||
scale - len(containers), detached, start
|
||||
))
|
||||
return containers
|
||||
|
||||
elif action == 'noop':
|
||||
def _execute_convergence_start(self, containers, scale, timeout, detached, start):
|
||||
if len(containers) > scale:
|
||||
self._downscale(containers[scale:], timeout)
|
||||
containers = containers[:scale]
|
||||
if start:
|
||||
parallel_execute(
|
||||
containers,
|
||||
lambda c: self.start_container_if_stopped(c, attach_logs=not detached),
|
||||
lambda c: c.name,
|
||||
"Starting"
|
||||
)
|
||||
if len(containers) < scale:
|
||||
containers.extend(self._execute_convergence_create(
|
||||
scale - len(containers), detached, start
|
||||
))
|
||||
return containers
|
||||
|
||||
def _downscale(self, containers, timeout=None):
|
||||
def stop_and_remove(container):
|
||||
container.stop(timeout=self.stop_timeout(timeout))
|
||||
container.remove()
|
||||
|
||||
parallel_execute(
|
||||
containers,
|
||||
stop_and_remove,
|
||||
lambda c: c.name,
|
||||
"Stopping and removing",
|
||||
)
|
||||
|
||||
def execute_convergence_plan(self, plan, timeout=None, detached=False,
|
||||
start=True, scale_override=None):
|
||||
(action, containers) = plan
|
||||
scale = scale_override if scale_override is not None else self.scale_num
|
||||
containers = sorted(containers, key=attrgetter('number'))
|
||||
|
||||
self.show_scale_warnings(scale)
|
||||
|
||||
if action == 'create':
|
||||
return self._execute_convergence_create(
|
||||
scale, detached, start
|
||||
)
|
||||
|
||||
if action == 'recreate':
|
||||
return self._execute_convergence_recreate(
|
||||
containers, scale, timeout, detached, start
|
||||
)
|
||||
|
||||
if action == 'start':
|
||||
return self._execute_convergence_start(
|
||||
containers, scale, timeout, detached, start
|
||||
)
|
||||
|
||||
if action == 'noop':
|
||||
if scale != len(containers):
|
||||
return self._execute_convergence_start(
|
||||
containers, scale, timeout, detached, start
|
||||
)
|
||||
for c in containers:
|
||||
log.info("%s is up-to-date" % c.name)
|
||||
|
||||
return containers
|
||||
|
||||
else:
|
||||
raise Exception("Invalid action: {}".format(action))
|
||||
raise Exception("Invalid action: {}".format(action))
|
||||
|
||||
def recreate_container(
|
||||
self,
|
||||
|
|
|
@ -1866,6 +1866,33 @@ class CLITestCase(DockerClientTestCase):
|
|||
self.assertEqual(len(project.get_service('simple').containers()), 0)
|
||||
self.assertEqual(len(project.get_service('another').containers()), 0)
|
||||
|
||||
def test_up_scale(self):
|
||||
self.base_dir = 'tests/fixtures/scale'
|
||||
project = self.project
|
||||
self.dispatch(['up', '-d'])
|
||||
assert len(project.get_service('web').containers()) == 2
|
||||
assert len(project.get_service('db').containers()) == 1
|
||||
|
||||
self.dispatch(['up', '-d', '--scale', 'web=1'])
|
||||
assert len(project.get_service('web').containers()) == 1
|
||||
assert len(project.get_service('db').containers()) == 1
|
||||
|
||||
self.dispatch(['up', '-d', '--scale', 'web=3'])
|
||||
assert len(project.get_service('web').containers()) == 3
|
||||
assert len(project.get_service('db').containers()) == 1
|
||||
|
||||
self.dispatch(['up', '-d', '--scale', 'web=1', '--scale', 'db=2'])
|
||||
assert len(project.get_service('web').containers()) == 1
|
||||
assert len(project.get_service('db').containers()) == 2
|
||||
|
||||
self.dispatch(['up', '-d'])
|
||||
assert len(project.get_service('web').containers()) == 2
|
||||
assert len(project.get_service('db').containers()) == 1
|
||||
|
||||
self.dispatch(['up', '-d', '--scale', 'web=0', '--scale', 'db=0'])
|
||||
assert len(project.get_service('web').containers()) == 0
|
||||
assert len(project.get_service('db').containers()) == 0
|
||||
|
||||
def test_port(self):
|
||||
self.base_dir = 'tests/fixtures/ports-composefile'
|
||||
self.dispatch(['up', '-d'], None)
|
||||
|
|
|
@ -0,0 +1,9 @@
|
|||
version: '2.2'
|
||||
services:
|
||||
web:
|
||||
image: busybox
|
||||
command: top
|
||||
scale: 2
|
||||
db:
|
||||
image: busybox
|
||||
command: top
|
|
@ -19,6 +19,7 @@ from compose.config.types import VolumeFromSpec
|
|||
from compose.config.types import VolumeSpec
|
||||
from compose.const import COMPOSEFILE_V2_0 as V2_0
|
||||
from compose.const import COMPOSEFILE_V2_1 as V2_1
|
||||
from compose.const import COMPOSEFILE_V2_2 as V2_2
|
||||
from compose.const import COMPOSEFILE_V3_1 as V3_1
|
||||
from compose.const import LABEL_PROJECT
|
||||
from compose.const import LABEL_SERVICE
|
||||
|
@ -1137,6 +1138,33 @@ class ProjectTest(DockerClientTestCase):
|
|||
containers = project.containers()
|
||||
self.assertEqual(len(containers), 1)
|
||||
|
||||
def test_project_up_config_scale(self):
|
||||
config_data = build_config(
|
||||
version=V2_2,
|
||||
services=[{
|
||||
'name': 'web',
|
||||
'image': 'busybox:latest',
|
||||
'command': 'top',
|
||||
'scale': 3
|
||||
}]
|
||||
)
|
||||
|
||||
project = Project.from_config(
|
||||
name='composetest', config_data=config_data, client=self.client
|
||||
)
|
||||
project.up()
|
||||
assert len(project.containers()) == 3
|
||||
|
||||
project.up(scale_override={'web': 2})
|
||||
assert len(project.containers()) == 2
|
||||
|
||||
project.up(scale_override={'web': 4})
|
||||
assert len(project.containers()) == 4
|
||||
|
||||
project.stop()
|
||||
project.up()
|
||||
assert len(project.containers()) == 3
|
||||
|
||||
@v2_only()
|
||||
def test_initialize_volumes(self):
|
||||
vol_name = '{0:x}'.format(random.getrandbits(32))
|
||||
|
|
Loading…
Reference in New Issue