浏览代码

[flink] Unify the session cache per user

Romain 5 年之前
父节点
当前提交
57b4129dbf
共有 1 个文件被更改,包括 7 次插入4 次删除
  1. 7 4
      desktop/libs/notebook/src/notebook/connectors/flink_sql.py

+ 7 - 4
desktop/libs/notebook/src/notebook/connectors/flink_sql.py

@@ -71,7 +71,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):
-    session = self.db.create_session()
+    session = self._get_session()
 
 
     response = {
     response = {
       'type': lang,
       'type': lang,
@@ -87,19 +87,22 @@ class FlinkSqlApi(Api):
     }
     }
 
 
     if session_key not in SESSIONS:
     if session_key not in SESSIONS:
-      SESSIONS[session_key] = self.create_session()
+      SESSIONS[session_key] = self.db.create_session()
 
 
     try:
     try:
-      self.db.session_heartbeat(session_id=SESSIONS[session_key]['id'])
+      self.db.session_heartbeat(session_id=SESSIONS[session_key]['session_id'])
     except Exception as e:
     except Exception as e:
       if 'Session: %(id)s does not exist' % SESSIONS[session_key] in str(e):
       if 'Session: %(id)s does not exist' % SESSIONS[session_key] in str(e):
         LOG.warn('Session: %(id)s does not exist, opening a new one' % SESSIONS[session_key])
         LOG.warn('Session: %(id)s does not exist, opening a new one' % SESSIONS[session_key])
-        SESSIONS[session_key] = self.create_session()
+        SESSIONS[session_key] = self.db.create_session()
       else:
       else:
         raise e
         raise e
 
 
+    SESSIONS[session_key]['id'] = SESSIONS[session_key]['session_id']
+
     return SESSIONS[session_key]
     return SESSIONS[session_key]
 
 
+
   @query_error_handler
   @query_error_handler
   def execute(self, notebook, snippet):
   def execute(self, notebook, snippet):
     global n
     global n