Unable to filter using Dataframe RuntimeError: There is no current event loop in thread

オープン
#265 コメント 0 件 リアクション 0 件 担当者 0 名 GitHub で見る

まだ誰も着手していません。

評価

難易度
4/5
見積もり時間
3〜5日
初心者へのやさしさ
25/100
issue の種類
バグ
明瞭さ
説明が足りない
活発さ
停滞
技術スタック
pandas, python

調査の方向性

Start with streamz/dataframe/core.py at getitem, then follow map_partitions in streamz/collection.py and Condition creation in streamz/core.py. Reproduce the filtering expression from the issue in a worker thread and determine what is required for it to run without the missing-event-loop error; done means the filtered rolling pipeline initializes successfully.

索引モデルが issue の本文から書いたものです。

説明

RunTimeError when trying to filter dataframe in my websocket client. The dataframe works without filtering - has anyone run into this issue before?

class Socket(threading.Thread):
    
    def __init__(self, symbol=None):
        super(Socket, self).__init__()
        self.wsURL = "wss://stream.socket.com"
        
        self.source = Stream()
        self.example = pd.DataFrame({'A': [], 'B': [], 'C': []}, index=[]).astype(float)
        self.sdf = DataFrame(self.source, example=self.example)

        self.source2 = Stream()
        self.example2 = pd.DataFrame({'A': []}, index=[])
        self.pdf = DataFrame(self.source2, example=self.example2)
        
    def run(self):
        def on_message(ws, message):
            j_msg = json.loads(message)
            if 'data' in j_msg: 
                data = j_msg['data']
                if j_msg['data']['event'] == "event1":
                    df = pd.DataFrame({'A': data['a'], 'B': data['b'], 'C': data['c']})
                    self.source.emit(df)
                elif j_msg['data']['event'] == "event2":
                    df = pd.DataFrame({'A': data['a']})
                    self.source2.emit(df)

        def on_error(ws, error):
            print(error)

        def on_close(ws):
            print("### closed ###")
            
        websocket.enableTrace(True)
        self.ws = websocket.WebSocketApp(self.wsURL,
                                    on_message = on_message,
                                    on_error = on_error,
                                    on_close = on_close)
        
        self.sdf[self.sdf.A == 1].B.rolling('60s').sum().stream.sink(print)
        self.pdf.A.rolling('60s').std().stream.sink(print)
        
        self.ws.keep_running = True
        self.ws.run_forever()
        self.ws.keep_running = False

Exception in thread Thread-4:
Traceback (most recent call last):
  File "/opt/conda/envs/env/lib/python3.5/threading.py", line 914, in _bootstrap_inner
    self.run()
  File "<ipython-input-2-65c47f0e1bee>", line 76, in run
    self.sdf[self.sdf.A == 1].qty.rolling('60s').sum().stream.sink(print)
  File "/opt/conda/envs/env/lib/python3.5/site-packages/streamz/dataframe/core.py", line 201, in __getitem__
    return self.map_partitions(operator.getitem, self, index)
  File "/opt/conda/envs/env/lib/python3.5/site-packages/streamz/collection.py", line 32, in map_partitions
    stream = type(streams[0].stream).zip(*[getattr(arg, 'stream', arg) for arg in args])
  File "/opt/conda/envs/env/lib/python3.5/site-packages/streamz/core.py", line 214, in wrapped
    return func(*args, **kwargs)
  File "/opt/conda/envs/env/lib/python3.5/site-packages/streamz/core.py", line 1007, in __init__
    self.condition = Condition()
  File "/opt/conda/envs/env/lib/python3.5/site-packages/tornado/locks.py", line 110, in __init__
    self.io_loop = ioloop.IOLoop.current()
  File "/opt/conda/envs/env/lib/python3.5/site-packages/tornado/ioloop.py", line 282, in current
    loop = asyncio.get_event_loop()
  File "/opt/conda/envs/env/lib/python3.5/asyncio/events.py", line 678, in get_event_loop
    return get_event_loop_policy().get_event_loop()
  File "/opt/conda/envs/env/lib/python3.5/asyncio/events.py", line 584, in get_event_loop
    % threading.current_thread().name)
RuntimeError: There is no current event loop in thread 'Thread-4'.

主要言語
Python
スター
1.3k
フォーク
149
PR マージ指標
30日以内にマージされた PR はありません

コントリビューションガイド

コントリビューションガイドを開く

はじめの一歩

  1. issue を最後まで読み、次にプロジェクトのコントリビューションガイドを読みます。
  2. 着手することを issue にコメントします — 二人が同じ作業をするのを防げます。
  3. リポジトリをフォークし、ブランチを切って変更します。
  4. issue 番号を参照したプルリクエストを送ります。

python-streamz/streamz のほかの issue

python-streamz/streamz の issue をすべて見る

似ている issue

Python の issue をもっと見る

新しい issue をメールで受け取る

初心者向けの GitHub issue を短くまとめたダイジェスト。