manager_client.py 11 KB


  1. #!/usr/bin/env python
  2. # -- coding: utf-8 --
  3. # Licensed to Cloudera, Inc. under one
  4. # or more contributor license agreements. See the NOTICE file
  5. # distributed with this work for additional information
  6. # regarding copyright ownership. Cloudera, Inc. licenses this file
  7. # to you under the Apache License, Version 2.0 (the
  8. # "License"); you may not use this file except in compliance
  9. # with the License. You may obtain a copy of the License at
  10. #
  11. # http://www.apache.org/licenses/LICENSE-2.0
  12. #
  13. # Unless required by applicable law or agreed to in writing, software
  14. # distributed under the License is distributed on an "AS IS" BASIS,
  15. # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  16. # See the License for the specific language governing permissions and
  17. # limitations under the License.
  18. import base64
  19. import json
  20. import logging
  21. import urllib
  22. import urllib2
  23. from django.core.cache import cache
  24. from django.utils.translation import ugettext as _
  25. from desktop.lib.rest.http_client import RestException, HttpClient
  26. from desktop.lib.rest.resource import Resource
  27. from desktop.lib.i18n import smart_unicode
  28. from metadata.conf import MANAGER, get_navigator_auth_username, get_navigator_auth_password
  29. LOG = logging.getLogger(__name__)
  30. VERSION = 'v19'
  31. class ManagerApiException(Exception):
  32. def __init__(self, message=None):
  33. self.message = message or _('No error message, please check the logs.')
  34. def __str__(self):
  35. return str(self.message)
  36. def __unicode__(self):
  37. return smart_unicode(self.message)
  38. class ManagerApi(object):
  39. """
  40. https://cloudera.github.io/cm_api/
  41. """
  42. def __init__(self, user=None, security_enabled=False, ssl_cert_ca_verify=False):
  43. self._api_url = '%s/%s' % (MANAGER.API_URL.get().strip('/'), VERSION)
  44. self._username = get_navigator_auth_username()
  45. self._password = get_navigator_auth_password()
  46. self.user = user
  47. self._client = HttpClient(self._api_url, logger=LOG)
  48. if security_enabled:
  49. self._client.set_kerberos_auth()
  50. else:
  51. self._client.set_basic_auth(self._username, self._password)
  52. self._client.set_verify(ssl_cert_ca_verify)
  53. self._root = Resource(self._client)
  54. def has_service(self, service_name, cluster_name=None):
  55. cluster = self._get_cluster(cluster_name)
  56. try:
  57. services = self._root.get('clusters/%(cluster_name)s/serviceTypes' % {
  58. 'cluster_name': cluster['name'],
  59. 'service_name': service_name
  60. })['items']
  61. return service_name in services
  62. except RestException, e:
  63. raise ManagerApiException(e)
  64. def get_spark_history_server_url(self, cluster_name=None):
  65. service_name = "SPARK_ON_YARN"
  66. shs_role_type = "SPARK_YARN_HISTORY_SERVER"
  67. try:
  68. cluster = self._get_cluster(cluster_name)
  69. services = self._root.get('clusters/%(cluster_name)s/services' % {
  70. 'cluster_name': cluster['name'],
  71. 'service_name': service_name
  72. })['items']
  73. service_display_names = [service['displayName'] for service in services if service['type'] == service_name]
  74. if service_display_names:
  75. spark_service_display_name = service_display_names[0]
  76. servers = self._root.get('clusters/%(cluster_name)s/services/%(spark_service_display_name)s/roles' % {
  77. 'cluster_name': cluster['name'],
  78. 'spark_service_display_name': spark_service_display_name
  79. })['items']
  80. shs_server_names = [server['name'] for server in servers if server['type'] == shs_role_type]
  81. shs_server_name = shs_server_names[0] if shs_server_names else None
  82. shs_server_hostRef = [server['hostRef'] for server in servers if server['type'] == shs_role_type]
  83. shs_server_hostId = shs_server_hostRef[0]['hostId'] if shs_server_hostRef else None
  84. if shs_server_name and shs_server_hostId:
  85. shs_server_configs = self._root.get('clusters/%(cluster_name)s/services/%(spark_service_display_name)s/roles/%(shs_server_name)s/config' % {
  86. 'cluster_name': cluster['name'],
  87. 'spark_service_display_name': spark_service_display_name,
  88. 'shs_server_name': shs_server_name
  89. }, params={'view': 'full'})['items']
  90. shs_ui_port = None
  91. shs_ssl_port = None
  92. shs_ssl_enabled = None
  93. for config in shs_server_configs:
  94. if 'relatedName' in config and 'default' in config:
  95. if config['relatedName'] == 'spark.history.ui.port':
  96. shs_ui_port = config['default']
  97. if config['relatedName'] == 'spark.ssl.historyServer.port':
  98. shs_ssl_port = config['default']
  99. if config['relatedName'] == 'spark.ssl.historyServer.enabled':
  100. shs_ssl_enabled = config['default']
  101. shs_ui_host = self._root.get('hosts/%(hostId)s' % {'hostId': shs_server_hostId})
  102. shs_ui_hostname = shs_ui_host['hostname'] if shs_ui_host else None
  103. return self.assemble_shs_url(shs_ui_hostname, shs_ui_port, shs_ssl_port, shs_ssl_enabled)
  104. except Exception, e:
  105. LOG.warn("Check Spark history server via ManangerAPI: %s" % e)
  106. return None
  107. def assemble_shs_url(self, shs_ui_hostname, shs_ui_port=None, shs_ssl_port=None, shs_ssl_enabled=None):
  108. if not shs_ui_hostname or not shs_ui_port or not shs_ssl_port or not shs_ssl_enabled:
  109. LOG.warn("Spark conf not found!")
  110. return None
  111. protocol = 'https' if shs_ssl_enabled.lower() == 'true' else 'http'
  112. shs_url = '%(protocol)s://%(hostname)s:%(port)s' % {
  113. 'protocol': protocol,
  114. 'hostname': shs_ui_hostname,
  115. 'port': shs_ssl_port if shs_ssl_enabled.lower() == 'true' else shs_ui_port,
  116. }
  117. return shs_url
  118. def tools_echo(self):
  119. try:
  120. params = (
  121. ('message', 'hello'),
  122. )
  123. LOG.info(params)
  124. return self._root.get('tools/echo', params=params)
  125. except RestException, e:
  126. raise ManagerApiException(e)
  127. def get_kafka_brokers(self, cluster_name=None):
  128. try:
  129. hosts = self._get_hosts('KAFKA', 'KAFKA_BROKER', cluster_name=cluster_name)
  130. brokers_hosts = [host['hostname'] + ':9092' for host in hosts]
  131. return ','.join(brokers_hosts)
  132. except RestException, e:
  133. raise ManagerApiException(e)
  134. def get_kudu_master(self, cluster_name=None):
  135. try:
  136. cluster = self._get_cluster(cluster_name)
  137. services = self._root.get('clusters/%(name)s/services' % cluster)['items']
  138. service = [service for service in services if service['type'] == 'KUDU'][0]
  139. master = self._get_roles(cluster['name'], service['name'], 'KUDU_MASTER')[0]
  140. master_host = self._root.get('hosts/%(hostId)s' % master['hostRef'])
  141. return master_host['hostname']
  142. except RestException, e:
  143. raise ManagerApiException(e)
  144. def get_kafka_topics(self, broker_host):
  145. try:
  146. client = HttpClient('http://%s:24042' % broker_host, logger=LOG)
  147. root = Resource(client)
  148. return root.get('/api/topics')
  149. except RestException, e:
  150. raise ManagerApiException(e)
  151. def update_flume_config(self, cluster_name, config_name, config_value):
  152. service = 'FLUME-1'
  153. cluster = self._get_cluster(cluster_name)
  154. roleConfigGroup = [role['roleConfigGroupRef']['roleConfigGroupName'] for role in self._get_roles(cluster['name'], service, 'AGENT')]
  155. data = {
  156. u'items': [{
  157. u'url': u'/api/v8/clusters/%(cluster_name)s/services/%(service)s/roleConfigGroups/%(roleConfigGroups)s/config?message=Updated%20service%20and%20role%20type%20configurations.'.replace('%(cluster_name)s', urllib.quote(cluster['name'])).replace('%(service)s', service).replace('%(roleConfigGroups)s', roleConfigGroup[0]),
  158. u'body': {
  159. u'items': [
  160. {u'name': config_name, u'value': config_value}
  161. ]
  162. },
  163. u'contentType': u'application/json',
  164. u'method': u'PUT'
  165. }]
  166. }
  167. return self.batch(
  168. items=data
  169. )
  170. def get_flume_agents(self, cluster_name=None):
  171. return [host['hostname'] for host in self._get_hosts('FLUME', 'AGENT', cluster_name=cluster_name)]
  172. def _get_hosts(self, service_name, role_name, cluster_name=None):
  173. try:
  174. cluster = self._get_cluster(cluster_name)
  175. services = self._root.get('clusters/%(name)s/services' % cluster)['items']
  176. service = [service for service in services if service['type'] == service_name][0]
  177. hosts = self._get_roles(cluster['name'], service['name'], role_name)
  178. hosts_ids = [host['hostRef']['hostId'] for host in hosts]
  179. hosts = self._root.get('hosts')['items']
  180. return [host for host in hosts if host['hostId'] in hosts_ids]
  181. except RestException, e:
  182. raise ManagerApiException(e)
  183. def refresh_flume(self, cluster_name, restart=False):
  184. service = 'FLUME-1'
  185. cluster = self._get_cluster(cluster_name)
  186. roles = [role['name'] for role in self._get_roles(cluster['name'], service, 'AGENT')]
  187. if restart:
  188. return self.restart_services(cluster['name'], service, roles)
  189. else:
  190. return self.refresh_configs(cluster['name'], service, roles)
  191. def refresh_configs(self, cluster_name, service=None, roles=None):
  192. try:
  193. if service is None:
  194. return self._root.post('clusters/%(cluster_name)s/commands/refresh' % {'cluster_name': cluster_name}, contenttype="application/json")
  195. elif roles is None:
  196. return self._root.post('clusters/%(cluster_name)s/services/%(service)s/roleCommands/refresh' % {'cluster_name': cluster_name, 'service': service}, contenttype="application/json")
  197. else:
  198. return self._root.post(
  199. 'clusters/%(cluster_name)s/services/%(service)s/roleCommands/refresh' % {'cluster_name': cluster_name, 'service': service},
  200. data=json.dumps({"items": roles}),
  201. contenttype="application/json"
  202. )
  203. except RestException, e:
  204. raise ManagerApiException(e)
  205. def restart_services(self, cluster_name, service=None, roles=None):
  206. try:
  207. if service is None:
  208. return self._root.post('clusters/%(cluster_name)s/commands/restart' % {'cluster_name': cluster_name}, contenttype="application/json")
  209. elif roles is None:
  210. return self._root.post('clusters/%(cluster_name)s/services/%(service)s/roleCommands/restart' % {'cluster_name': cluster_name, 'service': service}, contenttype="application/json")
  211. else:
  212. return self._root.post(
  213. 'clusters/%(cluster_name)s/services/%(service)s/roleCommands/restart' % {'cluster_name': cluster_name, 'service': service},
  214. data=json.dumps({"items": roles}),
  215. contenttype="application/json"
  216. )
  217. except RestException, e:
  218. raise ManagerApiException(e)
  219. def batch(self, items):
  220. try:
  221. return self._root.post('batch', data=json.dumps(items), contenttype='application/json')
  222. except RestException, e:
  223. raise ManagerApiException(e)
  224. def _get_cluster(self, cluster_name=None):
  225. clusters = self._root.get('clusters/')['items']
  226. if cluster_name is not None:
  227. cluster = [cluster for cluster in clusters if cluster['name'] == cluster_name][0]
  228. else:
  229. cluster = clusters[0]
  230. return cluster
  231. def _get_roles(self, cluster_name, service_name, role_type):
  232. roles = self._root.get('clusters/%(cluster_name)s/services/%(service_name)s/roles' % {'cluster_name': cluster_name, 'service_name': service_name})['items']
  233. return [role for role in roles if role['type'] == role_type]