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

HUE-2551 [oozie] Easier update of a scheduled workflow

krish 10 жил өмнө
parent
commit
136c060647

+ 16 - 0
apps/oozie/src/oozie/templates/dashboard/list_oozie_coordinator.mako

@@ -149,6 +149,13 @@ ${ layout.menubar(section='coordinators', dashboard=True) }
                   </button>
                 </li>
               % endif
+              <li class="white">
+                <button title="${ _('Sync Workflow') }" id="sync-wf-btn"
+                     data-sync-url="${ url('oozie:sync_coord_workflow', job_id=oozie_coordinator.id) }"
+                     class="btn btn-small sync-wf-btn" style="margin-bottom: 5px">
+                    ${ _('Sync Workflow') }
+                  </button>
+              </li>
             </ul>
           </div>
         </div>
@@ -865,6 +872,15 @@ ${ layout.menubar(section='coordinators', dashboard=True) }
       });
     });
 
+    $('#sync-wf-btn, .sync-wf-btn').click(function () {
+
+      $.get($(this).data("sync-url"), function (response) {
+        $('#rerun-coord-modal').html(response);
+        $('#rerun-coord-modal').modal('show');
+      });
+
+    });
+
     function refreshActionsPagination() {
       actionTableOffset = 1;
     }

+ 5 - 1
apps/oozie/src/oozie/templates/editor2/submit_job_popup.mako

@@ -26,7 +26,11 @@
   ${ csrf_token(request) | n,unicode }
   <div class="modal-header">
     <a href="#" class="close" data-dismiss="modal">&times;</a>
-    <h3>${ _('Submit %(job)s?') % {'job': name} }</h3>
+    % if header:
+      <h3>${header}</h3>
+    % else:
+      <h3>${ _('Submit %(job)s?') % {'job': name} }</h3>
+    % endif
   </div>
   <div class="modal-body">
 

+ 87 - 1
apps/oozie/src/oozie/tests.py

@@ -37,6 +37,7 @@ from desktop.lib.django_test_util import make_logged_in_client
 from desktop.lib.test_utils import grant_access, add_permission, add_to_group, reformat_json, reformat_xml
 from desktop.models import Document, Document2
 
+from hadoop import cluster as originalCluster
 from hadoop.pseudo_hdfs4 import is_live_cluster
 from jobsub.models import OozieDesign, OozieMapreduceAction
 from liboozie import oozie_api
@@ -65,7 +66,8 @@ class MockOozieApi:
                         {u'status': u'KILLED', u'run': 0, u'startTime': u'Mon, 30 Jul 2012 22:31:08 GMT', u'appName': u'WordCount2', u'lastModTime': u'Mon, 30 Jul 2012 22:32:20 GMT', u'actions': [], u'acl': None, u'appPath': None, u'externalId': '-', u'consoleUrl': u'http://runreal:11000/oozie?job=0000011-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Mon, 30 Jul 2012 22:31:08 GMT', u'toString': u'Workflow id[0000011-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Mon, 30 Jul 2012 22:32:20 GMT', u'id': u'0000011-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'},
                         {u'status': u'SUCCEEDED', u'run': 0, u'startTime': u'Mon, 30 Jul 2012 22:20:48 GMT', u'appName': u'WordCount3', u'lastModTime': u'Mon, 30 Jul 2012 22:22:00 GMT', u'actions': [], u'acl': None, u'appPath': None, u'externalId': '', u'consoleUrl': u'http://runreal:11000/oozie?job=0000009-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Mon, 30 Jul 2012 22:20:48 GMT', u'toString': u'Workflow id[0000009-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Mon, 30 Jul 2012 22:22:00 GMT', u'id': u'0000009-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'},
                         {u'status': u'SUCCEEDED', u'run': 0, u'startTime': u'Mon, 30 Jul 2012 22:16:58 GMT', u'appName': u'WordCount4', u'lastModTime': u'Mon, 30 Jul 2012 22:18:10 GMT', u'actions': [], u'acl': None, u'appPath': None, u'externalId': None, u'consoleUrl': u'http://runreal:11000/oozie?job=0000008-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Mon, 30 Jul 2012 22:16:58 GMT', u'toString': u'Workflow id[0000008-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Mon, 30 Jul 2012 22:18:10 GMT', u'id': u'0000008-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'},
