|
@@ -79,6 +79,7 @@ class FlinkSqlApi(Api):
|
|
|
|
|
|
|
|
@query_error_handler
|
|
@query_error_handler
|
|
|
def create_session(self, lang=None, properties=None):
|
|
def create_session(self, lang=None, properties=None):
|
|
|
|
|
+ LOG.info("Creating session for %s.", lang)
|
|
|
session = self._get_session()
|
|
session = self._get_session()
|
|
|
|
|
|
|
|
response = {
|
|
response = {
|
|
@@ -162,26 +163,30 @@ class FlinkSqlApi(Api):
|
|
|
session = self._get_session_info_from_user()
|
|
session = self._get_session_info_from_user()
|
|
|
|
|
|
|
|
if not session:
|
|
if not session:
|
|
|
- session = self.db.create_session()
|
|
|
|
|
- if self.default_database:
|
|
|
|
|
- self._use_database(self.default_catalog, self.default_database)
|
|
|
|
|
- elif self.default_catalog:
|
|
|
|
|
- self._use_catalog(self.default_catalog)
|
|
|
|
|
-
|
|
|
|
|
|
|
+ session = self._create_session()
|
|
|
try:
|
|
try:
|
|
|
- self.db.session_heartbeat(session_handle=session['sessionHandle'])
|
|
|
|
|
|
|
+ self.db.session_heartbeat(session_handle=session['id'])
|
|
|
except Exception as e:
|
|
except Exception as e:
|
|
|
if "Session '%(sessionHandle)s' does not exist" % session in str(e):
|
|
if "Session '%(sessionHandle)s' does not exist" % session in str(e):
|
|
|
LOG.warning('Session %(sessionHandle)s does not exist, opening a new one' % session)
|
|
LOG.warning('Session %(sessionHandle)s does not exist, opening a new one' % session)
|
|
|
- session = self.db.create_session()
|
|
|
|
|
|
|
+ session = self._create_session()
|
|
|
else:
|
|
else:
|
|
|
raise e
|
|
raise e
|
|
|
|
|
|
|
|
- session['id'] = session['sessionHandle']
|
|
|
|
|
self._set_session_info_to_user(session)
|
|
self._set_session_info_to_user(session)
|
|
|
|
|
|
|
|
return session
|
|
return session
|
|
|
|
|
|
|
|
|
|
+ def _create_session(self):
|
|
|
|
|
+ session = self.db.create_session()
|
|
|
|
|
+ session['id'] = session['sessionHandle']
|
|
|
|
|
+
|
|
|
|
|
+ if self.default_database:
|
|
|
|
|
+ self._use_database(session, self.default_catalog, self.default_database)
|
|
|
|
|
+ elif self.default_catalog:
|
|
|
|
|
+ self._use_catalog(session, self.default_catalog)
|
|
|
|
|
+ return session
|
|
|
|
|
+
|
|
|
@query_error_handler
|
|
@query_error_handler
|
|
|
def execute(self, notebook, snippet):
|
|
def execute(self, notebook, snippet):
|
|
|
session = self._get_session()
|
|
session = self._get_session()
|
|
@@ -461,19 +466,15 @@ class FlinkSqlApi(Api):
|
|
|
|
|
|
|
|
return [{'name': function[0]} for function in function_list]
|
|
return [{'name': function[0]} for function in function_list]
|
|
|
|
|
|
|
|
- def _use_catalog(self, catalog):
|
|
|
|
|
- session = self._get_session()
|
|
|
|
|
- self.db.configure_session(session_handle=(session['id']), statement="USE CATALOG `%s`" % catalog)
|
|
|
|
|
|
|
+ def _use_catalog(self, session, catalog):
|
|
|
|
|
+ self.db.configure_session(session['id'], "USE CATALOG `%s`" % catalog)
|
|
|
|
|
|
|
|
- def _use_database(self, catalog, database):
|
|
|
|
|
- session = self._get_session()
|
|
|
|
|
|
|
+ def _use_database(self, session, catalog, database):
|
|
|
if catalog:
|
|
if catalog:
|
|
|
- self.db.configure_session(session_handle=(session['id']),
|
|
|
|
|
- statement="USE `%(catalog)s`.`%(database)s`" % {'catalog': catalog,
|
|
|
|
|
- 'database': database})
|
|
|
|
|
|
|
+ self.db.configure_session(session['id'], "USE `%(catalog)s`.`%(database)s`" % {'catalog': catalog,
|
|
|
|
|
+ 'database': database})
|
|
|
else:
|
|
else:
|
|
|
- self.db.configure_session(session_handle=(session['id']),
|
|
|
|
|
- statement="USE `%(database)s`" % {'database': database})
|
|
|
|
|
|
|
+ self.db.configure_session(session['id'], "USE `%(database)s`" % {'database': database})
|
|
|
|
|
|
|
|
|
|
|
|
|
class FlinkSqlClient:
|
|
class FlinkSqlClient:
|