models.py 9.6 KB


  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. from __future__ import division
  18. from builtins import str
  19. from builtins import object
  20. import datetime
  21. import logging
  22. import math
  23. import functools
  24. import re
  25. from django.db import connection, models
  26. from django.urls import reverse
  27. from django.utils.html import escape
  28. from django.utils.translation import ugettext as _
  29. from desktop.auth.backend import is_admin
  30. from desktop.conf import REST_CONN_TIMEOUT
  31. from desktop.lib import i18n
  32. from desktop.lib.view_util import format_duration_in_millis, location_to_url
  33. from jobbrowser.conf import DISABLE_KILLING_JOBS
  34. LOG = logging.getLogger(__name__)
  35. def can_view_job(username, job):
  36. acl = get_acls(job).get('mapreduce.job.acl-view-job', '')
  37. return acl == '*' or username in acl.split(',')
  38. def can_modify_job(username, job):
  39. acl = get_acls(job).get('mapreduce.job.acl-modify-job', '')
  40. return acl == '*' or username in acl.split(',')
  41. def get_acls(job):
  42. if job.is_mr2:
  43. try:
  44. acls = job.acls
  45. except:
  46. LOG.exception('failed to get acls')
  47. acls = {}
  48. return acls
  49. else:
  50. return job.full_job_conf
  51. def can_kill_job(self, user):
  52. if DISABLE_KILLING_JOBS.get():
  53. return False
  54. if self.status.lower() not in ('running', 'pending', 'accepted'):
  55. return False
  56. if is_admin(user):
  57. return True
  58. if can_modify_job(user.username, self):
  59. return True
  60. return user.username == self.user
  61. # You'll have to do the following manually to clean this up:
  62. # * Rearrange models' order
  63. # * Make sure each model has one field with primary_key=True
  64. # * Make sure each ForeignKey has `on_delete` set to the desired behavior.
  65. # * Remove `managed = False` lines if you wish to allow Django to create, modify, and delete the table
  66. # Feel free to rename the models, but don't rename db_table values or field names.
  67. class HiveQuery(models.Model):
  68. # (mysql.E001) MySQL does not allow unique CharFields to have a max_length > 255.
  69. # query_id = models.CharField(unique=True, max_length=512, blank=True, null=True)
  70. id = models.IntegerField(unique=True, blank=True, null=False, primary_key=True)
  71. query_id = models.CharField(unique=True, max_length=255, blank=True, null=True)
  72. query = models.TextField(blank=True, null=True)
  73. query_fts = models.TextField(blank=True, null=True) # This field type is a guess.
  74. start_time = models.BigIntegerField(blank=True, null=True)
  75. end_time = models.BigIntegerField(blank=True, null=True)
  76. elapsed_time = models.BigIntegerField(blank=True, null=True)
  77. status = models.CharField(max_length=32, blank=True, null=True)
  78. queue_name = models.CharField(max_length=767, blank=True, null=True)
  79. user_id = models.CharField(max_length=256, blank=True, null=True)
  80. request_user = models.CharField(max_length=256, blank=True, null=True)
  81. cpu_time = models.BigIntegerField(blank=True, null=True)
  82. physical_memory = models.BigIntegerField(blank=True, null=True)
  83. virtual_memory = models.BigIntegerField(blank=True, null=True)
  84. data_read = models.BigIntegerField(blank=True, null=True)
  85. data_written = models.BigIntegerField(blank=True, null=True)
  86. operation_id = models.CharField(max_length=512, blank=True, null=True)
  87. client_ip_address = models.CharField(max_length=64, blank=True, null=True)
  88. hive_instance_address = models.CharField(max_length=512, blank=True, null=True)
  89. hive_instance_type = models.CharField(max_length=512, blank=True, null=True)
  90. session_id = models.CharField(max_length=512, blank=True, null=True)
  91. log_id = models.CharField(max_length=512, blank=True, null=True)
  92. thread_id = models.CharField(max_length=512, blank=True, null=True)
  93. execution_mode = models.CharField(max_length=16, blank=True, null=True)
  94. databases_used = models.TextField(blank=True, null=True) # This field type is a guess.
  95. tables_read = models.TextField(blank=True, null=True) # This field type is a guess.
  96. tables_written = models.TextField(blank=True, null=True) # This field type is a guess.
  97. domain_id = models.CharField(max_length=512, blank=True, null=True)
  98. llap_app_id = models.CharField(max_length=512, blank=True, null=True)
  99. used_cbo = models.CharField(max_length=16, blank=True, null=True)
  100. first_task_started_time = models.BigIntegerField(blank=True, null=True)
  101. waiting_time = models.BigIntegerField(blank=True, null=True)
  102. resource_utilization = models.BigIntegerField(blank=True, null=True)
  103. version = models.SmallIntegerField(blank=True, null=True)
  104. created_at = models.DateTimeField(blank=True, null=True)
  105. class Meta:
  106. managed = False
  107. db_table = 'hive_query'
  108. class QueryDetails(models.Model):
  109. hive_query = models.ForeignKey(HiveQuery, on_delete=models.CASCADE, unique=True, blank=True, null=True)
  110. explain_plan_raw = models.TextField(blank=True, null=True) # This field type is a guess.
  111. configuration_raw = models.TextField(blank=True, null=True) # This field type is a guess.
  112. perf = models.TextField(blank=True, null=True) # This field type is a guess.
  113. configuration_compressed = models.BinaryField(blank=True, null=True)
  114. explain_plan_compressed = models.BinaryField(blank=True, null=True)
  115. class Meta:
  116. managed = False
  117. db_table = 'query_details'
  118. class DagInfo(models.Model):
  119. # (mysql.E001) MySQL does not allow unique CharFields to have a max_length > 255.
  120. # dag_id = models.CharField(unique=True, max_length=512, blank=True, null=True)
  121. dag_id = models.CharField(unique=True, max_length=255, blank=True, null=True)
  122. dag_name = models.CharField(max_length=512, blank=True, null=True)
  123. application_id = models.CharField(max_length=512, blank=True, null=True)
  124. init_time = models.BigIntegerField(blank=True, null=True)
  125. start_time = models.BigIntegerField(blank=True, null=True)
  126. end_time = models.BigIntegerField(blank=True, null=True)
  127. time_taken = models.BigIntegerField(blank=True, null=True)
  128. status = models.CharField(max_length=64, blank=True, null=True)
  129. am_webservice_ver = models.CharField(max_length=16, blank=True, null=True)
  130. am_log_url = models.CharField(max_length=512, blank=True, null=True)
  131. queue_name = models.CharField(max_length=64, blank=True, null=True)
  132. caller_id = models.CharField(max_length=512, blank=True, null=True)
  133. caller_type = models.CharField(max_length=128, blank=True, null=True)
  134. hive_query = models.ForeignKey('HiveQuery', on_delete=models.CASCADE, blank=True, null=True)
  135. created_at = models.DateTimeField(blank=True, null=True)
  136. source_file = models.TextField(blank=True, null=True)
  137. class Meta:
  138. managed = False
  139. db_table = 'dag_info'
  140. class DagDetails(models.Model):
  141. dag_info = models.ForeignKey('DagInfo', on_delete=models.CASCADE, unique=True, blank=True, null=True)
  142. hive_query = models.ForeignKey('HiveQuery', on_delete=models.CASCADE, blank=True, null=True)
  143. dag_plan_raw = models.TextField(blank=True, null=True) # This field type is a guess.
  144. vertex_name_id_mapping_raw = models.TextField(blank=True, null=True) # This field type is a guess.
  145. diagnostics = models.TextField(blank=True, null=True)
  146. counters_raw = models.TextField(blank=True, null=True) # This field type is a guess.
  147. dag_plan_compressed = models.BinaryField(blank=True, null=True)
  148. vertex_name_id_mapping_compressed = models.BinaryField(blank=True, null=True)
  149. counters_compressed = models.BinaryField(blank=True, null=True)
  150. class Meta:
  151. managed = False
  152. db_table = 'dag_details'
  153. class LinkJobLogs(object):
  154. @classmethod
  155. def _make_hdfs_links(cls, log, is_embeddable=False):
  156. escaped_logs = escape(log)
  157. return re.sub('((?<= |;)/|hdfs://)[^ <&\t;,\n]+', functools.partial(LinkJobLogs._replace_hdfs_link, is_embeddable), escaped_logs)
  158. @classmethod
  159. def _make_mr_links(cls, log):
  160. escaped_logs = escape(log)
  161. return re.sub('(job_[0-9]{12,}_[0-9]+)', LinkJobLogs._replace_mr_link, escaped_logs)
  162. @classmethod
  163. def _make_links(cls, log, is_embeddable=False):
  164. escaped_logs = escape(log)
  165. hdfs_links = re.sub('((?<= |;)/|hdfs://)[^ <&\t;,\n]+', functools.partial(LinkJobLogs._replace_hdfs_link, is_embeddable), escaped_logs)
  166. return re.sub('(job_[0-9]{12,}_[0-9]+)', LinkJobLogs._replace_mr_link, hdfs_links)
  167. @classmethod
  168. def _replace_hdfs_link(self, is_embeddable=False, match=None):
  169. try:
  170. return '<a href="%s">%s</a>' % (location_to_url(match.group(0), strict=False, is_embeddable=is_embeddable), match.group(0))
  171. except:
  172. LOG.exception('failed to replace hdfs links: %s' % (match.groups(),))
  173. return match.group(0)
  174. @classmethod
  175. def _replace_mr_link(self, match):
  176. try:
  177. return '<a href="/hue%s">%s</a>' % (reverse('jobbrowser:jobbrowser.views.single_job', kwargs={'job': match.group(0)}), match.group(0))
  178. except:
  179. LOG.exception('failed to replace mr links: %s' % (match.groups(),))
  180. return match.group(0)
  181. def format_unixtime_ms(unixtime):
  182. """
  183. Format a unix timestamp in ms to a human readable string
  184. """
  185. if unixtime:
  186. return str(datetime.datetime.fromtimestamp(math.floor(unixtime / 1000)).strftime("%x %X %Z"))
  187. else:
  188. return ""