Skip to content

Commit a627e41

Browse files
committed
Add CsvReader, cleanup, and address comments
Signed-off-by: Arham Chopra <arham.chopra@cubistsystematic.com>
1 parent cc39ae1 commit a627e41

5 files changed

Lines changed: 301 additions & 228 deletions

File tree

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1 +1,5 @@
1-
from .filedrop import *
1+
try:
2+
from .adapter import *
3+
from .filedrop import *
4+
except ImportError:
5+
pass
Lines changed: 204 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,204 @@
1+
import csv
2+
import logging
3+
from dataclasses import dataclass
4+
from datetime import datetime
5+
from enum import Enum, auto
6+
from typing import Any, Dict, List, get_args, get_origin
7+
8+
import orjson
9+
import pyarrow.parquet as pq
10+
from csp import ts
11+
from csp.impl.pushadapter import PushInputAdapter
12+
from csp.impl.types.container_type_normalizer import ContainerTypeNormalizer
13+
from csp.impl.wiring import py_push_adapter_def
14+
from watchdog.events import FileSystemEvent, FileSystemEventHandler
15+
from watchdog.observers import Observer
16+
17+
__all__ = (
18+
"FileDropType",
19+
"FileDropAdapterConfiguration",
20+
"filedrop_adapter_def",
21+
)
22+
23+
24+
log = logging.getLogger(__name__)
25+
26+
27+
class FileDropType(Enum):
28+
CSV = auto()
29+
JSON = auto()
30+
PARQUET = auto()
31+
32+
33+
@dataclass
34+
class FileDropAdapterConfiguration:
35+
"""Configuration for the filedrop push adapter"""
36+
37+
# Path to the directory to monitor
38+
dir_path: str
39+
# Format of files to expect to load properly i.e parquet, json, csv
40+
filedrop_type: FileDropType
41+
# Map the data fields from the file to the fields of the structs
42+
field_map: Dict[str, str]
43+
# List of extensions to filter, empty list means all extensions are allowed
44+
extensions: List[str]
45+
# Extra args to the type adapter deserializer
46+
type_adapter_args: Dict[str, Any]
47+
48+
49+
class FileReaderBase:
50+
"""The base file reader that reads data from files and generates structs"""
51+
52+
def __init__(self, config: FileDropAdapterConfiguration, ts_typ: object):
53+
self.field_map = config.field_map
54+
self.extensions = config.extensions
55+
normalized_type = ContainerTypeNormalizer.normalize_type(ts_typ)
56+
self.is_list = get_origin(normalized_type) is list
57+
if self.is_list:
58+
inner_type = get_args(normalized_type)[0]
59+
type_adapter = inner_type.type_adapter()
60+
else:
61+
type_adapter = ts_typ.type_adapter()
62+
self.type_adapter = type_adapter
63+
if hasattr(config, "type_adapter_args"):
64+
self.context = config.type_adapter_args
65+
else:
66+
self.context = {}
67+
68+
def read(self, src_path: str) -> object:
69+
"""Generator to return stucts from a filepath"""
70+
71+
should_read = True
72+
if self.extensions:
73+
if not any([src_path.endswith(suffix) for suffix in self.extensions]):
74+
should_read = False
75+
if should_read:
76+
dicts = self.read_impl(src_path)
77+
structs = [self.deserialize_dict(self.apply_field_map(v)) for v in dicts]
78+
if self.is_list:
79+
yield structs
80+
else:
81+
for s in structs:
82+
yield s
83+
84+
def read_impl(self, src_path: str) -> List[dict]:
85+
"""File type specific implementation"""
86+
87+
raise Exception(f"read not implemented for {self}")
88+
89+
def apply_field_map(self, data: dict) -> dict:
90+
"""Convert the keys in the data to the field names of the struct"""
91+
92+
if self.field_map:
93+
new_data = {}
94+
for k, v in data.items():
95+
new_k = self.field_map.get(k, k)
96+
new_data[new_k] = v
97+
return new_data
98+
else:
99+
return data
100+
101+
def deserialize_dict(self, dict: dict) -> object:
102+
"""Convert a dict to struct"""
103+
104+
return self.type_adapter.validate_python(dict, context=self.context)
105+
106+
107+
class FileReaderCsv(FileReaderBase):
108+
"""File reader for json file type"""
109+
110+
def read_impl(self, src_path: str) -> List[dict]:
111+
data = []
112+
with open(src_path, "r") as f:
113+
reader = csv.DictReader(f)
114+
for row in reader:
115+
data.append(row)
116+
return data
117+
118+
119+
class FileReaderJson(FileReaderBase):
120+
"""File reader for json file type"""
121+
122+
def read_impl(self, src_path: str) -> List[dict]:
123+
with open(src_path, "rb") as f:
124+
data = orjson.loads(f.read())
125+
if isinstance(data, list):
126+
res = data
127+
else:
128+
res = [data]
129+
return res
130+
131+
132+
class FileReaderParquet(FileReaderBase):
133+
"""File reader for parquet file type"""
134+
135+
def read_impl(self, src_path: str) -> List[dict]:
136+
table = pq.read_table(src_path)
137+
return table.to_pylist()
138+
139+
140+
class EventHandlerCustom(FileSystemEventHandler):
141+
def __init__(self, adapter: PushInputAdapter, file_reader: FileReaderBase):
142+
self.file_reader = file_reader
143+
self.adapter = adapter
144+
self._created = set()
145+
self._opened = set()
146+
self._modified = set()
147+
self._closed = set()
148+
149+
def on_created(self, event: FileSystemEvent):
150+
self._created.add(event.src_path)
151+
152+
def on_opened(self, event: FileSystemEvent):
153+
if event.src_path in self._created:
154+
self._created.remove(event.src_path)
155+
self._opened.add(event.src_path)
156+
157+
def on_modified(self, event: FileSystemEvent):
158+
if event.src_path in self._opened:
159+
self._opened.remove(event.src_path)
160+
self._modified.add(event.src_path)
161+
162+
def on_closed(self, event: FileSystemEvent):
163+
if event.src_path in self._modified:
164+
self._modified.remove(event.src_path)
165+
file_path = event.src_path
166+
try:
167+
for data in self.file_reader.read(file_path):
168+
self.adapter.push_tick(data)
169+
except Exception as e:
170+
log.error(f"Failed to read data from {file_path} with exception: {e}, skipping")
171+
172+
173+
class _FileDropImpl(PushInputAdapter):
174+
FILEREADER_MAP = {
175+
FileDropType.CSV: FileReaderCsv,
176+
FileDropType.JSON: FileReaderJson,
177+
FileDropType.PARQUET: FileReaderParquet,
178+
}
179+
180+
def __init__(self, config: FileDropAdapterConfiguration, ts_typ: "T"): # noqa
181+
self.dir_path = config.dir_path
182+
self.observer = Observer()
183+
reader = self.FILEREADER_MAP[config.filedrop_type]
184+
file_reader = reader(config, ts_typ)
185+
self.event_handler = EventHandlerCustom(self, file_reader)
186+
self.observer.schedule(self.event_handler, self.dir_path, recursive=False)
187+
188+
def start(self, starttime: datetime, endtime: datetime):
189+
self.observer.start()
190+
self.observer_started = True
191+
192+
def stop(self):
193+
if self.observer_started:
194+
self.observer.stop()
195+
self.observer.join()
196+
197+
198+
filedrop_adapter_def = py_push_adapter_def(
199+
"filedrop_adapter_def",
200+
_FileDropImpl,
201+
ts["T"],
202+
config=FileDropAdapterConfiguration,
203+
ts_typ="T",
204+
)

0 commit comments

Comments
 (0)