forked from feldera/feldera
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy path_callback_runner.py
More file actions
64 lines (53 loc) · 1.94 KB
/
Copy path_callback_runner.py
File metadata and controls
64 lines (53 loc) · 1.94 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
from threading import Thread
from typing import Callable, Optional
import pandas as pd
from feldera import FelderaClient
from feldera._helpers import dataframe_from_response
from feldera.enums import PipelineFieldSelector
class CallbackRunner(Thread):
def __init__(
self,
client: FelderaClient,
pipeline_name: str,
view_name: str,
callback: Callable[[pd.DataFrame, int], None],
):
super().__init__()
self.daemon = True
self.client: FelderaClient = client
self.pipeline_name: str = pipeline_name
self.view_name: str = view_name
self.callback: Callable[[pd.DataFrame, int], None] = callback
self.schema: Optional[dict] = None
def run(self):
"""
The main loop of the thread. Listens for data and calls the callback function on each chunk of data received.
:meta private:
"""
pipeline = self.client.get_pipeline(
self.pipeline_name, PipelineFieldSelector.ALL
)
schemas = pipeline.tables + pipeline.views
for schema in schemas:
if schema.name == self.view_name:
self.schema = schema
break
if self.schema is None:
raise ValueError(
f"Table or View {self.view_name} not found in the pipeline schema."
)
gen_obj = self.client.listen_to_pipeline(
self.pipeline_name,
self.view_name,
format="json",
case_sensitive=self.schema.case_sensitive,
)
iterator = gen_obj()
for chunk in iterator:
chunk: dict = chunk
data: Optional[list[dict]] = chunk.get("json_data")
seq_no: Optional[int] = chunk.get("sequence_number")
if data is not None and seq_no is not None:
self.callback(
dataframe_from_response([data], self.schema.fields), seq_no
)