Эх сурвалжийг харах

[notebook] Add a Pig snippet

Romain Rigaux 10 жил өмнө
parent
commit
4bc5133b77

+ 41 - 0
apps/pig/src/pig/models.py

@@ -90,6 +90,47 @@ class PigScript(Document):
     return 'org.apache.pig.backend.hadoop.hbase.HBaseStorage' in script
 
 
+class PigScript2(object):
+
+  def __init__(self, attrs=None):
+    self.data = json.dumps({
+        'script': '',
+        'name': '',
+        'properties': [],
+        'job_id': None,
+        'parameters': [],
+        'resources': [],
+        'hadoopProperties': []
+    })
+
+    if attrs:
+      self.update_from_dict(attrs)
+
+  def update_from_dict(self, attrs):
+    data_dict = self.dict
+
+    data_dict.update(attrs)
+
+    self.data = json.dumps(data_dict)
+
+  @property
+  def dict(self):
+    return json.loads(self.data)
+
+  @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) or \
+        ('org.apache.hive.hcatalog.pig.HCatLoader' in script or 'org.apache.hive.hcatalog.pig.HCatStorer' in script) # New classes
+
+  @property
+  def use_hbase(self):
+    script = self.dict['script']
+
+    return 'org.apache.pig.backend.hadoop.hbase.HBaseStorage' in script
+
+
 def create_or_update_script(id, name, script, user, parameters, resources, hadoopProperties, is_design=True):
   try:
     pig_script = PigScript.objects.get(id=id)

+ 4 - 0
desktop/conf.dist/hue.ini

@@ -551,6 +551,10 @@
     name=Spark Submit Python
     interface=livy-batch
 
+    [[[pig]]]
+    name=Pig
+    interface=pig
+
     [[[text]]]
     name=Text
     interface=text

+ 4 - 0
desktop/conf/pseudo-distributed.ini.tmpl

@@ -565,6 +565,10 @@
     name=Spark Submit Python
     interface=livy-batch
 
+    [[[pig]]]
+    name=Pig
+    interface=pig
+
     [[[text]]]
     name=Text
     interface=text

+ 10 - 10
desktop/libs/notebook/src/notebook/api.py

@@ -50,7 +50,7 @@ def create_session(request):
     if any(old_session) and 'properties' in old_session[0]:
       properties = old_session[0]['properties']
 
-  response['session'] = get_api(request.user, session).create_session(lang=session['type'], properties=properties)
+  response['session'] = get_api(request.user, session, request.fs, request.jt).create_session(lang=session['type'], properties=properties)
   response['session']['properties'] = properties
   response['status'] = 0
 
@@ -65,7 +65,7 @@ def close_session(request):
 
   session = json.loads(request.POST.get('session', '{}'))
 
-  response['session'] = get_api(request.user, {'type': session['type']}).close_session(session=session)
+  response['session'] = get_api(request.user, {'type': session['type']}, request.fs, request.jt).close_session(session=session)
   response['status'] = 0
 
   return JsonResponse(response)
