Browse Source

HUE-8768 [kafka] Get live results via WS

Some refactoring coming in follow-up
Romain 5 years ago
parent
commit
cd72518698

+ 21 - 7
desktop/libs/kafka/src/kafka/ksql_client.py

@@ -130,6 +130,12 @@ class KSqlApi(object):
         # columns = line.keys()
         # data.append([line[col] for col in columns])
         if 'finalMessage' in line:
+          if has_channels() and channel_name:  # Send results via WS and empty results
+            _send_to_channel(
+                channel_name,
+                message_type='task.result',
+                message_data={'status': 'finalMessage', 'query_id': 1}
+            )
           break
         elif 'header' in line:
           continue
@@ -144,18 +150,26 @@ class KSqlApi(object):
             data.append(data_line['row']['columns'])
         else:
           data.append([line])
+
+        if has_channels() and channel_name:  # Send results via WS and empty results
+          _send_to_channel(
+              channel_name,
+              message_type='task.result',
+              message_data={'data': data, 'metadata': metadata, 'query_id': 1}
+          )
+          data = []  # TODO: special message when end of stream
     else:
       data, metadata = self._decode_result(
         self.ksql(statement)
       )
 
-    if has_channels() and channel_name:  # Send results via WS and empty results
-      _send_to_channel(
-          channel_name,
-          message_type='task.result',
-          message_data={'data': data, 'metadata': metadata, 'query_id': 1}
-      )
-      data = []  # TODO: special message when end of stream
+      if has_channels() and channel_name:  # Send results via WS and empty results
+        _send_to_channel(
+            channel_name,
+            message_type='task.result',
+            message_data={'data': data, 'metadata': metadata, 'query_id': 1}
+        )
+        data = []  # TODO: special message when end of stream
 
     return data, metadata
 

+ 5 - 1
desktop/libs/notebook/src/notebook/connectors/base.py

@@ -657,12 +657,16 @@ class ExecutionWrapper(object):
 
   def fetch(self, handle, start_over=None, rows=None):
     if start_over:
-      if not self.snippet['result'].get('handle') or not self.snippet['result']['handle'].get('guid') or not self.api.can_start_over(self.notebook, self.snippet):
+      if not self.snippet['result'].get('handle') \
+          or not self.snippet['result']['handle'].get('guid') \
+          or not self.api.can_start_over(self.notebook, self.snippet):
         start_over = False
         handle = self.api.execute(self.notebook, self.snippet)
         self.snippet['result']['handle'] = handle
+
         if self.callback and hasattr(self.callback, 'on_execute'):
           self.callback.on_execute(handle)
+
         self.should_close = True
         self._until_available()
 

+ 25 - 9
desktop/libs/notebook/src/notebook/connectors/ksql.py

@@ -23,6 +23,7 @@ from django.core.urlresolvers import reverse
 from django.utils.translation import ugettext as _
 
 from desktop.lib.i18n import force_unicode
+from desktop.conf import has_channels
 from kafka.ksql_client import KSqlApi as KSqlClientApi
 
 from notebook.connectors.base import Api, QueryError
@@ -31,6 +32,10 @@ from notebook.connectors.base import Api, QueryError
 LOG = logging.getLogger(__name__)
 
 
+if has_channels():
+  from notebook.consumer import _send_to_channel
+
+
 def query_error_handler(func):
   def decorator(*args, **kwargs):
     try:
@@ -53,24 +58,26 @@ class KSqlApi(Api):
 
   @query_error_handler
   def execute(self, notebook, snippet):
+    channel_name = notebook.get('editorWsChannel')
 
     data, description = self.db.query(
         snippet['statement'],
-        channel_name=notebook.get('editorWsChannel')
+        channel_name=channel_name
     )
     has_result_set = data is not None
 
     return {
-      'sync': True,
+      'sync': not (has_channels() and channel_name),
       'has_result_set': has_result_set,
       'result': {
-        'has_more': False,
-        'data': data if has_result_set else [],
-        'meta': [{
-          'name': col[0],
-          'type': col[1],
-          'comment': ''
-        } for col in description] if has_result_set else [],
+          'has_more': False,
+          'data': data if has_result_set else [],
+          'meta': [{
+            'name': col[0],
+            'type': col[1],
+            'comment': ''
+          } for col in description
+        ] if has_result_set else [],
         'type': 'table'
       }
     }
@@ -120,3 +127,12 @@ class KSqlApi(Api):
       response['error'] = e.message
 
     return response
+
+
+  def fetch_result(self, notebook, snippet, rows, start_over):
+    """Only called at the end of a live query."""
+    return {
+      'has_more': False,
+      'data': [],
+      'meta': []
+    }