Chromium Code Reviews
chromiumcodereview-hr@appspot.gserviceaccount.com (chromiumcodereview-hr) | Please choose your nickname with Settings | Help | Chromium Project | Gerrit Changes | Sign out
(121)

Unified Diff: third_party/gsutil/boto/boto/emr/connection.py

Issue 10199002: Upgrade gsutil to 3.4 (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Addressed comments Created 8 years, 8 months ago
Use n/p to move between diff chunks; N/P to move between comments. Draft comments are only viewable by you.
Jump to:
View side-by-side diff with in-line comments
Download patch
« no previous file with comments | « third_party/gsutil/boto/boto/emr/bootstrap_action.py ('k') | third_party/gsutil/boto/boto/emr/emrobject.py » ('j') | no next file with comments »
Expand Comments ('e') | Collapse Comments ('c') | Show Comments Hide Comments ('s')
Index: third_party/gsutil/boto/boto/emr/connection.py
diff --git a/third_party/gsutil/20110627/boto/boto/emr/connection.py b/third_party/gsutil/boto/boto/emr/connection.py
similarity index 53%
rename from third_party/gsutil/20110627/boto/boto/emr/connection.py
rename to third_party/gsutil/boto/boto/emr/connection.py
index 70ddf8247238461d152d3f8512fe698b22d80d8a..7e39af9cb366a3b042625f27221d9616135b9694 100644
--- a/third_party/gsutil/20110627/boto/boto/emr/connection.py
+++ b/third_party/gsutil/boto/boto/emr/connection.py
@@ -29,6 +29,7 @@ import boto
import boto.utils
from boto.ec2.regioninfo import RegionInfo
from boto.emr.emrobject import JobFlow, RunJobFlowResponse
+from boto.emr.emrobject import AddInstanceGroupsResponse, ModifyInstanceGroupsResponse
from boto.emr.step import JarStep
from boto.connection import AWSQueryConnection
from boto.exception import EmrResponseError
@@ -50,7 +51,8 @@ class EmrConnection(AWSQueryConnection):
proxy_user=None, proxy_pass=None, debug=0,
https_connection_factory=None, region=None, path='/'):
if not region:
- region = RegionInfo(self, self.DefaultRegionName, self.DefaultRegionEndpoint)
+ region = RegionInfo(self, self.DefaultRegionName,
+ self.DefaultRegionEndpoint)
self.region = region
AWSQueryConnection.__init__(self, aws_access_key_id,
aws_secret_access_key,
@@ -145,39 +147,134 @@ class EmrConnection(AWSQueryConnection):
return self.get_object(
'AddJobFlowSteps', params, RunJobFlowResponse, verb='POST')
- def run_jobflow(self, name, log_uri, ec2_keyname=None, availability_zone=None,
+ def add_instance_groups(self, jobflow_id, instance_groups):
+ """
+ Adds instance groups to a running cluster.
+
+ :type jobflow_id: str
+ :param jobflow_id: The id of the jobflow which will take the
+ new instance groups
+
+ :type instance_groups: list(boto.emr.InstanceGroup)
+ :param instance_groups: A list of instance groups to add to the job
+ """
+ if type(instance_groups) != types.ListType:
+ instance_groups = [instance_groups]
+ params = {}
+ params['JobFlowId'] = jobflow_id
+ params.update(self._build_instance_group_list_args(instance_groups))
+
+ return self.get_object('AddInstanceGroups', params,
+ AddInstanceGroupsResponse, verb='POST')
+
+ def modify_instance_groups(self, instance_group_ids, new_sizes):
+ """
+ Modify the number of nodes and configuration settings in an
+ instance group.
+
+ :type instance_group_ids: list(str)
+ :param instance_group_ids: A list of the ID's of the instance
+ groups to be modified
+
+ :type new_sizes: list(int)
+ :param new_sizes: A list of the new sizes for each instance group
+ """
+ if type(instance_group_ids) != types.ListType:
+ instance_group_ids = [instance_group_ids]
+ if type(new_sizes) != types.ListType:
+ new_sizes = [new_sizes]
+
+ instance_groups = zip(instance_group_ids, new_sizes)
+
+ params = {}
+ for k, ig in enumerate(instance_groups):
+ # could be wrong - the example amazon gives uses
+ # InstanceRequestCount, while the api documentation
+ # says InstanceCount
+ params['InstanceGroups.member.%d.InstanceGroupId' % (k+1) ] = ig[0]
+ params['InstanceGroups.member.%d.InstanceCount' % (k+1) ] = ig[1]
+
+ return self.get_object('ModifyInstanceGroups', params,
+ ModifyInstanceGroupsResponse, verb='POST')
+
+ def run_jobflow(self, name, log_uri=None, ec2_keyname=None,
+ availability_zone=None,
master_instance_type='m1.small',
slave_instance_type='m1.small', num_instances=1,
action_on_failure='TERMINATE_JOB_FLOW', keep_alive=False,
enable_debugging=False,
- hadoop_version='0.18',
+ hadoop_version=None,
steps=[],
- bootstrap_actions=[]):
+ bootstrap_actions=[],
+ instance_groups=None,
+ additional_info=None,
+ ami_version=None,
+ api_params=None):
"""
Runs a job flow
-
:type name: str
:param name: Name of the job flow
+
:type log_uri: str
:param log_uri: URI of the S3 bucket to place logs
+
:type ec2_keyname: str
:param ec2_keyname: EC2 key used for the instances
+
:type availability_zone: str
:param availability_zone: EC2 availability zone of the cluster
+
:type master_instance_type: str
:param master_instance_type: EC2 instance type of the master
+
:type slave_instance_type: str
:param slave_instance_type: EC2 instance type of the slave nodes
+
:type num_instances: int
:param num_instances: Number of instances in the Hadoop cluster
+
:type action_on_failure: str
:param action_on_failure: Action to take if a step terminates
+
:type keep_alive: bool
- :param keep_alive: Denotes whether the cluster should stay alive upon completion
+ :param keep_alive: Denotes whether the cluster should stay
+ alive upon completion
+
:type enable_debugging: bool
- :param enable_debugging: Denotes whether AWS console debugging should be enabled.
+ :param enable_debugging: Denotes whether AWS console debugging
+ should be enabled.
+
+ :type hadoop_version: str
+ :param hadoop_version: Version of Hadoop to use. This no longer
+ defaults to '0.20' and now uses the AMI default.
+
:type steps: list(boto.emr.Step)
:param steps: List of steps to add with the job
+
+ :type bootstrap_actions: list(boto.emr.BootstrapAction)
+ :param bootstrap_actions: List of bootstrap actions that run
+ before Hadoop starts.
+
+ :type instance_groups: list(boto.emr.InstanceGroup)
+ :param instance_groups: Optional list of instance groups to
+ use when creating this job.
+ NB: When provided, this argument supersedes num_instances
+ and master/slave_instance_type.
+
+ :type ami_version: str
+ :param ami_version: Amazon Machine Image (AMI) version to use
+ for instances. Values accepted by EMR are '1.0', '2.0', and
+ 'latest'; EMR currently defaults to '1.0' if you don't set
+ 'ami_version'.
+
+ :type additional_info: JSON str
+ :param additional_info: A JSON string for selecting additional features
+
+ :type api_params: dict
+ :param api_params: a dictionary of additional parameters to pass
+ directly to the EMR API (so you don't have to upgrade boto to
+ use new EMR features). You can also delete an API parameter
+ by setting it to None.
:rtype: str
:return: The jobflow id
@@ -185,14 +282,36 @@ class EmrConnection(AWSQueryConnection):
params = {}
if action_on_failure:
params['ActionOnFailure'] = action_on_failure
+ if log_uri:
+ params['LogUri'] = log_uri
params['Name'] = name
- params['LogUri'] = log_uri
- # Instance args
- instance_params = self._build_instance_args(ec2_keyname, availability_zone,
- master_instance_type, slave_instance_type,
- num_instances, keep_alive, hadoop_version)
- params.update(instance_params)
+ # Common instance args
+ common_params = self._build_instance_common_args(ec2_keyname,
+ availability_zone,
+ keep_alive,
+ hadoop_version)
+ params.update(common_params)
+
+ # NB: according to the AWS API's error message, we must
+ # "configure instances either using instance count, master and
+ # slave instance type or instance groups but not both."
+ #
+ # Thus we switch here on the truthiness of instance_groups.
+ if not instance_groups:
+ # Instance args (the common case)
+ instance_params = self._build_instance_count_and_type_args(
+ master_instance_type,
+ slave_instance_type,
+ num_instances)
+ params.update(instance_params)
+ else:
+ # Instance group args (for spot instances or a heterogenous cluster)
+ list_args = self._build_instance_group_list_args(instance_groups)
+ instance_params = dict(
+ ('Instances.%s' % k, v) for k, v in list_args.iteritems()
+ )
+ params.update(instance_params)
# Debugging step from EMR API docs
if enable_debugging:
@@ -212,10 +331,43 @@ class EmrConnection(AWSQueryConnection):
bootstrap_action_args = [self._build_bootstrap_action_args(bootstrap_action) for bootstrap_action in bootstrap_actions]
params.update(self._build_bootstrap_action_list(bootstrap_action_args))
+ if ami_version:
+ params['AmiVersion'] = ami_version
+
+ if additional_info is not None:
+ params['AdditionalInfo'] = additional_info
+
+ if api_params:
+ for key, value in api_params.iteritems():
+ if value is None:
+ params.pop(key, None)
+ else:
+ params[key] = value
+
response = self.get_object(
'RunJobFlow', params, RunJobFlowResponse, verb='POST')
return response.jobflowid
+ def set_termination_protection(self, jobflow_id,
+ termination_protection_status):
+ """
+ Set termination protection on specified Elastic MapReduce job flows
+
+ :type jobflow_ids: list or str
+ :param jobflow_ids: A list of job flow IDs
+
+ :type termination_protection_status: bool
+ :param termination_protection_status: Termination protection status
+ """
+ assert termination_protection_status in (True, False)
+
+ params = {}
+ params['TerminationProtected'] = (termination_protection_status and "true") or "false"
+ self.build_list_params(params, [jobflow_id], 'JobFlowIds.member')
+
+ return self.get_status('SetTerminationProtection', params, verb='POST')
+
+
def _build_bootstrap_action_args(self, bootstrap_action):
bootstrap_action_params = {}
bootstrap_action_params['ScriptBootstrapAction.Path'] = bootstrap_action.path
@@ -267,20 +419,69 @@ class EmrConnection(AWSQueryConnection):
params['Steps.member.%s.%s' % (i+1, key)] = value
return params
- def _build_instance_args(self, ec2_keyname, availability_zone, master_instance_type,
- slave_instance_type, num_instances, keep_alive, hadoop_version):
+ def _build_instance_common_args(self, ec2_keyname, availability_zone,
+ keep_alive, hadoop_version):
+ """
+ Takes a number of parameters used when starting a jobflow (as
+ specified in run_jobflow() above). Returns a comparable dict for
+ use in making a RunJobFlow request.
+ """
params = {
- 'Instances.MasterInstanceType' : master_instance_type,
- 'Instances.SlaveInstanceType' : slave_instance_type,
- 'Instances.InstanceCount' : num_instances,
'Instances.KeepJobFlowAliveWhenNoSteps' : str(keep_alive).lower(),
- 'Instances.HadoopVersion' : hadoop_version
}
+ if hadoop_version:
+ params['Instances.HadoopVersion'] = hadoop_version
if ec2_keyname:
params['Instances.Ec2KeyName'] = ec2_keyname
if availability_zone:
- params['Placement.AvailabilityZone'] = availability_zone
+ params['Instances.Placement.AvailabilityZone'] = availability_zone
+
+ return params
+
+ def _build_instance_count_and_type_args(self, master_instance_type,
+ slave_instance_type, num_instances):
+ """
+ Takes a master instance type (string), a slave instance type
+ (string), and a number of instances. Returns a comparable dict
+ for use in making a RunJobFlow request.
+ """
+ params = {
+ 'Instances.MasterInstanceType' : master_instance_type,
+ 'Instances.SlaveInstanceType' : slave_instance_type,
+ 'Instances.InstanceCount' : num_instances,
+ }
+ return params
+ def _build_instance_group_args(self, instance_group):
+ """
+ Takes an InstanceGroup; returns a dict that, when its keys are
+ properly prefixed, can be used for describing InstanceGroups in
+ RunJobFlow or AddInstanceGroups requests.
+ """
+ params = {
+ 'InstanceCount' : instance_group.num_instances,
+ 'InstanceRole' : instance_group.role,
+ 'InstanceType' : instance_group.type,
+ 'Name' : instance_group.name,
+ 'Market' : instance_group.market
+ }
+ if instance_group.market == 'SPOT':
+ params['BidPrice'] = instance_group.bidprice
return params
+ def _build_instance_group_list_args(self, instance_groups):
+ """
+ Takes a list of InstanceGroups, or a single InstanceGroup. Returns
+ a comparable dict for use in making a RunJobFlow or AddInstanceGroups
+ request.
+ """
+ if type(instance_groups) != types.ListType:
+ instance_groups = [instance_groups]
+
+ params = {}
+ for i, instance_group in enumerate(instance_groups):
+ ig_dict = self._build_instance_group_args(instance_group)
+ for key, value in ig_dict.iteritems():
+ params['InstanceGroups.member.%d.%s' % (i+1, key)] = value
+ return params
« no previous file with comments | « third_party/gsutil/boto/boto/emr/bootstrap_action.py ('k') | third_party/gsutil/boto/boto/emr/emrobject.py » ('j') | no next file with comments »

Powered by Google App Engine
This is Rietveld 408576698