@@ -80,7 +80,7 @@ def execute(request):
   notebook = json.loads(request.POST.get('notebook', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
-  response['handle'] = get_api(request.user, snippet).execute(notebook, snippet)
+  response['handle'] = get_api(request.user, snippet, request.fs, request.jt).execute(notebook, snippet)
   response['status'] = 0
 
   return JsonResponse(response)
@@ -95,7 +95,7 @@ def check_status(request):
   notebook = json.loads(request.POST.get('notebook', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
-  response['query_status'] = get_api(request.user, snippet).check_status(notebook, snippet)
+  response['query_status'] = get_api(request.user, snippet, request.fs, request.jt).check_status(notebook, snippet)
   response['status'] = 0
 
   return JsonResponse(response)
@@ -112,7 +112,7 @@ def fetch_result_data(request):
   rows = json.loads(request.POST.get('rows', 100))
   start_over = json.loads(request.POST.get('startOver', False))
 
-  response['result'] = get_api(request.user, snippet).fetch_result(notebook, snippet, rows, start_over)
+  response['result'] = get_api(request.user, snippet, request.fs, request.jt).fetch_result(notebook, snippet, rows, start_over)
   response['status'] = 0
 
   return JsonResponse(response)
@@ -127,7 +127,7 @@ def fetch_result_metadata(request):
   notebook = json.loads(request.POST.get('notebook', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
-  response['result'] = get_api(request.user, snippet).fetch_result_metadata(notebook, snippet)
+  response['result'] = get_api(request.user, snippet, request.fs, request.jt).fetch_result_metadata(notebook, snippet)
   response['status'] = 0
 
   return JsonResponse(response)
@@ -142,7 +142,7 @@ def cancel_statement(request):
   notebook = json.loads(request.POST.get('notebook', '{}'))
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
-  response['result'] = get_api(request.user, snippet).cancel(notebook, snippet)
+  response['result'] = get_api(request.user, snippet, request.fs, request.jt).cancel(notebook, snippet)
   response['status'] = 0
 
   return JsonResponse(response)
@@ -163,7 +163,7 @@ def get_logs(request):
   size = request.POST.get('size')
   size = int(size) if size else None
 
-  db = get_api(request.user, snippet)
+  db = get_api(request.user, snippet, request.fs, request.jt)
   response['logs'] = db.get_log(notebook, snippet, startFrom=startFrom, size=size)
   response['progress'] = db._progress(snippet, response['logs']) if snippet['status'] != 'available' and snippet['status'] != 'success' else 100
   response['job_urls'] = [{
@@ -226,7 +226,7 @@ def close_notebook(request):
 
   for session in notebook['sessions']:
     try:
-      response['result'].append(get_api(request.user, session).close_session(session))
+      response['result'].append(get_api(request.user, session, request.fs, request.jt).close_session(session))
     except QueryExpired:
       pass
     except Exception, e:
@@ -248,7 +248,7 @@ def close_statement(request):
   snippet = json.loads(request.POST.get('snippet', '{}'))
 
   try:
-    response['result'] = get_api(request.user, snippet).close_statement(snippet)
+    response['result'] = get_api(request.user, snippet, request.fs, request.jt).close_statement(snippet)
   except QueryExpired:
     pass
 

+ 14 - 5
desktop/libs/notebook/src/notebook/connectors/base.py

@@ -77,13 +77,15 @@ class Notebook():
     return '\n\n'.join([snippet['statement_raw'] for snippet in self.get_data()['snippets']])
 
 
-def get_api(user, snippet):
-  from notebook.connectors.hiveserver2 import HS2Api  
-  from notebook.connectors.mysql import MySqlApi
+def get_api(user, snippet, fs, jt):
+  from notebook.connectors.hiveserver2 import HS2Api
   from notebook.connectors.jdbc import JDBCApi
-  from notebook.connectors.text import TextApi
+  from notebook.connectors.mysql import MySqlApi
+  from notebook.connectors.pig_batch import PigApi
   from notebook.connectors.spark_shell import SparkApi
   from notebook.connectors.spark_batch import SparkBatchApi
+  from notebook.connectors.text import TextApi
+
 
   interface = [interpreter for interpreter in get_interpreters() if interpreter['type'] == snippet['type']]
   if not interface:
@@ -102,6 +104,8 @@ def get_api(user, snippet):
     return MySqlApi(user)
   elif interface == 'jdbc':
     return JDBCApi(user)
+  elif interface == 'pig':
+    return PigApi(user, fs, jt)
   else:
     raise PopupException(_('Notebook connector interface not recognized: %s') % interface)
 
@@ -114,8 +118,10 @@ def _get_snippet_session(notebook, snippet):
 
 class Api(object):
 
-  def __init__(self, user):
+  def __init__(self, user, fs=None, jt=None):
     self.user = user
+    self.fs = fs
+    self.jt = jt
 
   def create_session(self, lang, properties=None):
     return {
@@ -130,5 +136,8 @@ class Api(object):
   def fetch_result(self, notebook, snippet, rows, start_over):
     pass
 
+  def download(self, notebook, snippet, format):
+    pass
+
   def get_log(self, notebook, snippet, startFrom=None, size=None):
     return 'No logs'

+ 132 - 0
desktop/libs/notebook/src/notebook/connectors/pig_batch.py

@@ -0,0 +1,132 @@
+#!/usr/bin/env python
+# Licensed to Cloudera, Inc. under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  Cloudera, Inc. licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#     http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+
+import logging
+import json
+
+from django.utils.translation import ugettext as _
+from django.core.urlresolvers import reverse
+
+from notebook.connectors.base import Api, QueryError
+
+
+LOG = logging.getLogger(__name__)
+
+
+try:
+  from pig import api
+  from pig.models import PigScript2, get_workflow_output, hdfs_link
+  from oozie.views.dashboard import check_job_access_permission, check_job_edition_permission
+except ImportError, e:
+  LOG.exception('Pig application is not enabled')
+
+
+class PigApi(Api):
+
+  def execute(self, notebook, snippet):
+
+    attrs = {
+      'script': snippet['statement'],
+      'name': snippet['properties'].get('name', 'Pig Snippet'),
+      'parameters': snippet['properties'].get('parameters'),
+      'resources': snippet['properties'].get('resources'),
+      'hadoopProperties': snippet['properties'].get('hadoopProperties')
+    }
+
+    pig_script = PigScript2(attrs)
+
+    params = json.dumps([])
+    oozie_id = api.get(self.fs, self.jt, self.user).submit(pig_script, params)
+
+    return {
+      'id': oozie_id,
+      'watchUrl': reverse('pig:watch', kwargs={'job_id': oozie_id}) + '?format=python'
+    }
+
+  def check_status(self, notebook, snippet):
+    job_id = snippet['result']['handle']['id']
+    request = MockRequest(self.user, self.fs, self.jt)
+
+    oozie_workflow = check_job_access_permission(request, job_id)
+    logs, workflow_actions, is_really_done = api.get(self.jt, self.jt, self.user).get_log(request, oozie_workflow)
+
+    if is_really_done and not oozie_workflow.is_running():
+      if oozie_workflow.status in ('KILLED', 'FAILED'):
+        raise QueryError(_('The script failed to run and was stopped'))
+      status = 'available'
+    elif oozie_workflow.is_running():
+      status = 'running'
+    else:
+      status = 'failed'
+
+    return {
+        'status': status
+    }
+
+  def fetch_result(self, notebook, snippet, rows, start_over):
+    job_id = snippet['result']['handle']['id']
+
+    oozie_workflow = check_job_access_permission(MockRequest(self.user, self.fs, self.jt), job_id)
+    output = get_workflow_output(oozie_workflow, self.fs)
+
+    return {
+        'data':  [hdfs_link(output)],
+        'meta': [{'name': 'Header', 'type': 'STRING_TYPE', 'comment': ''}],
+        'type': 'text'
+    }
+
+  def cancel(self, notebook, snippet):
+    job_id = snippet['result']['handle']['id']
+
+    job = check_job_access_permission(self, job_id)
+    check_job_edition_permission(job, self.user)
+
+    api.get(self.fs, self.jt, self.user).stop(job_id)
+
+    return {'status': 0}
+
+  def get_log(self, notebook, snippet, startFrom=0, size=None):
+    job_id = snippet['result']['handle']['id']
+    request = MockRequest(self.user, self.fs, self.jt)
+
+    oozie_workflow = check_job_access_permission(MockRequest(self.user, self.fs, self.jt), job_id)
+    logs, workflow_actions, is_really_done = api.get(self.jt, self.jt, self.user).get_log(request, oozie_workflow)
+
+    return logs
+
+  def _progress(self, snippet, logs):
+    job_id = snippet['result']['handle']['id']
+
+    oozie_workflow = check_job_access_permission(MockRequest(self.user, self.fs, self.jt), job_id)
+    return oozie_workflow.get_progress(),
+
+  def close_statement(self, snippet):
+    pass
+
+  def close_session(self, session):
+    pass
+
+  def _get_jobs(self, log):
+    return []
+
+
+class MockRequest():
+
+  def __init__(self, user, fs, jt):
+    self.user = user
+    self.fs = fs
+    self.js = jt

+ 3 - 2
desktop/libs/notebook/src/notebook/static/notebook/js/notebook.ko.js

@@ -134,9 +134,10 @@ var getDefaultSnippetProperties = function (snippetType) {
     properties['files'] = [];
   }
   else if (snippetType == 'pig') {
+    properties['script'] = '';
     properties['parameters'] = [];
-    properties['hadoop_properties'] = [];
-    properties['files'] = [];
+    properties['hadoopProperties'] = [];
+    properties['resources'] = [];
   }
 
   return properties;