-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathhello_world.py
88 lines (67 loc) · 2.77 KB
/
hello_world.py
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
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
import faust
from faust import Stream
from models import Customer, CustomerEvent
app = faust.App('hello-world', broker='kafka://localhost:29092', store="memory://")
customers = app.Table('customers', partitions=1, value_type=Customer)
customer_became_lead_topic = app.topic('customer_became_lead', partitions=1, value_type=CustomerEvent)
customer_subscribed_topic = app.topic('customer_subscribed', partitions=1, value_type=CustomerEvent)
customer_cancelled_topic = app.topic('customer_cancelled', partitions=1, value_type=CustomerEvent)
@app.agent(customer_became_lead_topic)
async def customer_became_lead(events: Stream[CustomerEvent]):
async for event in events:
try:
customer = customers[f'{event.business_unit}_{event.customer_id}']
customer.updated_at = event.timestamp
customer.status = 'lead'
except KeyError:
customer = Customer(
id=event.customer_id,
business_unit=event.business_unit,
created_at=event.timestamp,
updated_at=event.timestamp,
)
customers[customer.key] = customer
print(customer)
@app.agent(customer_subscribed_topic)
async def customer_subscribed(events: Stream[CustomerEvent]):
async for event in events:
try:
customer = customers[f'{event.business_unit}_{event.customer_id}']
customer.status = 'active'
customer.subscribed_at = event.timestamp
customer.updated_at = event.timestamp
except KeyError:
customer = Customer(
id=event.customer_id,
business_unit=event.business_unit,
status='active',
subscribed_at=event.timestamp,
updated_at=event.timestamp,
)
customers[customer.key] = customer
print(customer)
@app.agent(customer_cancelled_topic)
async def customer_cancelled(events: Stream[CustomerEvent]):
async for event in events:
try:
customer = customers[f'{event.business_unit}_{event.customer_id}']
customer.status = 'former'
customer.cancelled_at = event.timestamp
except KeyError:
customer = Customer(
id=event.customer_id,
business_unit=event.business_unit,
status='active',
cancelled_at=event.timestamp,
updated_at=event.timestamp,
)
customers[customer.key] = customer
print(customer)
@app.page('/customers')
async def list_customers(self, request):
return self.json(list(customers.values()))
@app.page('/customers')
async def list_customers(self, request):
return self.json(list(customers.values()))
if __name__ == '__main__':
app.main()