conf.py 14 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445
  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. from collections import OrderedDict
  20. from django.test.client import Client
  21. from django.urls import reverse
  22. from django.utils.translation import ugettext_lazy as _t, ugettext as _
  23. from desktop import appmanager
  24. from desktop.conf import is_oozie_enabled, has_connectors, is_cm_managed
  25. from desktop.lib.conf import Config, UnspecifiedConfigSection, ConfigSection, coerce_json_dict, coerce_bool, coerce_csv
  26. LOG = logging.getLogger(__name__)
  27. # Not used when connector are on
  28. INTERPRETERS_CACHE = None
  29. SHOW_NOTEBOOKS = Config(
  30. key="show_notebooks",
  31. help=_t("Show the notebook menu or not"),
  32. type=coerce_bool,
  33. default=True
  34. )
  35. def _remove_duplications(a_list):
  36. return list(OrderedDict.fromkeys(a_list))
  37. def check_has_missing_permission(user, interpreter, user_apps=None):
  38. # TODO: port to cluster config
  39. if user_apps is None:
  40. user_apps = appmanager.get_apps_dict(user) # Expensive method
  41. return (interpreter == 'hive' and 'hive' not in user_apps) or \
  42. (interpreter == 'impala' and 'impala' not in user_apps) or \
  43. (interpreter == 'pig' and 'pig' not in user_apps) or \
  44. (interpreter == 'solr' and 'search' not in user_apps) or \
  45. (interpreter in ('spark', 'pyspark', 'r', 'jar', 'py', 'sparksql') and 'spark' not in user_apps) or \
  46. (interpreter in ('java', 'spark2', 'mapreduce', 'shell', 'sqoop1', 'distcp') and 'oozie' not in user_apps)
  47. def _connector_to_iterpreter(connector):
  48. return {
  49. 'name': connector['nice_name'],
  50. 'type': connector['name'], # Aka id
  51. 'dialect': connector['dialect'],
  52. 'category': connector['category'],
  53. 'is_sql': connector['dialect_properties']['is_sql'],
  54. 'interface': connector['interface'],
  55. 'options': {setting['name']: setting['value'] for setting in connector['settings']},
  56. 'dialect_properties': connector['dialect_properties'],
  57. }
  58. def get_ordered_interpreters(user=None):
  59. global INTERPRETERS_CACHE
  60. if has_connectors():
  61. from desktop.lib.connectors.api import _get_installed_connectors
  62. interpreters = [
  63. _connector_to_iterpreter(connector)
  64. for connector in _get_installed_connectors(categories=['editor', 'catalogs'], user=user)
  65. ]
  66. else:
  67. if INTERPRETERS_CACHE is None:
  68. none_user = None # for getting full list of interpreters
  69. if is_cm_managed():
  70. extra_interpreters = INTERPRETERS.get() # Combine the other apps interpreters
  71. _default_interpreters(none_user)
  72. else:
  73. extra_interpreters = {}
  74. if not INTERPRETERS.get():
  75. _default_interpreters(none_user)
  76. INTERPRETERS_CACHE = INTERPRETERS.get()
  77. INTERPRETERS_CACHE.update(extra_interpreters)
  78. user_apps = appmanager.get_apps_dict(user)
  79. user_interpreters = []
  80. for interpreter in INTERPRETERS_CACHE:
  81. if check_has_missing_permission(user, interpreter, user_apps=user_apps):
  82. pass # Not allowed
  83. else:
  84. user_interpreters.append(interpreter)
  85. interpreters_shown_on_wheel = _remove_duplications(INTERPRETERS_SHOWN_ON_WHEEL.get())
  86. unknown_interpreters = set(interpreters_shown_on_wheel) - set(user_interpreters)
  87. if unknown_interpreters:
  88. # Just filtering it out might be better than failing for this user
  89. raise ValueError("Interpreters from interpreters_shown_on_wheel is not in the list of Interpreters %s" % unknown_interpreters)
  90. reordered_interpreters = interpreters_shown_on_wheel + [i for i in user_interpreters if i not in interpreters_shown_on_wheel]
  91. interpreters = [{
  92. 'name': INTERPRETERS_CACHE[i].NAME.get(),
  93. 'type': i,
  94. 'interface': INTERPRETERS_CACHE[i].INTERFACE.get(),
  95. 'options': INTERPRETERS_CACHE[i].OPTIONS.get()
  96. }
  97. for i in reordered_interpreters
  98. ]
  99. return [{
  100. "name": i.get('nice_name', i['name']),
  101. "type": i['type'],
  102. "interface": i['interface'],
  103. "options": i['options'],
  104. 'dialect': i.get('dialect', i['name']).lower(),
  105. 'dialect_properties': i.get('dialect_properties') or {}, # Empty when connectors off
  106. 'category': i.get('category', 'editor'),
  107. "is_sql": i.get('is_sql') or \
  108. i['interface'] in ["hiveserver2", "rdbms", "jdbc", "solr", "sqlalchemy", "ksql", "flink"] or \
  109. i['type'] == 'sql',
  110. "is_catalog": i['interface'] in ["hms",],
  111. }
  112. for i in interpreters
  113. ]
  114. # cf. admin wizard too
  115. INTERPRETERS = UnspecifiedConfigSection(
  116. "interpreters",
  117. help="One entry for each type of snippet.",
  118. each=ConfigSection(
  119. help=_t("Define the name and how to connect and execute the language."),
  120. members=dict(
  121. NAME=Config(
  122. "name",
  123. help=_t("The name of the snippet."),
  124. default="SQL",
  125. type=str,
  126. ),
  127. INTERFACE=Config(
  128. "interface",
  129. help="The backend connection to use to communicate with the server.",
  130. default="hiveserver2",
  131. type=str,
  132. ),
  133. OPTIONS=Config(
  134. key='options',
  135. help=_t('Specific options for connecting to the server.'),
  136. type=coerce_json_dict,
  137. default='{}'
  138. )
  139. )
  140. )
  141. )
  142. INTERPRETERS_SHOWN_ON_WHEEL = Config(
  143. key="interpreters_shown_on_wheel",
  144. help=_t("Comma separated list of interpreters that should be shown on the wheel. "
  145. "This list takes precedence over the order in which the interpreter entries appear. "
  146. "Only the first 5 interpreters will appear on the wheel."),
  147. type=coerce_csv,
  148. default=[]
  149. )
  150. DEFAULT_LIMIT = Config(
  151. "default_limit",
  152. help="Default limit to use in SELECT statements if not present. Set to 0 to disable.",
  153. default=5000,
  154. type=int
  155. )
  156. ENABLE_DBPROXY_SERVER = Config(
  157. key="enable_dbproxy_server",
  158. help=_t("Main flag to override the automatic starting of the DBProxy server."),
  159. type=coerce_bool,
  160. default=True
  161. )
  162. DBPROXY_EXTRA_CLASSPATH = Config(
  163. key="dbproxy_extra_classpath",
  164. help=_t("Additional classes to put on the dbproxy classpath when starting. Values separated by ':'"),
  165. type=str,
  166. default=''
  167. )
  168. ENABLE_QUERY_BUILDER = Config(
  169. key="enable_query_builder",
  170. help=_t("Flag to enable the SQL query builder of the table assist."),
  171. type=coerce_bool,
  172. default=True
  173. )
  174. ENABLE_NOTEBOOK_2 = Config(
  175. key="enable_notebook_2",
  176. help=_t("Feature flag to enable Notebook 2."),
  177. type=coerce_bool,
  178. default=False
  179. )
  180. # Note: requires Oozie app
  181. ENABLE_QUERY_SCHEDULING = Config(
  182. key="enable_query_scheduling",
  183. help=_t("Flag to enable the creation of a coordinator for the current SQL query."),
  184. type=coerce_bool,
  185. default=False
  186. )
  187. ENABLE_EXTERNAL_STATEMENT = Config(
  188. key="enable_external_statements",
  189. help=_t("Flag to enable the selection of queries from files, saved queries into the editor or as snippet."),
  190. type=coerce_bool,
  191. default=True
  192. )
  193. ENABLE_BATCH_EXECUTE = Config(
  194. key="enable_batch_execute",
  195. help=_t("Flag to enable the bulk submission of queries as a background task through Oozie."),
  196. type=coerce_bool,
  197. dynamic_default=is_oozie_enabled
  198. )
  199. ENABLE_SQL_INDEXER = Config(
  200. key="enable_sql_indexer",
  201. help=_t("Flag to turn on the SQL indexer."),
  202. type=coerce_bool,
  203. default=False
  204. )
  205. ENABLE_PRESENTATION = Config(
  206. key="enable_presentation",
  207. help=_t("Flag to turn on the Presentation mode of the editor."),
  208. type=coerce_bool,
  209. default=True
  210. )
  211. ENABLE_QUERY_ANALYSIS = Config(
  212. key="enable_query_analysis",
  213. help=_t("Flag to turn on the built-in hints on Impala queries in the editor."),
  214. type=coerce_bool,
  215. default=False
  216. )
  217. EXAMPLES = ConfigSection(
  218. key='examples',
  219. help=_t('Define which query and table examples can be automatically setup for the available dialects.'),
  220. members=dict(
  221. AUTO_LOAD=Config(
  222. 'auto_load',
  223. help=_t('If installing the examples automatically at startup.'),
  224. type=coerce_bool,
  225. default=False
  226. ),
  227. QUERIES=Config(
  228. 'queries',
  229. help='Names of the saved queries to install. All if empty.',
  230. type=coerce_csv,
  231. default=[]
  232. ),
  233. TABLES=Config(
  234. key='tables',
  235. help=_t('Names of the tables to install. All if empty.'),
  236. type=coerce_csv,
  237. default=[]
  238. )
  239. )
  240. )
  241. def _default_interpreters(user):
  242. interpreters = []
  243. apps = appmanager.get_apps_dict(user)
  244. if 'hive' in apps:
  245. from beeswax.hive_site import get_hive_execution_engine
  246. interpreter_name = 'Impala' if get_hive_execution_engine() == 'impala' else 'Hive' # Until using a proper dialect for 'FENG'
  247. interpreters.append(('hive', {
  248. 'name': interpreter_name, 'interface': 'hiveserver2', 'options': {}
  249. }),)
  250. if 'impala' in apps:
  251. interpreters.append(('impala', {
  252. 'name': 'Impala', 'interface': 'hiveserver2', 'options': {}
  253. }),)
  254. if 'pig' in apps:
  255. interpreters.append(('pig', {
  256. 'name': 'Pig', 'interface': 'oozie', 'options': {}
  257. }))
  258. if 'oozie' in apps and 'jobsub' in apps:
  259. interpreters.extend((
  260. ('java', {
  261. 'name': 'Java', 'interface': 'oozie', 'options': {}
  262. }),
  263. ('spark2', {
  264. 'name': 'Spark', 'interface': 'oozie', 'options': {}
  265. }),
  266. ('mapreduce', {
  267. 'name': 'MapReduce', 'interface': 'oozie', 'options': {}
  268. }),
  269. ('shell', {
  270. 'name': 'Shell', 'interface': 'oozie', 'options': {}
  271. }),
  272. ('sqoop1', {
  273. 'name': 'Sqoop 1', 'interface': 'oozie', 'options': {}
  274. }),
  275. ('distcp', {
  276. 'name': 'Distcp', 'interface': 'oozie', 'options': {}
  277. }),
  278. ))
  279. from dashboard.conf import get_properties # Cyclic dependency
  280. dashboards = get_properties()
  281. if dashboards.get('solr') and dashboards['solr']['analytics']:
  282. interpreters.append(('solr', {
  283. 'name': 'Solr SQL', 'interface': 'solr', 'options': {}
  284. }),)
  285. from desktop.models import Cluster # Cyclic dependency
  286. cluster = Cluster(user)
  287. if cluster and cluster.get_type() == 'dataeng':
  288. interpreters.append(('dataeng', {
  289. 'name': 'DataEng', 'interface': 'dataeng', 'options': {}
  290. }))
  291. if 'spark' in apps:
  292. interpreters.extend((
  293. ('spark', {
  294. 'name': 'Scala', 'interface': 'livy', 'options': {}
  295. }),
  296. ('pyspark', {
  297. 'name': 'PySpark', 'interface': 'livy', 'options': {}
  298. }),
  299. ('r', {
  300. 'name': 'R', 'interface': 'livy', 'options': {}
  301. }),
  302. ('jar', {
  303. 'name': 'Spark Submit Jar', 'interface': 'livy-batch', 'options': {}
  304. }),
  305. ('py', {
  306. 'name': 'Spark Submit Python', 'interface': 'livy-batch', 'options': {}
  307. }),
  308. ('text', {
  309. 'name': 'Text', 'interface': 'text', 'options': {}
  310. }),
  311. ('markdown', {
  312. 'name': 'Markdown', 'interface': 'text', 'options': {}
  313. })
  314. ))
  315. INTERPRETERS.set_for_testing(OrderedDict(interpreters))
  316. def config_validator(user, interpreters=None):
  317. res = []
  318. if not has_connectors():
  319. return res
  320. client = Client()
  321. client.force_login(user=user)
  322. if not user.is_authenticated:
  323. res.append(('Editor', _('Could not authenticate with user %s to validate interpreters') % user))
  324. if interpreters is None:
  325. interpreters = get_ordered_interpreters(user=user)
  326. for interpreter in interpreters:
  327. if interpreter.get('is_sql'):
  328. connector_id = interpreter['type']
  329. try:
  330. response = _excute_test_query(client, connector_id, interpreter=interpreter)
  331. data = json.loads(response.content)
  332. if data['status'] != 0:
  333. raise Exception(data)
  334. except Exception as e:
  335. trace = str(e)
  336. msg = "Testing the connector connection failed."
  337. if 'Error validating the login' in trace or 'TSocket read 0 bytes' in trace:
  338. msg += ' Failed to authenticate, check authentication configurations.'
  339. LOG.exception(msg)
  340. res.append(
  341. (
  342. '%(name)s - %(dialect)s (%(type)s)' % interpreter,
  343. _(msg) + (' %s' % trace[:100] + ('...' if len(trace) > 50 else ''))
  344. )
  345. )
  346. return res
  347. def _excute_test_query(client, connector_id, interpreter=None):
  348. '''
  349. Helper utils until the API gets simplified.
  350. '''
  351. notebook_json = """
  352. {
  353. "selectedSnippet": "hive",
  354. "showHistory": false,
  355. "description": "Test Query",
  356. "name": "Test Query",
  357. "sessions": [
  358. {
  359. "type": "hive",
  360. "properties": [],
  361. "id": null
  362. }
  363. ],
  364. "type": "hive",
  365. "id": null,
  366. "snippets": [{"id":"2b7d1f46-17a0-30af-efeb-33d4c29b1055","type":"%(connector_id)s","status":"running","statement":"select * from web_logs","properties":{"settings":[],"variables":[],"files":[],"functions":[]},"result":{"id":"b424befa-f4f5-8799-a0b4-79753f2552b1","type":"table","handle":{"log_context":null,"statements_count":1,"end":{"column":21,"row":0},"statement_id":0,"has_more_statements":false,"start":{"column":0,"row":0},"secret":"rVRWw7YPRGqPT7LZ/TeFaA==an","has_result_set":true,"statement":"select * from web_logs","operation_type":0,"modified_row_count":null,"guid":"7xm6+epkRx6dyvYvGNYePA==an"}},"lastExecuted": 1462554843817,"database":"default"}],
  367. "uuid": "d9efdee1-ef25-4d43-b8f9-1a170f69a05a"
  368. }
  369. """ % {
  370. 'connector_id': connector_id,
  371. }
  372. snippet = json.loads(notebook_json)['snippets'][0]
  373. snippet['interpreter'] = interpreter
  374. return client.post(
  375. reverse('notebook:api_sample_data', kwargs={'database': 'default', 'table': 'default'}), {
  376. 'notebook': notebook_json,
  377. 'snippet': json.dumps(snippet),
  378. 'is_async': json.dumps(True),
  379. 'operation': json.dumps('hello')
  380. })