forked from nats-io/nats.py
-
Notifications
You must be signed in to change notification settings - Fork 0
/
direct_get.py
110 lines (88 loc) · 2.83 KB
/
direct_get.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
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
import asyncio
import nats
from nats.js import api
from nats.js import errors
async def main():
nc = await nats.connect("localhost")
# Create JetStream context.
js = nc.jetstream()
# Clean-up KV in case it already exists
try:
await js.delete_key_value("demo")
except errors.NotFoundError:
pass
# Create a KV store
kv = await js.create_key_value(
api.KeyValueConfig(bucket="demo", history=64, direct=True)
)
# Set a key
await kv.put("device.a", b"bar")
# Set another key
await kv.put("device.b", b"foo")
# Update a key
await kv.put("device.a", b"baz")
# Set a key
await kv.put("device.c", b"qux")
# Use case: Read latest value for all devices with a batch size
# Usage 1: Async iterator which can be reentered
async with js.direct_get(
"KV_demo",
batch_size=10,
multi_last=["$KV.demo.device.*"],
) as response:
# Because a batch size is used,
# we need to loop until response
# is no longer pending.
while response.pending():
# Iterate over a batch of messages
async for msg in response:
print(f"Pending: {response.pending()}")
print(f"Data: {msg.data}")
# Print response result. It's only possible
# to access result when a batch has been fully
# processed.
result = response.result()
print(f"Result: {result}")
# Usage 2: Async iterator which does not need to be reentered
async with js.direct_get(
"KV_demo",
batch_size=1,
multi_last=["$KV.demo.device.*"],
continue_on_eob=True,
) as response:
async for msg in response:
print(f"Pending: {response.pending()}")
print(f"Data: {msg.data}")
result = response.result()
print(f"Result: {result}")
# Usage 3: Using a callback to process messages until batch is ended
async def cb(msg: api.RawStreamMsg) -> None:
print(f"Msg sequence: {msg.seq}")
print(f"Data: {msg.data}")
# Read latest value using a callback
result = await js.direct_get(
"KV_demo",
batch_size=1,
cb=cb,
multi_last=["$KV.demo.device.*"],
)
print(result)
# Continue reading from last sequence
while result.num_pending:
result = await js.direct_get(
"KV_demo",
batch_size=1,
cb=cb,
multi_last=["$KV.demo.device.*"],
continue_from=result,
)
# Usage 4: Using a callback to process messages until whole response is received
result = await js.direct_get(
"KV_demo",
batch_size=10,
cb=cb,
multi_last=["$KV.demo.device.*"],
continue_on_eob=True,
)
if __name__ == "__main__":
asyncio.run(main())