job_server_api.py 4.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155
  1. #!/usr/bin/env python
  2. # Licensed to Cloudera, Inc. under one
  3. # or more contributor license agreements. See the NOTICE file
  4. # distributed with this work for additional information
  5. # regarding copyright ownership. Cloudera, Inc. licenses this file
  6. # to you under the Apache License, Version 2.0 (the
  7. # "License"); you may not use this file except in compliance
  8. # with the License. You may obtain a copy of the License at
  9. #
  10. # http://www.apache.org/licenses/LICENSE-2.0
  11. #
  12. # Unless required by applicable law or agreed to in writing, software
  13. # distributed under the License is distributed on an "AS IS" BASIS,
  14. # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  15. # See the License for the specific language governing permissions and
  16. # limitations under the License.
  17. import logging
  18. import json
  19. import posixpath
  20. import threading
  21. from desktop.lib.rest.http_client import HttpClient
  22. from desktop.lib.rest.resource import Resource
  23. from spark.conf import get_livy_server_url, SECURITY_ENABLED
  24. LOG = logging.getLogger(__name__)
  25. DEFAULT_USER = 'hue'
  26. _API_VERSION = 'v1'
  27. _JSON_CONTENT_TYPE = 'application/json'
  28. _BINARY_CONTENT_TYPE = 'application/octet-stream'
  29. _TEXT_CONTENT_TYPE = 'text/plain'
  30. API_CACHE = None
  31. API_CACHE_LOCK = threading.Lock()
  32. def get_api(user):
  33. global API_CACHE
  34. if API_CACHE is None:
  35. API_CACHE_LOCK.acquire()
  36. try:
  37. if API_CACHE is None:
  38. API_CACHE = JobServerApi(get_livy_server_url())
  39. finally:
  40. API_CACHE_LOCK.release()
  41. API_CACHE.setuser(user)
  42. return API_CACHE
  43. class JobServerApi(object):
  44. def __init__(self, oozie_url):
  45. self._url = posixpath.join(oozie_url)
  46. self._client = HttpClient(self._url, logger=LOG)
  47. self._root = Resource(self._client)
  48. self._security_enabled = SECURITY_ENABLED.get()
  49. self._thread_local = threading.local()
  50. if self.security_enabled:
  51. self._client.set_kerberos_auth()
  52. def __str__(self):
  53. return "JobServerApi at %s" % (self._url,)
  54. @property
  55. def url(self):
  56. return self._url
  57. @property
  58. def security_enabled(self):
  59. return self._security_enabled
  60. @property
  61. def user(self):
  62. return self._thread_local.user
  63. def setuser(self, user):
  64. if hasattr(user, 'username'):
  65. self._thread_local.user = user.username
  66. else:
  67. self._thread_local.user = user
  68. def get_status(self):
  69. return self._root.get('sessions')
  70. def get_log(self, uuid, startFrom=None, size=None):
  71. params = {}
  72. if startFrom is not None:
  73. params['from'] = startFrom
  74. if size is not None:
  75. params['size'] = size
  76. response = self._root.get('sessions/%s/log' % uuid, params=params)
  77. return '\n'.join(response['log'])
  78. def create_session(self, **properties):
  79. properties['proxyUser'] = self.user
  80. return self._root.post('sessions', data=json.dumps(properties), contenttype=_JSON_CONTENT_TYPE)
  81. def get_session(self, uuid):
  82. return self._root.get('sessions/%s' % uuid)
  83. def submit_statement(self, uuid, statement):
  84. data = {'code': statement}
  85. return self._root.post('sessions/%s/statements' % uuid, data=json.dumps(data), contenttype=_JSON_CONTENT_TYPE)
  86. def inspect(self, uuid, statement):
  87. data = {'code': statement}
  88. return self._root.post('sessions/%s/inspect' % uuid, data=json.dumps(data), contenttype=_JSON_CONTENT_TYPE)
  89. def fetch_data(self, session, statement):
  90. return self._root.get('sessions/%s/statements/%s' % (session, statement))
  91. def cancel(self, session):
  92. return self._root.post('sessions/%s/interrupt' % session)
  93. def close(self, uuid):
  94. return self._root.delete('sessions/%s' % uuid)
  95. def get_batches(self):
  96. return self._root.get('batches')
  97. def submit_batch(self, properties):
  98. properties['proxyUser'] = self.user
  99. return self._root.post('batches', data=json.dumps(properties), contenttype=_JSON_CONTENT_TYPE)
  100. def get_batch(self, uuid):
  101. return self._root.get('batches/%s' % uuid)
  102. def get_batch_status(self, uuid):
  103. response = self._root.get('batches/%s/state' % uuid)
  104. return response['state']
  105. def get_batch_log(self, uuid, startFrom=None, size=None):
  106. params = {}
  107. if startFrom is not None:
  108. params['from'] = startFrom
  109. if size is not None:
  110. params['size'] = size
  111. response = self._root.get('batches/%s/log' % uuid, params=params)
  112. return '\n'.join(response['log'])
  113. def close_batch(self, uuid):
  114. return self._root.delete('batches/%s' % uuid)