Unable to filter using Dataframe RuntimeError: There is no current event loop in thread
Nessuno ha ancora preso questa issue.
Valutazione
- Difficoltà
- 4/5
- Tempo stimato
- 3-5 giorni
- Idoneità per principianti
- 25/100
- Tipo di issue
- Bug
- Chiarezza
- Da chiarire
- Stato di attività
- Ferma
- Stack tecnologico
- pandas, python
- Ambito
- data, stream-processing
Direzione di ricerca
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.
Scritto dal modello di indicizzazione a partire dal testo della issue.
Descrizione
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'.
- Lingua principale
- Python
- Stelle
- 1.3k
- Fork
- 149
- Merge medio
- 17h 39m
- PR unite (30g)
- 1
Guida per i contributori
Apri la guida per i contributori
Come iniziare
- Leggi tutta la issue e poi la guida ai contributi del progetto.
- Commenta sulla issue per dire che te ne occupi tu — evita che due persone facciano lo stesso lavoro.
- Fai un fork del repository e lavora su un branch.
- Apri una pull request che faccia riferimento al numero della issue.
Altre issue di python-streamz/streamz
-
pkg_resources warning Aperta
Difficoltà 3/5 1-2 giorni Idoneità per principianti 35/100
python-streamz/streamz#481 · 4 commenti ·
-
Combining the streamz.Stream.filenames() and streamz.Stream.from_textfile() using dask scatter? Aperta
Difficoltà 5/5 Più di una settimana Idoneità per principianti 25/100
python-streamz/streamz#480 · 2 commenti ·
-
Compile the code into c++ Aperta
Difficoltà 5/5 Più di una settimana Idoneità per principianti 20/100
python-streamz/streamz#479 · 1 commento ·
-
Difficoltà 5/5 Più di una settimana Idoneità per principianti 20/100
python-streamz/streamz#478 · 6 commenti ·
-
Difficoltà 4/5 3-5 giorni Idoneità per principianti 10/100
python-streamz/streamz#476 · 17 commenti · 2 reazioni ·
Tutte le issue di python-streamz/streamz
Issue simili
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 75/100
anthropics/skills#1811 · 1 commento ·
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 75/100
speaches-ai/speaches#678 ·
-
bug
Difficoltà 2/5 1-3 ore Idoneità per principianti 75/100
datalayer/mcp-compose#42 ·
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 75/100
conda-forge/spacy-feedstock#177 ·
-
Difficoltà 2/5 1-3 ore Idoneità per principianti 70/100
UKGovernmentBEIS/inspect_evals#2523 ·