-                        {u'status': u'SUCCEEDED', u'run': 0, u'startTime': u'Fri, 24 Jul 2015 14:56:08 AEST', u'appName': u'WordCount5', u'lastModTime': u'Fri, 24 Jul 2015 14:57:17 AEST', u'actions': [], u'acl': None, u'appPath': None, u'externalId': None, u'consoleUrl': u'http://runreal:11000/oozie?job=0000008-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Fri, 24 Jul 2015 14:56:08 AEST', u'toString': u'Workflow id[0000007-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Fri, 24 Jul 2015 14:57:17 AEST', u'id': u'0000007-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'}]
+                        {u'status': u'SUCCEEDED', u'run': 0, u'startTime': u'Fri, 24 Jul 2015 14:56:08 AEST', u'appName': u'WordCount5', u'lastModTime': u'Fri, 24 Jul 2015 14:57:17 AEST', u'actions': [], u'acl': None, u'appPath': None, u'externalId': None, u'consoleUrl': u'http://runreal:11000/oozie?job=0000008-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Fri, 24 Jul 2015 14:56:08 AEST', u'toString': u'Workflow id[0000007-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Fri, 24 Jul 2015 14:57:17 AEST', u'id': u'0000007-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'},
+                        {u'status': u'RUNNING', u'run': 0, u'startTime': u'Mon, 30 Jul 2012 22:35:48 GMT', u'appName': u'WordCount1', u'lastModTime': u'Mon, 30 Jul 2012 22:37:00 GMT', u'actions': [], u'acl': None, u'appPath': None, u'externalId': 'job_201208072118_0044', u'consoleUrl': u'http://runreal:11000/oozie?job=0000006-120725142744176-oozie-oozi-W', u'conf': None, u'parentId': None, u'createdTime': u'Mon, 30 Jul 2012 22:35:48 GMT', u'toString': u'Workflow id[0000006-120725142744176-oozie-oozi-W] status[SUCCEEDED]', u'endTime': u'Mon, 30 Jul 2012 22:37:00 GMT', u'id': u'0000006-120725142744176-oozie-oozi-W', u'group': None, u'user': u'test'}]
 
   WORKFLOW_IDS = [wf['id'] for wf in JSON_WORKFLOW_LIST]
   WORKFLOW_DICT = dict([(wf['id'], wf) for wf in JSON_WORKFLOW_LIST])
@@ -232,6 +234,27 @@ class MockOozieApi:
     return {'status': "RUNNING"}
 
 
+class MockFs():
+  def __init__(self, logical_name=None):
+
+    self.fs_defaultfs = 'hdfs://curacao:8020'
+    self.logical_name = logical_name if logical_name else ''
+
+  def setuser(self, user):
+    pass
+
+  def join(self, path1, path2):
+    return path1 + path2
+
+  def do_as_user(self, username, fn, *args, **kwargs):
+    return ''
+
+  def read(self, length=1024*1024):
+    return 'data'
+
+  def exists(self, path):
+    return True
+
 class OozieMockBase(object):
 
   def setUp(self):
@@ -249,14 +272,20 @@ class OozieMockBase(object):
     add_to_group("test")
     self.user = User.objects.get(username='test')
     self.wf = create_workflow(self.c, self.user)
+    self.original_fs = originalCluster.FS_CACHE["default"]
+    originalCluster.FS_CACHE["default"] = MockFs()
 
 
   def tearDown(self):
     oozie_api.OozieApi = oozie_api.OriginalOozieApi
+    if originalCluster.FS_CACHE is None:
+      originalCluster.FS_CACHE = {}
+    originalCluster.FS_CACHE["default"] = self.original_fs
     Workflow.objects.check_workspace = Workflow.objects.original_check_workspace
     oozie_api._api_cache = None
 
     History.objects.all().delete()
