api.py 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241
  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 json
  18. import logging
  19. import time
  20. from django.core.urlresolvers import reverse
  21. from django.utils.translation import ugettext as _
  22. from desktop.lib.i18n import smart_str
  23. from desktop.lib.view_util import format_duration_in_millis
  24. from liboozie.oozie_api import get_oozie
  25. from oozie.models import Workflow, Pig
  26. from oozie.views.api import get_log as get_workflow_logs
  27. from oozie.views.editor import _submit_workflow
  28. LOG = logging.getLogger(__name__)
  29. def get(fs, jt, user):
  30. return OozieApi(fs, jt, user)
  31. class OozieApi(object):
  32. """
  33. Oozie submission.
  34. """
  35. WORKFLOW_NAME = 'pig-app-hue-script'
  36. LOG_START_PATTERN = '(Pig script \[(?:[\w.-]+)\] content:.+)'
  37. LOG_END_PATTERN = '(<<< Invocation of Pig command completed <<<|<<< Invocation of Main class completed <<<)'
  38. MAX_DASHBOARD_JOBS = 100
  39. def __init__(self, fs, jt, user):
  40. self.oozie_api = get_oozie(user)
  41. self.fs = fs
  42. self.jt = jt
  43. self.user = user
  44. def submit(self, pig_script, params):
  45. workflow = None
  46. try:
  47. workflow = self._create_workflow(pig_script, params)
  48. mapping = dict([(param['name'], param['value']) for param in workflow.get_parameters()])
  49. oozie_wf = _submit_workflow(self.user, self.fs, self.jt, workflow, mapping)
  50. finally:
  51. if workflow:
  52. workflow.delete(skip_trash=True)
  53. return oozie_wf
  54. def _create_workflow(self, pig_script, params):
  55. workflow = Workflow.objects.new_workflow(self.user)
  56. workflow.schema_version = 'uri:oozie:workflow:0.5'
  57. workflow.name = OozieApi.WORKFLOW_NAME
  58. workflow.is_history = True
  59. if pig_script.use_hcatalog:
  60. workflow.add_parameter("oozie.action.sharelib.for.pig", "pig,hcatalog")
  61. workflow.save()
  62. Workflow.objects.initialize(workflow, self.fs)
  63. script_path = workflow.deployment_dir + '/script.pig'
  64. if self.fs: # For testing, difficult to mock
  65. self.fs.do_as_user(self.user.username, self.fs.create, script_path, data=smart_str(pig_script.dict['script']))
  66. files = []
  67. archives = []
  68. popup_params = json.loads(params)
  69. popup_params_names = [param['name'] for param in popup_params]
  70. pig_params = self._build_parameters(popup_params)
  71. if pig_script.isV2:
  72. pig_params += [{"type": "argument", "value": param} for param in pig_script.dict['parameters']]
  73. job_properties = [{"name": prop.split('=', 1)[0], "value": prop.split('=', 1)[1]} for prop in pig_script.dict['hadoopProperties']]
  74. for resource in pig_script.dict['resources']:
  75. if resource.endswith('.zip') or resource.endswith('.tgz') or resource.endswith('.tar') or resource.endswith('.gz'):
  76. archives.append({"dummy": "", "name": resource})
  77. else:
  78. files.append(resource)
  79. else:
  80. script_params = [param for param in pig_script.dict['parameters'] if param['name'] not in popup_params_names]
  81. pig_params += self._build_parameters(script_params)
  82. job_properties = [{"name": prop['name'], "value": prop['value']} for prop in pig_script.dict['hadoopProperties']]
  83. for resource in pig_script.dict['resources']:
  84. if resource['type'] == 'file':
  85. files.append(resource['value'])
  86. if resource['type'] == 'archive':
  87. archives.append({"dummy": "", "name": resource['value']})
  88. action = Pig.objects.create(
  89. name='pig-5760',
  90. script_path=script_path,
  91. workflow=workflow,
  92. node_type='pig',
  93. params=json.dumps(pig_params),
  94. files=json.dumps(files),
  95. archives=json.dumps(archives),
  96. job_properties=json.dumps(job_properties)
  97. )
  98. credentials = []
  99. if pig_script.use_hcatalog and self.oozie_api.security_enabled:
  100. credentials.append({'name': 'hcat', 'value': True})
  101. if pig_script.use_hbase and self.oozie_api.security_enabled:
  102. credentials.append({'name': 'hbase', 'value': True})
  103. if credentials:
  104. action.credentials = credentials # Note, action.credentials is a @setter here
  105. action.save()
  106. action.add_node(workflow.end)
  107. start_link = workflow.start.get_link()
  108. start_link.child = action
  109. start_link.save()
  110. return workflow
  111. def _build_parameters(self, params):
  112. pig_params = []
  113. for param in params:
  114. if param['name'].startswith('-'):
  115. pig_params.append({"type": "argument", "value": "%(name)s" % param})
  116. if param['value']:
  117. pig_params.append({"type": "argument", "value": "%(value)s" % param})
  118. else:
  119. # Simpler way and backward compatibility for parameters
  120. pig_params.append({"type": "argument", "value": "-param"})
  121. pig_params.append({"type": "argument", "value": "%(name)s=%(value)s" % param})
  122. return pig_params
  123. def stop(self, job_id):
  124. return self.oozie_api.job_control(job_id, 'kill')
  125. def get_jobs(self):
  126. kwargs = {'cnt': OozieApi.MAX_DASHBOARD_JOBS,}
  127. kwargs['filters'] = [
  128. ('user', self.user.username),
  129. ('name', OozieApi.WORKFLOW_NAME)
  130. ]
  131. return self.oozie_api.get_workflows(**kwargs).jobs
  132. def get_log(self, request, oozie_workflow, make_links=True):
  133. return get_workflow_logs(request, oozie_workflow, make_links=make_links, log_start_pattern=self.LOG_START_PATTERN,
  134. log_end_pattern=self.LOG_END_PATTERN)
  135. def massaged_jobs_for_json(self, request, oozie_jobs, hue_jobs):
  136. jobs = []
  137. hue_jobs = dict([(script.dict.get('job_id'), script) for script in hue_jobs if script.dict.get('job_id')])
  138. for job in oozie_jobs:
  139. if job.is_running():
  140. job = self.oozie_api.get_job(job.id)
  141. get_copy = request.GET.copy() # Hacky, would need to refactor JobBrowser get logs
  142. get_copy['format'] = 'python'
  143. request.GET = get_copy
  144. try:
  145. logs, workflow_action, is_really_done = self.get_log(request, job)
  146. progress = workflow_action[0]['progress']
  147. except:
  148. LOG.exception('failed to get progress')
  149. progress = 0
  150. else:
  151. progress = 100
  152. hue_pig = hue_jobs.get(job.id) and hue_jobs.get(job.id) or None
  153. massaged_job = {
  154. 'id': job.id,
  155. 'lastModTime': hasattr(job, 'lastModTime') and job.lastModTime and format_time(job.lastModTime) or None,
  156. 'kickoffTime': hasattr(job, 'kickoffTime') and job.kickoffTime or None,
  157. 'timeOut': hasattr(job, 'timeOut') and job.timeOut or None,
  158. 'endTime': job.endTime and format_time(job.endTime) or None,
  159. 'status': job.status,
  160. 'isRunning': job.is_running(),
  161. 'duration': job.endTime and job.startTime and format_duration_in_millis(( time.mktime(job.endTime) - time.mktime(job.startTime) ) * 1000) or None,
  162. 'appName': hue_pig and hue_pig.dict['name'] or _('Unsaved script'),
  163. 'scriptId': hue_pig and hue_pig.id or -1,
  164. 'scriptContent': hue_pig and hue_pig.dict['script'] or '',
  165. 'progress': progress,
  166. 'progressPercent': '%d%%' % progress,
  167. 'user': job.user,
  168. 'absoluteUrl': job.get_absolute_url(),
  169. 'canEdit': has_job_edition_permission(job, self.user),
  170. 'killUrl': reverse('oozie:manage_oozie_jobs', kwargs={'job_id':job.id, 'action':'kill'}),
  171. 'watchUrl': reverse('pig:watch', kwargs={'job_id': job.id}) + '?format=python',
  172. 'created': hasattr(job, 'createdTime') and job.createdTime and job.createdTime and ((job.type == 'Bundle' and job.createdTime) or format_time(job.createdTime)),
  173. 'startTime': hasattr(job, 'startTime') and format_time(job.startTime) or None,
  174. 'run': hasattr(job, 'run') and job.run or 0,
  175. 'frequency': hasattr(job, 'frequency') and job.frequency or None,
  176. 'timeUnit': hasattr(job, 'timeUnit') and job.timeUnit or None,
  177. }
  178. jobs.append(massaged_job)
  179. return jobs
  180. def format_time(st_time):
  181. if st_time is None:
  182. return '-'
  183. else:
  184. return time.strftime("%a, %d %b %Y %H:%M:%S", st_time)
  185. def has_job_edition_permission(oozie_job, user):
  186. return user.is_superuser or oozie_job.user == user.username