Quellcode durchsuchen

HUE-8259 [core] Integrating Celery

Prakash Ranade vor 6 Jahren
Ursprung
Commit
eaa8c1c

+ 27 - 0
celery_rabbitmq_install.txt

@@ -0,0 +1,27 @@
+Hue Celery/RabbitMQ Install
+
+Installation tested on Ubuntu 16.04
+
+echo "deb https://packages.erlang-solutions.com/ubuntu $(lsb_release -sc) contrib" >> /etc/apt/sources.list.d/erlang.list
+wget https://packages.erlang-solutions.com/ubuntu/erlang_solutions.asc
+sudo apt-key add erlang_solutions.asc
+
+echo "deb https://dl.bintray.com/rabbitmq/debian $(lsb_release -sc) main" >> /etc/apt/sources.list.d/rabbitmq.list
+wget -O- https://dl.bintray.com/rabbitmq/Keys/rabbitmq-release-signing-key.asc | sudo apt-key add -
+wget -O- https://www.rabbitmq.com/rabbitmq-release-signing-key.asc | sudo apt-key add -
+
+
+apt update -y
+apt-get install -y erlang
+apt-get install -y rabbitmq-server
+
+systemctl enable rabbitmq-server
+systemctl start rabbitmq-server
+
+rabbitmq-plugins enable rabbitmq_management
+rabbitmqctl cluster_status
+rabbitmqctl add_user hueuser cloudera
+rabbitmqctl add_vhost huevhost
+rabbitmqctl set_user_tags hueuser administrator
+rabbitmqctl set_permissions -p huevhost hueuser ".*" ".*" ".*"
+rabbitmqctl delete_user guest

+ 17 - 0
desktop/core/src/desktop/settings.py

@@ -203,6 +203,7 @@ INSTALLED_APPS = [
     # App that keeps track of failed logins.
     'axes',
     'webpack_loader',
+    'django_celery_results',
 ]
 
 WEBPACK_LOADER = {
@@ -216,6 +217,22 @@ LOCALE_PATHS = [
   get_desktop_root('core/src/desktop/locale')
 ]
 
+# Celery
+CELERY_BROKER_URL='pyamqp://hueuser:cloudera@localhost:5672/huevhost/'
+#CELERY_BIN=
+#CELERY_RESULT_BACKEND='django-db'
+CELERY_APP="desktop"
+CELERYD_OPTS="--time-limit=300 --concurrency=8"
+
+# %n will be replaced with the first part of the nodename.
+CELERYD_LOG_FILE="/var/log/celery/%n%I.log"
+CELERYD_PID_FILE="/var/run/celery/%n.pid"
+CELERY_CREATE_DIRS=1
+CELERYD_USER="root"
+CELERYD_GROUP="root"
+
+CELERY_RESULT_BACKEND = 'django-db'
+
 # Keep default values up to date
 GTEMPLATE_CONTEXT_PROCESSORS = (
   'django.contrib.auth.context_processors.auth',

+ 41 - 0
desktop/libs/notebook/src/notebook/tasks.py

@@ -0,0 +1,41 @@
+from __future__ import absolute_import, unicode_literals
+import os
+import django
+# set the default Django settings module for the 'celery' program.
+os.environ.setdefault("DJANGO_SETTINGS_MODULE", "desktop.settings")
+django.setup()
+from django.conf import settings
+
+from celery import Celery
+app = Celery("desktop")
+app.config_from_object('django.conf:settings', namespace='CELERY')
+
+import json
+import logging
+
+from desktop import conf
+from desktop.conf import ENABLE_DOWNLOAD, USE_NEW_EDITOR
+
+from celery.utils.log import get_task_logger
+logger = get_task_logger(__name__)
+from django.http import HttpRequest
+from django.contrib.auth.models import User
+
+from notebook.connectors.base import get_api, _get_snippet_name
+from django.contrib.auth import authenticate
+
+@app.task()
+def download(postdict, notebook, snippet, file_format):
+    request = HttpRequest()
+    request.POST = postdict
+    user = authenticate(username="admin",password="admin")
+    request.user = user
+    response = get_api(request, snippet).download(notebook, snippet, file_format)
+    f=open("/tmp/foo","w")
+    for data in response.streaming_content:
+      f.write(data)
+    f.close()
+    return 0
+
+if __name__ == '__main__':
+    task = download.s(notebook, snippet, file_format).delay()

+ 23 - 1
desktop/libs/notebook/src/notebook/views.py

@@ -40,6 +40,7 @@ from notebook.decorators import check_editor_access_permission, check_document_a
 from notebook.management.commands.notebook_setup import Command
 from notebook.models import make_notebook
 
+import notebook.tasks as ntasks
 
 LOG = logging.getLogger(__name__)
 
@@ -313,7 +314,7 @@ def copy(request):
 
 
 @check_document_access_permission()
-def download(request):
+def download_new(request):
   if not ENABLE_DOWNLOAD.get():
     return serve_403_error(request)
 
@@ -332,6 +333,27 @@ def download(request):
 
   return response
 
+@check_document_access_permission()
+def download(request):
+  if not ENABLE_DOWNLOAD.get():
+    return serve_403_error(request)
+
+  notebook = json.loads(request.POST.get('notebook', '{}'))
+  snippet = json.loads(request.POST.get('snippet', '{}'))
+  file_format = request.POST.get('format', 'csv')
+
+  ntasks.download.delay(request.POST, notebook, snippet, file_format)
+  #response = get_api(request, snippet).download(notebook, snippet, file_format, user_agent=request.META.get('HTTP_USER_AGENT'))
+  response = {}
+
+  if response:
+    request.audit = {
+      'operation': 'DOWNLOAD',
+      'operationText': 'User %s downloaded results from %s as %s' % (request.user.username, _get_snippet_name(notebook), file_format),
+      'allowed': True
+    }
+
+  return response
 
 def install_examples(request):
   response = {'status': -1, 'message': ''}