+
     for coordinator in Coordinator.objects.all():
       coordinator.delete(skip_trash=True)
     for bundle in Bundle.objects.all():
@@ -3283,6 +3312,15 @@ class TestDashboard(OozieMockBase):
     response = self.c.get(reverse('oozie:rerun_oozie_coord', args=[MockOozieApi.WORKFLOW_IDS[0], '/path']))
     assert_true('Rerun' in response.content, response.content)
 
+  def test_sync_coord_workflow(self):
+    wf_doc = save_temp_workflow(MockOozieApi.JSON_WORKFLOW_LIST[5], self.user)
+    reset = ENABLE_V2.set_for_testing(True)
+    try:
+      response = self.c.get(reverse('oozie:sync_coord_workflow', args=[MockOozieApi.WORKFLOW_IDS[5]]))
+      assert_equal([{'name':'Dryrun', 'value': False}, {'name':'ls_arg', 'value': '-l'}], response.context['params_form'].initial)
+    finally:
+      wf_doc.delete()
+      reset()
 
   def test_rerun_coordinator_permissions(self):
     post_data = {
@@ -3914,3 +3952,51 @@ def synchronize_workflow_attributes(workflow_json, correct_workflow_json):
     workflow_dict['attributes']['deployment_dir'] = correct_workflow_dict['attributes']['deployment_dir']
 
   return reformat_json(workflow_dict)
+
+def save_temp_workflow(wf, user):
+    data = json.dumps({
+          'layout': [{
+              "size":12, "rows":[
+                  {"widgets":[{"size":12, "name":"Start", "id":"3f107997-04cc-8733-60a9-a4bb62cebffc", "widgetType":"start-widget", "properties":{}, "offset":0, "isLoading":False, "klass":"card card-widget span12"}]},
+                  {"widgets":[{"size":12, "name":"End", "id":"33430f0f-ebfa-c3ec-f237-3e77efa03d0a", "widgetType":"end-widget", "properties":{}, "offset":0, "isLoading":False, "klass":"card card-widget span12"}]},
+                  {"widgets":[{"size":12, "name":"Kill", "id":"17c9c895-5a16-7443-bb81-f34b30b21548", "widgetType":"kill-widget", "properties":{}, "offset":0, "isLoading":False, "klass":"card card-widget span12"}]}
+              ],
+              "drops":[ "temp"],
+              "klass":"card card-home card-column span12"
+          }],
+          'workflow': {
+              "id": None,
+              "uuid": None,
+              "name": "My Workflow",
+              "properties": {
+                  "deployment_dir": "",
+                  "description": "",
+                  "job_xml": "",
+                  "sla_enabled": False,
+                  "schema_version": "uri:oozie:workflow:0.5",
+                  "properties": [],
+                  "parameters": [{'name':'Dryrun', 'value': False}, {'name':'ls_arg', 'value': '-l'}],
+                  "sla":  [
+      {'key': 'enabled', 'value': False}, # Always first element
+      {'key': 'nominal-time', 'value': '${nominal_time}'},
+      {'key': 'should-start', 'value': ''},
+      {'key': 'should-end', 'value': '${30 * MINUTES}'},
+      {'key': 'max-duration', 'value': ''},
+      {'key': 'alert-events', 'value': ''},
+      {'key': 'alert-contact', 'value': ''},
+      {'key': 'notification-msg', 'value': ''},
+      {'key': 'upstream-apps', 'value': ''},
+  ],
+                  "show_arrows": True,
+                  "wf1_id": None
+              },
+              "nodes":[
+                  {"id":"3f107997-04cc-8733-60a9-a4bb62cebffc","name":"Start","type":"start-widget","properties":{},"children":[{'to': '33430f0f-ebfa-c3ec-f237-3e77efa03d0a'}]},
+                  {"id":"33430f0f-ebfa-c3ec-f237-3e77efa03d0a","name":"End","type":"end-widget","properties":{},"children":[]},
+                  {"id":"17c9c895-5a16-7443-bb81-f34b30b21548","name":"Kill","type":"kill-widget","properties":{'message': 'Action failed, error message[${wf:errorMessage(wf:lastErrorNode())}]'},"children":[]}
+              ]
+          }
+      })
+    workflow_doc = Document2.objects.create(name='test', type='oozie-workflow2', owner=user, data=data)
+    wf[u'conf'] = u'<configuration><property><name>hue-id-w</name><value>' + str(workflow_doc.id) + u'</value></property></configuration>'
+    return workflow_doc

+ 1 - 0
apps/oozie/src/oozie/urls.py

@@ -133,6 +133,7 @@ urlpatterns += patterns(
   url(r'^rerun_oozie_job/(?P<job_id>[-\w]+)/(?P<app_path>.+?)$', 'rerun_oozie_job', name='rerun_oozie_job'),
   url(r'^rerun_oozie_coord/(?P<job_id>[-\w]+)/(?P<app_path>.+?)$', 'rerun_oozie_coordinator', name='rerun_oozie_coord'),
   url(r'^rerun_oozie_bundle/(?P<job_id>[-\w]+)/(?P<app_path>.+?)$', 'rerun_oozie_bundle', name='rerun_oozie_bundle'),
+  url(r'^sync_coord_workflow/(?P<job_id>[-\w]+)$', 'sync_coord_workflow', name='sync_coord_workflow'),
   url(r'^manage_oozie_jobs/(?P<job_id>[-\w]+)/(?P<action>(start|suspend|resume|kill|rerun|change|ignore))$', 'manage_oozie_jobs', name='manage_oozie_jobs'),
   url(r'^bulk_manage_oozie_jobs/$', 'bulk_manage_oozie_jobs', name='bulk_manage_oozie_jobs'),
 

+ 44 - 1
apps/oozie/src/oozie/views/dashboard.py

@@ -41,7 +41,7 @@ from desktop.models import Document, Document2
 
 from liboozie.oozie_api import get_oozie
 from liboozie.credentials import Credentials
-from liboozie.submittion import Submission
+from liboozie.submission2 import Submission
 from liboozie.types import Workflow as OozieWorkflow, Coordinator as CoordinatorWorkflow, Bundle as BundleWorkflow
 
 from oozie.conf import OOZIE_JOBS_COUNT, ENABLE_CRON_SCHEDULING, ENABLE_V2
@@ -638,6 +638,49 @@ def massaged_sla_for_json(sla, request):
   return massaged_sla
 
 
+@show_oozie_error
+def sync_coord_workflow(request, job_id):
+  ParametersFormSet = formset_factory(ParameterForm, extra=0)
+  job = check_job_access_permission(request, job_id)
+  check_job_edition_permission(job, request.user)
+
+  hue_coord = get_history().get_coordinator_from_config(job.conf_dict)
+  hue_wf = (hue_coord and hue_coord.workflow) or get_history().get_workflow_from_config(job.conf_dict)
+
+  if request.method == 'POST':
+    params_form = ParametersFormSet(request.POST)
+    if params_form.is_valid():
+      mapping = dict([(param['name'], param['value']) for param in params_form.cleaned_data])
+
+      submission = Submission(user=request.user, job=hue_wf, fs=request.fs, jt=request.jt, properties=mapping)
+      submission._sync_definition(hue_wf.deployment_dir, mapping)
+
+      request.info(_('Successfully updated Workflow definition'))
+      return redirect(reverse('oozie:list_oozie_coordinator', kwargs={'job_id': job_id}))
+    else:
+      request.error(_('Invalid submission form: %s' % params_form.errors))
+  else:
+    parameters = hue_wf and hue_wf.find_all_parameters() or []
+    params_dict = dict([(param['name'], param['value']) for param in parameters])
+
+    submission = Submission(user=request.user, job=hue_wf, fs=request.fs, jt=request.jt, properties=None)
+    prev_properties = hue_wf and hue_wf.deployment_dir and \
+                      submission.get_external_parameters(request.fs.join(hue_wf.deployment_dir, hue_wf.XML_FILE_NAME)) or {}
+
+    for key, value in params_dict.iteritems():
+      params_dict[key] = prev_properties[key] if key in prev_properties.keys() else params_dict[key]
+
+    initial_params = ParameterForm.get_initial_params(params_dict)
+    params_form = ParametersFormSet(initial=initial_params)
+
+  popup = render('editor2/submit_job_popup.mako', request, {
+             'params_form': params_form,
+             'name': _('Job'),
+             'header': _('Sync Workflow definition?'),
+             'action': reverse('oozie:sync_coord_workflow', kwargs={'job_id': job_id})
+           }, force_template=True).content
+  return JsonResponse(popup, safe=False)
+
 @show_oozie_error
 def rerun_oozie_job(request, job_id, app_path):
   ParametersFormSet = formset_factory(ParameterForm, extra=0)

+ 17 - 6
desktop/libs/liboozie/src/liboozie/submission2.py

@@ -277,13 +277,9 @@ class Submission(object):
     Copy XML and the jar_path files from Java or MR actions to the deployment directory.
     This should run as the workflow user.
     """
-    xml_path = self.fs.join(deployment_dir, self.job.XML_FILE_NAME)
-    self.fs.create(xml_path, overwrite=True, permission=0644, data=smart_str(oozie_xml))
-    LOG.debug("Created %s" % (xml_path,))
 
-    properties_path = self.fs.join(deployment_dir, 'job.properties')
-    self.fs.create(properties_path, overwrite=True, permission=0644, data=smart_str('\n'.join(['%s=%s' % (key, val) for key, val in oozie_properties.iteritems()])))
-    LOG.debug("Created %s" % (properties_path,))
+    self._create_file(deployment_dir, self.job.XML_FILE_NAME, oozie_xml)
+    self._create_file(deployment_dir, 'job.properties', data='\n'.join(['%s=%s' % (key, val) for key, val in oozie_properties.iteritems()]))
 
     # List jar files
     files = []
@@ -337,7 +333,22 @@ class Submission(object):
     from oozie.models2 import Coordinator
     return Coordinator.PROPERTY_APP_PATH in self.properties
 
+  def _create_file(self, deployment_dir, file_name, data, do_as=False):
+   file_path = self.fs.join(deployment_dir, file_name)
+   if do_as:
+     self.fs.do_as_user(self.user, self.fs.create, file_path, overwrite=True, permission=0644, data=smart_str(data))
+   else:
+     self.fs.create(file_path, overwrite=True, permission=0644, data=smart_str(data))
+   LOG.debug("Created/Updated %s" % (file_path,))
 
+  def _sync_definition(self, deployment_dir, mapping):
+    """ This is helper function for 'Sync Workflow' functionality in a Coordinator.
+      It copies updated workflow changes into HDFS """
+
+    self._create_file(deployment_dir, self.job.XML_FILE_NAME, self.job.to_xml(mapping=mapping), do_as=True)
+
+    data_properties = smart_str('\n'.join(['%s=%s' % (key, val) for key, val in mapping.iteritems()]))
+    self._create_file(deployment_dir, 'job.properties', data_properties, do_as=True)
 
 def create_directories(fs, directory_list=[]):
   # If needed, create the remote home, deployment and data directories

+ 4 - 0
desktop/libs/liboozie/src/liboozie/submittion2_tests.py

@@ -135,6 +135,10 @@ def test_copy_files():
     assert_not_equal(stats_udf5['fileId'], cluster.fs.stats(deployment_dir + '/udf5.jar')['fileId'])
     assert_equal(stats_udf6['fileId'], cluster.fs.stats(deployment_dir + '/udf6.jar')['fileId'])
 
+    # Test _create_file()
+    submission._create_file(deployment_dir, 'test.txt', data='Test data')
+    assert_true(cluster.fs.exists(deployment_dir + '/test.txt'), list_dir_workspace)
+
   finally:
     try:
       cluster.fs.rmtree(prefix)