-
Notifications
You must be signed in to change notification settings - Fork 11
/
ctx.clj
166 lines (95 loc) · 3.95 KB
/
ctx.clj
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
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
(ns dvlopt.kstreams.ctx
"Contexts are needed for low-level processors."
{:author "Adam Helinski"}
(:refer-clojure :exclude [partition])
(:require [dvlopt.kafka :as K]
[dvlopt.kafka.-interop.clj :as K.-interop.clj]
[dvlopt.kafka.-interop.java :as K.-interop.java]
[dvlopt.kstreams :as KS])
(:import (org.apache.kafka.streams.processor ProcessorContext
To)
(org.apache.kafka.streams.state KeyValueStore
WindowStore
SessionStore)))
;;;;;;;;;;
(defn commit
"Manually requests a commit of the current progress.
Commits are handle automatically by the library. This function request the next commit to happen as soon
as possible."
^ProcessorContext
[^ProcessorContext ctx]
(.commit ctx)
ctx)
(defn forward
"Forwards a key-value to child nodes. Can be used in the :dvlopt.kstreams/processor.on-record function of a processor
or during the execution of a function added with `schedule`.
A map of options may be given :
:dvlopt.kafka/key
Unserialized key.
:dvlopt.kafka/timestamp
Chosen timestamp.
If missing, selects the timestamp of the record being processed by :dvlopt.kstreams/processor.on-record or the
moment the scheduled function executes.
:dvlopt.kafka/value
Unserialized value.
:dvlopt.streams/child
Name of a specific child node.
If missing, forwards to all children."
^ProcessorContext
[^ProcessorContext ctx options]
(let [to (if-let [child (::KS/child options)]
(To/child child)
(To/all))]
(some->> (::K/timestamp options)
(.withTimestamp to))
(.forward ctx
(::K/key options)
(::K/value options)
to))
ctx)
(defn schedule
"Schedules a periodic operation for processors (may be used during :dvlopt.kstreams/processor.init and/or
:dvlopt.kstreams/processor.on-record). Can be called multiple times on the same context.
The time interval must be as described in `dvlopt.kafka` (best effort for millisecond precision).
The time type is either :
:stream-time
Time advances following timestamps extracted from records.
An operation is skipped if stream time advances more than the interval.
:wall-clock-time
Times advances following system time (best effort).
An operation is skipped if garbage-collection halts the world for too long or if the current operation
takes more time to complete than the interval.
The callback accepts only 1 argument, the timestamp of \"when\" it is called (depending on time type).
Returns a no-op function for cancelling the operation.
Ex. (schedule ctx
[5 :seconds]
:stream-time
(fn callback [timestamp]
(do-stuff-like-forwarding-records ...)))"
^ProcessorContext
[^ProcessorContext ctx interval time-type callback]
(K.-interop.clj/cancellable (.schedule ctx
(K.-interop.java/to-milliseconds interval)
(K.-interop.java/punctuation-type time-type)
(K.-interop.java/punctuator callback))))
(defn kv-store
"Retrieves a writable key-value store.
Cf. `dvlopt.kstreams.store`"
^KeyValueStore
[^ProcessorContext ctx store-name]
(.getStateStore ctx
store-name))
(defn window-store
"Retrieves a writable window store.
Cf. `dvlopt.kstreams.store`"
^WindowStore
[^ProcessorContext ctx store-name]
(.getStateStore ctx
store-name))
(defn session-store
"Retrieves a writable session store.
Cf. `dvlopt.kstreams.store`"
^SessionStore
[^ProcessorContext ctx store-name]
(.getStateStore ctx
store-name))