api.py 27 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817
  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. import json
  18. import logging
  19. import sqlparse
  20. import urllib
  21. from django.urls import reverse
  22. from django.db.models import Q
  23. from django.utils.translation import ugettext as _
  24. from django.views.decorators.http import require_GET, require_POST
  25. from desktop.api2 import __paginate
  26. from desktop.lib.i18n import smart_str
  27. from desktop.lib.django_util import JsonResponse
  28. from desktop.models import Document2, Document
  29. from indexer.file_format import HiveFormat
  30. from indexer.fields import Field
  31. from notebook.connectors.base import get_api, Notebook, QueryExpired, SessionExpired, QueryError, _get_snippet_name
  32. from notebook.connectors.dataeng import DataEngApi
  33. from notebook.connectors.hiveserver2 import HS2Api
  34. from notebook.connectors.oozie_batch import OozieApi
  35. from notebook.decorators import api_error_handler, check_document_access_permission, check_document_modify_permission
  36. from notebook.models import escape_rows, make_notebook
  37. from notebook.views import upgrade_session_properties
  38. LOG = logging.getLogger(__name__)
  39. DEFAULT_HISTORY_NAME = ''
  40. @require_POST
  41. @api_error_handler
  42. def create_notebook(request):
  43. response = {'status': -1}
  44. editor_type = request.POST.get('type', 'notebook')
  45. directory_uuid = request.POST.get('directory_uuid')
  46. editor = Notebook()
  47. data = editor.get_data()
  48. if editor_type != 'notebook':
  49. data['name'] = ''
  50. data['type'] = 'query-%s' % editor_type # TODO: Add handling for non-SQL types
  51. data['directoryUuid'] = directory_uuid
  52. editor.data = json.dumps(data)
  53. response['notebook'] = editor.get_data()
  54. response['status'] = 0
  55. return JsonResponse(response)
  56. @require_POST
  57. @check_document_access_permission()
  58. @api_error_handler
  59. def create_session(request):
  60. response = {'status': -1}
  61. notebook = json.loads(request.POST.get('notebook', '{}'))
  62. session = json.loads(request.POST.get('session', '{}'))
  63. properties = session.get('properties', [])
  64. response['session'] = get_api(request, session).create_session(lang=session['type'], properties=properties)
  65. response['status'] = 0
  66. return JsonResponse(response)
  67. @require_POST
  68. @check_document_access_permission()
  69. @api_error_handler
  70. def close_session(request):
  71. response = {'status': -1}
  72. session = json.loads(request.POST.get('session', '{}'))
  73. response['session'] = get_api(request, {'type': session['type']}).close_session(session=session)
  74. response['status'] = 0
  75. return JsonResponse(response)
  76. def _execute_notebook(request, notebook, snippet):
  77. response = {'status': -1}
  78. result = None
  79. history = None
  80. historify = (notebook['type'] != 'notebook' or snippet.get('wasBatchExecuted')) and not notebook.get('skipHistorify')
  81. try:
  82. try:
  83. if historify:
  84. history = _historify(notebook, request.user)
  85. notebook = Notebook(document=history).get_data()
  86. response['handle'] = get_api(request, snippet).execute(notebook, snippet)
  87. # Retrieve and remove the result from the handle
  88. if response['handle'].get('sync'):
  89. result = response['handle'].pop('result')
  90. finally:
  91. if historify:
  92. _snippet = [s for s in notebook['snippets'] if s['id'] == snippet['id']][0]
  93. if 'handle' in response: # No failure
  94. _snippet['result']['handle'] = response['handle']
  95. _snippet['result']['statements_count'] = response['handle'].get('statements_count', 1)
  96. _snippet['result']['statement_id'] = response['handle'].get('statement_id', 0)
  97. _snippet['result']['handle']['statement'] = response['handle'].get('statement', snippet['statement']).strip() # For non HS2, as non multi query yet
  98. else:
  99. _snippet['status'] = 'failed'
  100. if history: # If _historify failed, history will be None
  101. history.update_data(notebook)
  102. history.save()
  103. response['history_id'] = history.id
  104. response['history_uuid'] = history.uuid
  105. if notebook['isSaved']: # Keep track of history of saved queries
  106. response['history_parent_uuid'] = history.dependencies.filter(type__startswith='query-').latest('last_modified').uuid
  107. except QueryError, ex: # We inject the history information from _historify() to the failed queries
  108. if response.get('history_id'):
  109. ex.extra['history_id'] = response['history_id']
  110. if response.get('history_uuid'):
  111. ex.extra['history_uuid'] = response['history_uuid']
  112. if response.get('history_parent_uuid'):
  113. ex.extra['history_parent_uuid'] = response['history_parent_uuid']
  114. raise ex
  115. # Inject and HTML escape results
  116. if result is not None:
  117. response['result'] = result
  118. response['result']['data'] = escape_rows(result['data'])
  119. response['status'] = 0
  120. return response
  121. @require_POST
  122. @check_document_access_permission()
  123. @api_error_handler
  124. def execute(request, engine=None):
  125. notebook = json.loads(request.POST.get('notebook', '{}'))
  126. snippet = json.loads(request.POST.get('snippet', '{}'))
  127. response = _execute_notebook(request, notebook, snippet)
  128. print response
  129. return JsonResponse(response)
  130. @require_POST
  131. @check_document_access_permission()
  132. @api_error_handler
  133. def check_status(request):
  134. response = {'status': -1}
  135. notebook = json.loads(request.POST.get('notebook', '{}'))
  136. snippet = json.loads(request.POST.get('snippet', '{}'))
  137. if not snippet:
  138. nb_doc = Document2.objects.get_by_uuid(user=request.user, uuid=notebook['id'])
  139. notebook = Notebook(document=nb_doc).get_data()
  140. snippet = notebook['snippets'][0]
  141. try:
  142. response['query_status'] = get_api(request, snippet).check_status(notebook, snippet)
  143. response['status'] = 0
  144. except SessionExpired:
  145. response['status'] = 'expired'
  146. raise
  147. except QueryExpired:
  148. response['status'] = 'expired'
  149. raise
  150. finally:
  151. if response['status'] == 0 and snippet['status'] != response['query_status']:
  152. status = response['query_status']['status']
  153. elif response['status'] == 'expired':
  154. status = 'expired'
  155. else:
  156. status = 'failed'
  157. if notebook['type'].startswith('query') or notebook.get('isManaged'):
  158. nb_doc = Document2.objects.get(id=notebook['id'])
  159. if nb_doc.can_write(request.user):
  160. nb = Notebook(document=nb_doc).get_data()
  161. if status != nb['snippets'][0]['status']:
  162. nb['snippets'][0]['status'] = status
  163. nb_doc.update_data(nb)
  164. nb_doc.save()
  165. return JsonResponse(response)
  166. @require_POST
  167. @check_document_access_permission()
  168. @api_error_handler
  169. def fetch_result_data(request):
  170. response = {'status': -1}
  171. notebook = json.loads(request.POST.get('notebook', '{}'))
  172. snippet = json.loads(request.POST.get('snippet', '{}'))
  173. rows = json.loads(request.POST.get('rows', '100'))
  174. start_over = json.loads(request.POST.get('startOver', 'false'))
  175. response['result'] = get_api(request, snippet).fetch_result(notebook, snippet, rows, start_over)
  176. # Materialize and HTML escape results
  177. if response['result'].get('data') and response['result'].get('type') == 'table':
  178. response['result']['data'] = escape_rows(response['result']['data'])
  179. response['status'] = 0
  180. return JsonResponse(response)
  181. @require_POST
  182. @check_document_access_permission()
  183. @api_error_handler
  184. def fetch_result_metadata(request):
  185. response = {'status': -1}
  186. notebook = json.loads(request.POST.get('notebook', '{}'))
  187. snippet = json.loads(request.POST.get('snippet', '{}'))
  188. response['result'] = get_api(request, snippet).fetch_result_metadata(notebook, snippet)
  189. response['status'] = 0
  190. return JsonResponse(response)
  191. @require_POST
  192. @check_document_access_permission()
  193. @api_error_handler
  194. def fetch_result_size(request):
  195. response = {'status': -1}
  196. notebook = json.loads(request.POST.get('notebook', '{}'))
  197. snippet = json.loads(request.POST.get('snippet', '{}'))
  198. response['result'] = get_api(request, snippet).fetch_result_size(notebook, snippet)
  199. response['status'] = 0
  200. return JsonResponse(response)
  201. @require_POST
  202. @check_document_access_permission()
  203. @api_error_handler
  204. def cancel_statement(request):
  205. response = {'status': -1}
  206. notebook = json.loads(request.POST.get('notebook', '{}'))
  207. snippet = json.loads(request.POST.get('snippet', '{}'))
  208. response['result'] = get_api(request, snippet).cancel(notebook, snippet)
  209. response['status'] = 0
  210. return JsonResponse(response)
  211. @require_POST
  212. @check_document_access_permission()
  213. @api_error_handler
  214. def get_logs(request):
  215. response = {'status': -1}
  216. notebook = json.loads(request.POST.get('notebook', '{}'))
  217. snippet = json.loads(request.POST.get('snippet', '{}'))
  218. startFrom = request.POST.get('from')
  219. startFrom = int(startFrom) if startFrom else None
  220. size = request.POST.get('size')
  221. size = int(size) if size else None
  222. db = get_api(request, snippet)
  223. full_log = smart_str(request.POST.get('full_log', ''))
  224. logs = smart_str(db.get_log(notebook, snippet, startFrom=startFrom, size=size))
  225. full_log += logs
  226. jobs = db.get_jobs(notebook, snippet, full_log)
  227. response['logs'] = logs.strip()
  228. response['progress'] = min(db.progress(snippet, full_log), 99) if snippet['status'] != 'available' and snippet['status'] != 'success' else 100
  229. response['jobs'] = jobs
  230. response['isFullLogs'] = isinstance(db, (OozieApi, DataEngApi))
  231. response['status'] = 0
  232. return JsonResponse(response)
  233. def _save_notebook(notebook, user):
  234. notebook_type = notebook.get('type', 'notebook')
  235. save_as = False
  236. if notebook.get('parentSavedQueryUuid'): # We save into the original saved query, not into the query history
  237. notebook_doc = Document2.objects.get_by_uuid(user=user, uuid=notebook['parentSavedQueryUuid'])
  238. elif notebook.get('id'):
  239. notebook_doc = Document2.objects.get(id=notebook['id'])
  240. else:
  241. notebook_doc = Document2.objects.create(name=notebook['name'], uuid=notebook['uuid'], type=notebook_type, owner=user)
  242. Document.objects.link(notebook_doc, owner=notebook_doc.owner, name=notebook_doc.name, description=notebook_doc.description, extra=notebook_type)
  243. save_as = True
  244. if notebook.get('directoryUuid'):
  245. notebook_doc.parent_directory = Document2.objects.get_by_uuid(user=user, uuid=notebook.get('directoryUuid'), perm_type='write')
  246. else:
  247. notebook_doc.parent_directory = Document2.objects.get_home_directory(user)
  248. notebook['isSaved'] = True
  249. notebook['isHistory'] = False
  250. notebook['id'] = notebook_doc.id
  251. _clear_sessions(notebook)
  252. notebook_doc1 = notebook_doc._get_doc1(doc2_type=notebook_type)
  253. notebook_doc.update_data(notebook)
  254. notebook_doc.search = _get_statement(notebook)
  255. notebook_doc.name = notebook_doc1.name = notebook['name']
  256. notebook_doc.description = notebook_doc1.description = notebook['description']
  257. notebook_doc.save()
  258. notebook_doc1.save()
  259. return notebook_doc, save_as
  260. @api_error_handler
  261. @require_POST
  262. @check_document_modify_permission()
  263. def save_notebook(request):
  264. response = {'status': -1}
  265. notebook = json.loads(request.POST.get('notebook', '{}'))
  266. notebook_doc, save_as = _save_notebook(notebook, request.user)
  267. response['status'] = 0
  268. response['save_as'] = save_as
  269. response.update(notebook_doc.to_dict())
  270. response['message'] = request.POST.get('editorMode') == 'true' and _('Query saved successfully') or _('Notebook saved successfully')
  271. return JsonResponse(response)
  272. def _clear_sessions(notebook):
  273. notebook['sessions'] = [_s for _s in notebook['sessions'] if _s['type'] in ('scala', 'spark', 'pyspark', 'sparkr')]
  274. def _historify(notebook, user):
  275. query_type = notebook['type']
  276. name = notebook['name'] if (notebook['name'] and notebook['name'].strip() != '') else DEFAULT_HISTORY_NAME
  277. is_managed = notebook.get('isManaged') == True # Prevents None
  278. if is_managed and Document2.objects.filter(uuid=notebook['uuid']).exists():
  279. history_doc = Document2.objects.get(uuid=notebook['uuid'])
  280. else:
  281. history_doc = Document2.objects.create(
  282. name=name,
  283. type=query_type,
  284. owner=user,
  285. is_history=True,
  286. is_managed=is_managed
  287. )
  288. # Link history of saved query
  289. if notebook['isSaved']:
  290. parent_doc = Document2.objects.get(uuid=notebook.get('parentSavedQueryUuid') or notebook['uuid']) # From previous history query or initial saved query
  291. notebook['parentSavedQueryUuid'] = parent_doc.uuid
  292. history_doc.dependencies.add(parent_doc)
  293. if not is_managed:
  294. Document.objects.link(
  295. history_doc,
  296. name=history_doc.name,
  297. owner=history_doc.owner,
  298. description=history_doc.description,
  299. extra=query_type
  300. )
  301. notebook['uuid'] = history_doc.uuid
  302. _clear_sessions(notebook)
  303. history_doc.update_data(notebook)
  304. history_doc.search = _get_statement(notebook)
  305. history_doc.save()
  306. return history_doc
  307. def _get_statement(notebook):
  308. if notebook['snippets'] and len(notebook['snippets']) > 0:
  309. return Notebook.statement_with_variables(notebook['snippets'][0])
  310. return ''
  311. @require_GET
  312. @api_error_handler
  313. @check_document_access_permission()
  314. def get_history(request):
  315. response = {'status': -1}
  316. doc_type = request.GET.get('doc_type')
  317. doc_text = request.GET.get('doc_text')
  318. page = min(int(request.GET.get('page', 1)), 100)
  319. limit = min(int(request.GET.get('limit', 50)), 100)
  320. is_notification_manager = request.GET.get('is_notification_manager', 'false') == 'true'
  321. if is_notification_manager:
  322. docs = Document2.objects.get_tasks_history(user=request.user)
  323. else:
  324. docs = Document2.objects.get_history(doc_type='query-%s' % doc_type, user=request.user)
  325. if doc_text:
  326. docs = docs.filter(Q(name__icontains=doc_text) | Q(description__icontains=doc_text) | Q(search__icontains=doc_text))
  327. # Paginate
  328. docs = docs.order_by('-last_modified')
  329. response['count'] = docs.count()
  330. docs = __paginate(page, limit, queryset=docs)['documents']
  331. history = []
  332. for doc in docs:
  333. notebook = Notebook(document=doc).get_data()
  334. if 'snippets' in notebook:
  335. statement = notebook['description'] if is_notification_manager else _get_statement(notebook)
  336. history.append({
  337. 'name': doc.name,
  338. 'id': doc.id,
  339. 'uuid': doc.uuid,
  340. 'type': doc.type,
  341. 'data': {
  342. 'statement': statement[:1001] if statement else '',
  343. 'lastExecuted': notebook['snippets'][0].get('lastExecuted', -1),
  344. 'status': notebook['snippets'][0]['status'],
  345. 'parentSavedQueryUuid': notebook.get('parentSavedQueryUuid', '')
  346. } if notebook['snippets'] else {},
  347. 'absoluteUrl': doc.get_absolute_url(),
  348. })
  349. else:
  350. LOG.error('Incomplete History Notebook: %s' % notebook)
  351. response['history'] = sorted(history, key=lambda row: row['data']['lastExecuted'], reverse=True)
  352. response['message'] = _('History fetched')
  353. response['status'] = 0
  354. return JsonResponse(response)
  355. @require_POST
  356. @api_error_handler
  357. @check_document_modify_permission()
  358. def clear_history(request):
  359. response = {'status': -1}
  360. notebook = json.loads(request.POST.get('notebook'), '{}')
  361. doc_type = request.POST.get('doc_type')
  362. is_notification_manager = request.POST.get('is_notification_manager', 'false') == 'true'
  363. if is_notification_manager:
  364. history = Document2.objects.get_tasks_history(user=request.user)
  365. else:
  366. history = Document2.objects.get_history(doc_type='query-%s' % doc_type, user=request.user)
  367. response['updated'] = history.delete()
  368. response['message'] = _('History cleared !')
  369. response['status'] = 0
  370. return JsonResponse(response)
  371. @require_GET
  372. @check_document_access_permission()
  373. def open_notebook(request):
  374. response = {'status': -1}
  375. notebook_id = request.GET.get('notebook')
  376. notebook = Notebook(document=Document2.objects.get(id=notebook_id))
  377. notebook = upgrade_session_properties(request, notebook)
  378. response['status'] = 0
  379. response['notebook'] = notebook.get_json()
  380. response['message'] = _('Notebook loaded successfully')
  381. @require_POST
  382. @check_document_access_permission()
  383. def close_notebook(request):
  384. response = {'status': -1, 'result': []}
  385. notebook = json.loads(request.POST.get('notebook', '{}'))
  386. for session in [_s for _s in notebook['sessions'] if _s['type'] in ('scala', 'spark', 'pyspark', 'sparkr')]:
  387. try:
  388. response['result'].append(get_api(request, session).close_session(session))
  389. except QueryExpired:
  390. pass
  391. except Exception, e:
  392. LOG.exception('Error closing session %s' % str(e))
  393. for snippet in [_s for _s in notebook['snippets'] if _s['type'] in ('hive', 'impala')]:
  394. try:
  395. if snippet['status'] != 'running':
  396. response['result'].append(get_api(request, snippet).close_statement(snippet))
  397. else:
  398. LOG.info('Not closing SQL snippet as still running.')
  399. except QueryExpired:
  400. pass
  401. except Exception, e:
  402. LOG.exception('Error closing statement %s' % str(e))
  403. response['status'] = 0
  404. response['message'] = _('Notebook closed successfully')
  405. return JsonResponse(response)
  406. @require_POST
  407. @check_document_access_permission()
  408. def close_statement(request):
  409. response = {'status': -1}
  410. # Passed by check_document_access_permission but unused by APIs
  411. notebook = json.loads(request.POST.get('notebook', '{}'))
  412. snippet = json.loads(request.POST.get('snippet', '{}'))
  413. try:
  414. response['result'] = get_api(request, snippet).close_statement(snippet)
  415. except QueryExpired:
  416. pass
  417. response['status'] = 0
  418. response['message'] = _('Statement closed !')
  419. return JsonResponse(response)
  420. @require_POST
  421. @check_document_access_permission()
  422. @api_error_handler
  423. def autocomplete(request, server=None, database=None, table=None, column=None, nested=None):
  424. response = {'status': -1}
  425. # Passed by check_document_access_permission but unused by APIs
  426. notebook = json.loads(request.POST.get('notebook', '{}'))
  427. snippet = json.loads(request.POST.get('snippet', '{}'))
  428. try:
  429. autocomplete_data = get_api(request, snippet).autocomplete(snippet, database, table, column, nested)
  430. response.update(autocomplete_data)
  431. except QueryExpired:
  432. pass
  433. response['status'] = 0
  434. return JsonResponse(response)
  435. @require_POST
  436. @check_document_access_permission()
  437. @api_error_handler
  438. def get_sample_data(request, server=None, database=None, table=None, column=None):
  439. response = {'status': -1}
  440. # Passed by check_document_access_permission but unused by APIs
  441. notebook = json.loads(request.POST.get('notebook', '{}'))
  442. snippet = json.loads(request.POST.get('snippet', '{}'))
  443. async = json.loads(request.POST.get('async', 'false'))
  444. sample_data = get_api(request, snippet).get_sample_data(snippet, database, table, column, async=async)
  445. response.update(sample_data)
  446. response['status'] = 0
  447. return JsonResponse(response)
  448. @require_POST
  449. @check_document_access_permission()
  450. @api_error_handler
  451. def explain(request):
  452. response = {'status': -1}
  453. notebook = json.loads(request.POST.get('notebook', '{}'))
  454. snippet = json.loads(request.POST.get('snippet', '{}'))
  455. response = get_api(request, snippet).explain(notebook, snippet)
  456. return JsonResponse(response)
  457. @require_POST
  458. @api_error_handler
  459. def format(request):
  460. response = {'status': 0}
  461. statements = request.POST.get('statements', '')
  462. response['formatted_statements'] = sqlparse.format(statements, reindent=True, keyword_case='upper') # SQL only currently
  463. return JsonResponse(response)
  464. @require_POST
  465. @check_document_access_permission()
  466. @api_error_handler
  467. def export_result(request):
  468. response = {'status': -1, 'message': _('Success')}
  469. # Passed by check_document_access_permission but unused by APIs
  470. notebook = json.loads(request.POST.get('notebook', '{}'))
  471. snippet = json.loads(request.POST.get('snippet', '{}'))
  472. data_format = json.loads(request.POST.get('format', 'hdfs-file'))
  473. destination = urllib.unquote(json.loads(request.POST.get('destination', '')))
  474. overwrite = json.loads(request.POST.get('overwrite', 'false'))
  475. is_embedded = json.loads(request.POST.get('is_embedded', 'false'))
  476. start_time = json.loads(request.POST.get('start_time', '-1'))
  477. api = get_api(request, snippet)
  478. if data_format == 'hdfs-file': # Blocking operation, like downloading
  479. if request.fs.isdir(destination):
  480. if notebook.get('name'):
  481. destination += '/%(name)s.csv' % notebook
  482. else:
  483. destination += '/%(type)s-%(id)s.csv' % notebook
  484. if overwrite and request.fs.exists(destination):
  485. request.fs.do_as_user(request.user.username, request.fs.rmtree, destination)
  486. response['watch_url'] = api.export_data_as_hdfs_file(snippet, destination, overwrite)
  487. response['status'] = 0
  488. request.audit = {
  489. 'operation': 'EXPORT',
  490. 'operationText': 'User %s exported to HDFS destination: %s' % (request.user.username, destination),
  491. 'allowed': True
  492. }
  493. elif data_format == 'hive-table':
  494. if is_embedded:
  495. sql, success_url = api.export_data_as_table(notebook, snippet, destination)
  496. task = make_notebook(
  497. name=_('Export %s query to table %s') % (snippet['type'], destination),
  498. description=_('Query %s to %s') % (_get_snippet_name(notebook), success_url),
  499. editor_type=snippet['type'],
  500. statement=sql,
  501. status='ready',
  502. database=snippet['database'],
  503. on_success_url=success_url,
  504. last_executed=start_time,
  505. is_task=True
  506. )
  507. response = task.execute(request)
  508. else:
  509. notebook_id = notebook['id'] or request.GET.get('editor', request.GET.get('notebook'))
  510. response['watch_url'] = reverse('notebook:execute_and_watch') + '?action=save_as_table&notebook=' + str(notebook_id) + '&snippet=0&destination=' + destination
  511. response['status'] = 0
  512. request.audit = {
  513. 'operation': 'EXPORT',
  514. 'operationText': 'User %s exported to Hive table: %s' % (request.user.username, destination),
  515. 'allowed': True
  516. }
  517. elif data_format == 'hdfs-directory':
  518. if is_embedded:
  519. sql, success_url = api.export_large_data_to_hdfs(notebook, snippet, destination)
  520. task = make_notebook(
  521. name=_('Export %s query to directory') % snippet['type'],
  522. description=_('Query %s to %s') % (_get_snippet_name(notebook), success_url),
  523. editor_type=snippet['type'],
  524. statement=sql,
  525. status='ready-execute',
  526. database=snippet['database'],
  527. on_success_url=success_url,
  528. last_executed=start_time,
  529. is_task=True
  530. )
  531. response = task.execute(request)
  532. else:
  533. notebook_id = notebook['id'] or request.GET.get('editor', request.GET.get('notebook'))
  534. response['watch_url'] = reverse('notebook:execute_and_watch') + '?action=insert_as_query&notebook=' + str(notebook_id) + '&snippet=0&destination=' + destination
  535. response['status'] = 0
  536. request.audit = {
  537. 'operation': 'EXPORT',
  538. 'operationText': 'User %s exported to HDFS directory: %s' % (request.user.username, destination),
  539. 'allowed': True
  540. }
  541. elif data_format in ('search-index', 'dashboard'):
  542. # Open the result in the Dashboard via a SQL sub-query or the Import wizard (quick vs scalable)
  543. if is_embedded:
  544. notebook_id = notebook['id'] or request.GET.get('editor', request.GET.get('notebook'))
  545. if data_format == 'dashboard':
  546. engine = notebook['type'].replace('query-', '')
  547. response['watch_url'] = reverse('dashboard:browse', kwargs={'name': notebook_id}) + '?source=query&engine=%(engine)s' % {'engine': engine}
  548. response['status'] = 0
  549. else:
  550. sample = get_api(request, snippet).fetch_result(notebook, snippet, rows=4, start_over=True)
  551. for col in sample['meta']:
  552. col['type'] = HiveFormat.FIELD_TYPE_TRANSLATE.get(col['type'], 'string')
  553. response['status'] = 0
  554. response['id'] = notebook_id
  555. response['name'] = _get_snippet_name(notebook)
  556. response['source_type'] = 'query'
  557. response['target_type'] = 'index'
  558. response['target_path'] = destination
  559. response['sample'] = list(sample['data'])
  560. response['columns'] = [
  561. Field(col['name'], col['type']).to_dict() for col in sample['meta']
  562. ]
  563. else:
  564. notebook_id = notebook['id'] or request.GET.get('editor', request.GET.get('notebook'))
  565. response['watch_url'] = reverse('notebook:execute_and_watch') + '?action=index_query&notebook=' + str(notebook_id) + '&snippet=0&destination=' + destination
  566. response['status'] = 0
  567. if response.get('status') != 0:
  568. response['message'] = _('Exporting result failed.')
  569. return JsonResponse(response)
  570. @require_POST
  571. @check_document_access_permission()
  572. @api_error_handler
  573. def statement_risk(request):
  574. response = {'status': -1, 'message': ''}
  575. notebook = json.loads(request.POST.get('notebook', '{}'))
  576. snippet = json.loads(request.POST.get('snippet', '{}'))
  577. api = HS2Api(request.user, snippet)
  578. response['query_complexity'] = api.statement_risk(notebook, snippet)
  579. response['status'] = 0
  580. return JsonResponse(response)
  581. @require_POST
  582. @check_document_access_permission()
  583. @api_error_handler
  584. def statement_compatibility(request):
  585. response = {'status': -1, 'message': ''}
  586. notebook = json.loads(request.POST.get('notebook', '{}'))
  587. snippet = json.loads(request.POST.get('snippet', '{}'))
  588. source_platform = request.POST.get('sourcePlatform')
  589. target_platform = request.POST.get('targetPlatform')
  590. api = get_api(request, snippet)
  591. response['query_compatibility'] = api.statement_compatibility(notebook, snippet, source_platform=source_platform, target_platform=target_platform)
  592. response['status'] = 0
  593. return JsonResponse(response)
  594. @require_POST
  595. @check_document_access_permission()
  596. @api_error_handler
  597. def statement_similarity(request):
  598. response = {'status': -1, 'message': ''}
  599. notebook = json.loads(request.POST.get('notebook', '{}'))
  600. snippet = json.loads(request.POST.get('snippet', '{}'))
  601. source_platform = request.POST.get('sourcePlatform')
  602. api = get_api(request, snippet)
  603. response['statement_similarity'] = api.statement_similarity(notebook, snippet, source_platform=source_platform)
  604. response['status'] = 0
  605. return JsonResponse(response)
  606. @require_POST
  607. @check_document_access_permission()
  608. @api_error_handler
  609. def get_external_statement(request):
  610. response = {'status': -1, 'message': ''}
  611. notebook = json.loads(request.POST.get('notebook', '{}'))
  612. snippet = json.loads(request.POST.get('snippet', '{}'))
  613. if snippet.get('statementType') == 'file':
  614. response['statement'] = _get_statement_from_file(request.user, request.fs, snippet)
  615. elif snippet.get('statementType') == 'document':
  616. notebook = Notebook(Document2.objects.get_by_uuid(user=request.user, uuid=snippet['associatedDocumentUuid'], perm_type='read'))
  617. response['statement'] = notebook.get_str()
  618. response['status'] = 0
  619. return JsonResponse(response)
  620. def _get_statement_from_file(user, fs, snippet):
  621. script_path = snippet['statementPath']
  622. if script_path:
  623. script_path = script_path.replace('hdfs://', '')
  624. if fs.do_as_user(user, fs.isfile, script_path):
  625. return fs.do_as_user(user, fs.read, script_path, 0, 16 * 1024 ** 2)