-
Notifications
You must be signed in to change notification settings - Fork 5
Expand file tree
/
Copy path1_device_data.py
More file actions
50 lines (36 loc) · 1.02 KB
/
1_device_data.py
File metadata and controls
50 lines (36 loc) · 1.02 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
from datetime import datetime
from pydantic import BaseModel
from chalk.features import DataFrame, has_many, features, FeatureTime
from chalk.streams import stream, KafkaSource
@features
class Measurement:
device_id: str
lat: float
long: float
voltage: float
temp: float
timestamp: FeatureTime
@features
class Sensor:
id: str
measurements: DataFrame[Measurement] = has_many(lambda: Measurement.device_id == Sensor.id)
source = KafkaSource(name="sensor_stream")
class DeviceDataJson(BaseModel):
latitude: float
longitude: float
voltage: float
temperature: float
class Message(BaseModel):
device_id: str
timestamp: datetime
data: DeviceDataJson
@stream(source=source)
def read_message(message: Message) -> Measurement:
return Measurement(
device_id=message.device_id,
timestamp=message.timestamp,
lat=message.data.latitude,
long=message.data.longitude,
voltage=message.data.voltage,
temp=message.data.temperature,
)