Pārlūkot izejas kodu

[ksql] Handle generator raised StopIteration on query result

Romain 5 gadi atpakaļ
vecāks
revīzija
40dd8ab87a
1 mainītis faili ar 32 papildinājumiem un 26 dzēšanām
  1. 32 26
      desktop/libs/kafka/src/kafka/ksql_client.py

+ 32 - 26
desktop/libs/kafka/src/kafka/ksql_client.py

@@ -128,38 +128,44 @@ class KSqlApi(object):
             message_data={'status': 'running', 'query_id': 1}
         )
 
-      for line in result:
-        # columns = line.keys()
-        # data.append([line[col] for col in columns])
-        if 'finalMessage' in line:
+      try:
+        for line in result:
+          # 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
+          else:
+            line = line.strip()[:-1]
+
+          if is_select:
+            data_line = json.loads(line)
+            if data_line.get('@type') == 'statement_error':
+              raise KSqlApiException(data_line['message'])
+            if data_line['row']:  # If limit not reached
+              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={'status': 'finalMessage', 'query_id': 1}
+                message_data={'data': data, 'meta': metadata, 'query_id': 1}
             )
-          break
-        elif 'header' in line:
-          continue
-        else:
-          line = line.strip()[:-1]
-
-        if is_select:
-          data_line = json.loads(line)
-          if data_line.get('@type') == 'statement_error':
-            raise KSqlApiException(data_line['message'])
-          if data_line['row']:  # If limit not reached
-            data.append(data_line['row']['columns'])
+            data = []  # TODO: special message when end of stream
+      except RuntimeError as e:
+        if 'generator raised StopIteration' in str(e):
+          return data, metadata
         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, 'meta': metadata, 'query_id': 1}
-          )
-          data = []  # TODO: special message when end of stream
+          raise e
     else:
       data, metadata = self._decode_result(
         self.ksql(statement)