浏览代码

[pig] Support HCatalog transparently

Fix 500 error in Job Browser when the attemps of past jobs are gone.
Romain Rigaux 11 年之前
父节点
当前提交
97f80a4

+ 1 - 1
apps/jobbrowser/src/jobbrowser/views.py

@@ -279,7 +279,7 @@ def job_single_logs(request, job):
     if recent_tasks:
       task = recent_tasks[0]
 
-  if task is None:
+  if task is None or not task.taskAttemptIds:
     raise PopupException(_("No tasks found for job %(id)s.") % {'id': job.jobId})
 
   return single_task_attempt_logs(request, **{'job': job.jobId, 'taskid': task.taskId, 'attemptid': task.taskAttemptIds[-1]})

+ 5 - 1
apps/jobbrowser/src/jobbrowser/yarn_models.py

@@ -159,7 +159,11 @@ class Task:
   def attempts(self):
     # We can cache as we deal with history server
     if not hasattr(self, '_attempts'):
-      self._attempts = [Attempt(self, attempt) for attempt in self.job.api.task_attempts(self.job.id, self.id)['taskAttempts']['taskAttempt']]
+      task_attempts = self.job.api.task_attempts(self.job.id, self.id)['taskAttempts']
+      if task_attempts:
+        self._attempts = [Attempt(self, attempt) for attempt in task_attempts['taskAttempt']]
+      else:
+        self._attempts = []
     return self._attempts
 
   @property

+ 5 - 0
apps/oozie/src/oozie/models.py

@@ -182,6 +182,11 @@ class Job(models.Model):
   def get_parameters(self):
     return json.loads(self.parameters)
 
+  def add_parameter(self, name, value):
+    oozie_parameters = self.get_parameters()
+    oozie_parameters.append({"name": name, "value": value})
+    self.parameters = json.dumps(oozie_parameters)
+
   @property
   def parameters_escapejs(self):
     return self._escapejs_parameters_list(self.parameters)

+ 5 - 5
apps/pig/src/pig/api.py

@@ -54,14 +54,11 @@ class OozieApi:
     self.user = user
 
   def submit(self, pig_script, params):
-    mapping = {
-      'oozie.use.system.libpath':  'true',
-    }
-
     workflow = None
 
     try:
       workflow = self._create_workflow(pig_script, params)
+      mapping = dict([(param['name'], param['value']) for param in workflow.get_parameters()])
       oozie_wf = _submit_workflow(self.user, self.fs, self.jt, workflow, mapping)
     finally:
       if workflow:
@@ -73,11 +70,14 @@ class OozieApi:
     workflow = Workflow.objects.new_workflow(self.user)
     workflow.name = OozieApi.WORKFLOW_NAME
     workflow.is_history = True
+    if pig_script.use_hcatalog:
+      workflow.add_parameter("oozie.action.sharelib.for.pig", "pig,hcatalog")
     workflow.save()
     Workflow.objects.initialize(workflow, self.fs)
 
     script_path = workflow.deployment_dir + '/script.pig'
-    self.fs.do_as_user(self.user.username, self.fs.create, script_path, data=pig_script.dict['script'])
+    if self.fs: # For testing, difficult to mock
+      self.fs.do_as_user(self.user.username, self.fs.create, script_path, data=pig_script.dict['script'])
 
     files = []
     archives = []

+ 6 - 4
apps/pig/src/pig/models.py

@@ -15,10 +15,7 @@
 # See the License for the specific language governing permissions and
 # limitations under the License.
 
-try:
-  import json
-except ImportError:
-  import simplejson as json
+import json
 import posixpath
 
 from django.db import models
@@ -81,6 +78,11 @@ class PigScript(Document):
   def get_absolute_url(self):
     return reverse('pig:index') + '#edit/%s' % self.id
 
+  @property
+  def use_hcatalog(self):
+    script = self.dict['script']
+    return 'org.apache.hcatalog.pig.HCatStorer' in script or 'org.apache.hcatalog.pig.HCatLoader' in script
+
 
 def create_or_update_script(id, name, script, user, parameters, resources, hadoopProperties, is_design=True):
   try:

+ 1 - 1
apps/pig/src/pig/templates/app.mako

@@ -349,7 +349,7 @@ ${ commonheader(None, "pig", user) | n,unicode }
           <br/>
           <br/>
 
-          <h4>${ _('Parameters') } &nbsp; <i id="parameters-dyk" class="fa fa-question-circle"></i></h4>
+          <h4>${ _('Pig parameters') } &nbsp; <i id="parameters-dyk" class="fa fa-question-circle"></i></h4>
           <div id="parameters-dyk-content" class="hide">
             <ul style="text-align: left;">
               <li>input /user/data</li>

+ 20 - 1
apps/pig/src/pig/tests.py

@@ -22,7 +22,7 @@ import time
 from django.contrib.auth.models import User
 from django.core.urlresolvers import reverse
 
-from nose.tools import assert_true, assert_equal
+from nose.tools import assert_true, assert_equal, assert_false
 
 from desktop.lib.django_test_util import make_logged_in_client
 from desktop.lib.test_utils import grant_access
@@ -128,6 +128,25 @@ class TestMock(TestPigBase):
     pig_script = self.create_script()
     assert_equal('Test', pig_script.dict['name'])
 
+  def test_check_hcatalogs_sharelib(self):
+    api = get(None, None, self.user)
+    pig_script = self.create_script()
+
+    # Regular
+    wf = api._create_workflow(pig_script, '[]')
+    assert_false({'name': u'oozie.action.sharelib.for.pig', 'value': u'pig,hcatalog'} in wf.find_all_parameters(), wf.find_all_parameters())
+
+    # With HCat
+    pig_script.update_from_dict({
+        'script':"""
+           a = LOAD 'sample_07' USING org.apache.hcatalog.pig.HCatLoader();
+           dump a;
+    """})
+    pig_script.save()
+
+    wf = api._create_workflow(pig_script, '[]')
+    assert_true({'name': u'oozie.action.sharelib.for.pig', 'value': u'pig,hcatalog'} in wf.find_all_parameters(), wf.find_all_parameters())
+
   def test_editor_view(self):
     response = self.c.get(reverse('pig:app'))
     assert_true('Unsaved script' in response.content)

+ 2 - 1
apps/pig/static/js/pig.ko.js

@@ -546,7 +546,8 @@ var PigViewModel = function (props) {
       function (data) {
         $(document).trigger("stopped");
         $("#stopModal").modal("hide");
-      }, "json");
+      }, "json"
+    );
   }
 
   function callCopy(script) {