|
|
@@ -212,11 +212,11 @@ class FlinkSqlApi(Api):
|
|
|
response = {}
|
|
|
|
|
|
if database is None:
|
|
|
- response['databases'] = self.show_databases()
|
|
|
+ response['databases'] = self._show_databases()
|
|
|
elif table is None:
|
|
|
- response['tables_meta'] = self.show_tables(database)
|
|
|
+ response['tables_meta'] = self._show_tables(database)
|
|
|
elif column is None:
|
|
|
- columns = self.get_columns(database, table)
|
|
|
+ columns = self._get_columns(database, table)
|
|
|
response['columns'] = [col['name'] for col in columns]
|
|
|
response['extended_columns'] = [{
|
|
|
'comment': col.get('comment'),
|
|
|
@@ -248,7 +248,29 @@ class FlinkSqlApi(Api):
|
|
|
return response
|
|
|
|
|
|
|
|
|
- def show_databases(self):
|
|
|
+ def cancel(self, notebook, snippet):
|
|
|
+ session = self._get_session()
|
|
|
+ statement_id = snippet['result']['handle']['guid']
|
|
|
+
|
|
|
+ try:
|
|
|
+ self.db.close_statement(session_id=session['id'], job_id=statement_id)
|
|
|
+ except Exception as e:
|
|
|
+ if 'does not exist in current session:' in str(e):
|
|
|
+ return {'status': -1} # skipped
|
|
|
+ else:
|
|
|
+ raise e
|
|
|
+
|
|
|
+ return {'status': 0}
|
|
|
+
|
|
|
+
|
|
|
+ def close_session(self, session):
|
|
|
+ # Avoid closing session on page refresh or editor close for now
|
|
|
+ pass
|
|
|
+ # session = self._get_session()
|
|
|
+ # self.db.close_session(session['id'])
|
|
|
+
|
|
|
+
|
|
|
+ def _show_databases(self):
|
|
|
session = self._get_session()
|
|
|
session_id = session['id']
|
|
|
|
|
|
@@ -257,7 +279,7 @@ class FlinkSqlApi(Api):
|
|
|
return [db[0] for db in resp['results'][0]['data']]
|
|
|
|
|
|
|
|
|
- def show_tables(self, database):
|
|
|
+ def _show_tables(self, database):
|
|
|
session = self._get_session()
|
|
|
session_id = session['id']
|
|
|
|
|
|
@@ -267,7 +289,7 @@ class FlinkSqlApi(Api):
|
|
|
return [table[0] for table in resp['results'][0]['data']]
|
|
|
|
|
|
|
|
|
- def get_columns(self, database, table):
|
|
|
+ def _get_columns(self, database, table):
|
|
|
session = self._get_session()
|
|
|
session_id = session['id']
|
|
|
|
|
|
@@ -285,28 +307,6 @@ class FlinkSqlApi(Api):
|
|
|
]
|
|
|
|
|
|
|
|
|
- def cancel(self, notebook, snippet):
|
|
|
- session = self._get_session()
|
|
|
- statement_id = snippet['result']['handle']['guid']
|
|
|
-
|
|
|
- try:
|
|
|
- self.db.close_statement(session_id=session['id'], job_id=statement_id)
|
|
|
- except Exception as e:
|
|
|
- if 'does not exist in current session:' in str(e):
|
|
|
- return {'status': -1} # skipped
|
|
|
- else:
|
|
|
- raise e
|
|
|
-
|
|
|
- return {'status': 0}
|
|
|
-
|
|
|
-
|
|
|
- def close_session(self, session):
|
|
|
- # Avoid closing session on page refresh or editor close for now
|
|
|
- pass
|
|
|
- # session = self._get_session()
|
|
|
- # self.db.close_session(session['id'])
|
|
|
-
|
|
|
-
|
|
|
class FlinkSqlClient():
|
|
|
'''
|
|
|
Implements https://github.com/ververica/flink-sql-gateway
|