File tree
41 files changed
+125
-141
lines changed- rocketmq-streams-channel-db/src/main/java/org/apache/rocketmq/streams/db/sink
- rocketmq-streams-channel-rocketmq/src/main/java/org/apache/rocketmq/streams/source
- rocketmq-streams-clients/src/main/java/org/apache/rocketmq/streams/client
- source
- transform
- rocketmq-streams-commons/src/main/java/org/apache/rocketmq/streams/common
- channel/impl/memory
- configurable
- functions
- topology
- builder
- stages
- udf
- rocketmq-streams-configurable/src/main/java/org/apache/rocketmq/streams
- configuable/service
- configurable/service
- rocketmq-streams-dim/src/main/java/org/apache/rocketmq/streams/dim/intelligence
- rocketmq-streams-examples/src/main/java/org/apache/rocketmq/streams/examples/rocketmqsource
- rocketmq-streams-filter/src/main/java/org/apache/rocketmq/streams/filter/operator
- action/impl
- rocketmq-streams-script/src/main/java/org/apache/rocketmq/streams/script/operator/impl
- rocketmq-streams-window/src/main/java/org/apache/rocketmq/streams/window/operator
Some content is hidden
Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.
41 files changed
+125
-141
lines changedLines changed: 2 additions & 2 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
17 | 17 |
| |
18 | 18 |
| |
19 | 19 |
| |
20 |
| - | |
| 20 | + | |
21 | 21 |
| |
22 | 22 |
| |
23 | 23 |
| |
24 | 24 |
| |
25 | 25 |
| |
26 | 26 |
| |
27 |
| - | |
| 27 | + | |
28 | 28 |
| |
29 | 29 |
| |
30 | 30 |
| |
|
Lines changed: 34 additions & 30 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
18 | 18 |
| |
19 | 19 |
| |
20 | 20 |
| |
| 21 | + | |
21 | 22 |
| |
22 | 23 |
| |
23 | 24 |
| |
24 | 25 |
| |
| 26 | + | |
25 | 27 |
| |
26 | 28 |
| |
27 | 29 |
| |
| |||
54 | 56 |
| |
55 | 57 |
| |
56 | 58 |
| |
| 59 | + | |
57 | 60 |
| |
58 | 61 |
| |
59 | 62 |
| |
| |||
72 | 75 |
| |
73 | 76 |
| |
74 | 77 |
| |
75 |
| - | |
| 78 | + | |
76 | 79 |
| |
77 |
| - | |
| 80 | + | |
| 81 | + | |
78 | 82 |
| |
79 |
| - | |
80 |
| - | |
| 83 | + | |
81 | 84 |
| |
82 | 85 |
| |
83 | 86 |
| |
| |||
93 | 96 |
| |
94 | 97 |
| |
95 | 98 |
| |
96 |
| - | |
| 99 | + | |
97 | 100 |
| |
98 | 101 |
| |
99 | 102 |
| |
| |||
108 | 111 |
| |
109 | 112 |
| |
110 | 113 |
| |
111 |
| - | |
112 |
| - | |
113 | 114 |
| |
114 |
| - | |
| 115 | + | |
115 | 116 |
| |
116 | 117 |
| |
117 | 118 |
| |
118 | 119 |
| |
119 |
| - | |
120 |
| - | |
121 |
| - | |
| 120 | + | |
| 121 | + | |
| 122 | + | |
122 | 123 |
| |
123 | 124 |
| |
124 | 125 |
| |
125 | 126 |
| |
126 |
| - | |
| 127 | + | |
127 | 128 |
| |
128 | 129 |
| |
129 | 130 |
| |
| |||
143 | 144 |
| |
144 | 145 |
| |
145 | 146 |
| |
146 |
| - | |
| 147 | + | |
147 | 148 |
| |
148 | 149 |
| |
149 | 150 |
| |
| |||
159 | 160 |
| |
160 | 161 |
| |
161 | 162 |
| |
| 163 | + | |
162 | 164 |
| |
163 |
| - | |
| 165 | + | |
164 | 166 |
| |
165 | 167 |
| |
166 | 168 |
| |
| |||
176 | 178 |
| |
177 | 179 |
| |
178 | 180 |
| |
179 |
| - | |
| 181 | + | |
180 | 182 |
| |
181 | 183 |
| |
182 | 184 |
| |
183 | 185 |
| |
184 | 186 |
| |
185 |
| - | |
| 187 | + | |
186 | 188 |
| |
187 | 189 |
| |
188 | 190 |
| |
189 | 191 |
| |
190 | 192 |
| |
191 | 193 |
| |
| 194 | + | |
192 | 195 |
| |
193 | 196 |
| |
194 | 197 |
| |
195 | 198 |
| |
196 | 199 |
| |
197 | 200 |
| |
198 | 201 |
| |
199 |
| - | |
200 |
| - | |
201 |
| - | |
202 |
| - | |
203 |
| - | |
| 202 | + | |
| 203 | + | |
| 204 | + | |
| 205 | + | |
| 206 | + | |
204 | 207 |
| |
205 | 208 |
| |
206 | 209 |
| |
| |||
212 | 215 |
| |
213 | 216 |
| |
214 | 217 |
| |
| 218 | + | |
215 | 219 |
| |
216 | 220 |
| |
217 | 221 |
| |
218 | 222 |
| |
219 | 223 |
| |
220 | 224 |
| |
221 | 225 |
| |
222 |
| - | |
223 |
| - | |
| 226 | + | |
| 227 | + | |
224 | 228 |
| |
225 | 229 |
| |
226 | 230 |
| |
227 | 231 |
| |
228 |
| - | |
229 |
| - | |
| 232 | + | |
| 233 | + | |
230 | 234 |
| |
231 | 235 |
| |
232 | 236 |
| |
| |||
309 | 313 |
| |
310 | 314 |
| |
311 | 315 |
| |
312 |
| - | |
| 316 | + | |
313 | 317 |
| |
314 | 318 |
| |
315 | 319 |
| |
| |||
367 | 371 |
| |
368 | 372 |
| |
369 | 373 |
| |
370 |
| - | |
371 |
| - | |
| 374 | + | |
| 375 | + | |
372 | 376 |
| |
373 | 377 |
| |
374 |
| - | |
375 |
| - | |
| 378 | + | |
| 379 | + | |
376 | 380 |
| |
377 | 381 |
|
Lines changed: 3 additions & 9 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
17 | 17 |
| |
18 | 18 |
| |
19 | 19 |
| |
20 |
| - | |
21 |
| - | |
22 |
| - | |
23 |
| - | |
24 | 20 |
| |
25 | 21 |
| |
26 | 22 |
| |
| |||
29 | 25 |
| |
30 | 26 |
| |
31 | 27 |
| |
32 |
| - | |
33 | 28 |
| |
34 | 29 |
| |
35 | 30 |
| |
36 |
| - | |
37 | 31 |
| |
38 | 32 |
| |
39 | 33 |
| |
| |||
48 | 42 |
| |
49 | 43 |
| |
50 | 44 |
| |
51 |
| - | |
| 45 | + | |
52 | 46 |
| |
53 | 47 |
| |
54 | 48 |
| |
| |||
67 | 61 |
| |
68 | 62 |
| |
69 | 63 |
| |
70 |
| - | |
| 64 | + | |
71 | 65 |
| |
72 | 66 |
| |
73 | 67 |
| |
74 | 68 |
| |
75 |
| - | |
| 69 | + | |
76 | 70 |
| |
77 | 71 |
| |
78 | 72 |
|
rocketmq-streams-clients/src/main/java/org/apache/rocketmq/streams/client/transform/DataStream.java
Lines changed: 8 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
66 | 66 |
| |
67 | 67 |
| |
68 | 68 |
| |
| 69 | + | |
| 70 | + | |
| 71 | + | |
| 72 | + | |
| 73 | + | |
69 | 74 |
| |
70 | 75 |
| |
71 | 76 |
| |
| |||
75 | 80 |
| |
76 | 81 |
| |
77 | 82 |
| |
| 83 | + | |
78 | 84 |
| |
79 | 85 |
| |
80 | 86 |
| |
| |||
358 | 364 |
| |
359 | 365 |
| |
360 | 366 |
| |
| 367 | + | |
| 368 | + | |
361 | 369 |
| |
362 | 370 |
| |
363 | 371 |
| |
|
Lines changed: 2 additions & 2 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
21 | 21 |
| |
22 | 22 |
| |
23 | 23 |
| |
24 |
| - | |
| 24 | + | |
25 | 25 |
| |
26 | 26 |
| |
27 | 27 |
| |
28 |
| - | |
| 28 | + | |
29 | 29 |
| |
30 | 30 |
| |
31 | 31 |
| |
|
Lines changed: 2 additions & 2 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
17 | 17 |
| |
18 | 18 |
| |
19 | 19 |
| |
20 |
| - | |
| 20 | + | |
21 | 21 |
| |
22 | 22 |
| |
23 |
| - | |
| 23 | + | |
24 | 24 |
| |
25 | 25 |
| |
26 | 26 |
| |
|
Lines changed: 1 addition & 1 deletion
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
16 | 16 |
| |
17 | 17 |
| |
18 | 18 |
| |
19 |
| - | |
| 19 | + | |
20 | 20 |
| |
21 | 21 |
| |
22 | 22 |
| |
|
Lines changed: 1 addition & 3 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
16 | 16 |
| |
17 | 17 |
| |
18 | 18 |
| |
19 |
| - | |
20 |
| - | |
21 |
| - | |
| 19 | + | |
22 | 20 |
| |
23 | 21 |
| |
24 | 22 |
|
Lines changed: 1 addition & 1 deletion
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
19 | 19 |
| |
20 | 20 |
| |
21 | 21 |
| |
22 |
| - | |
| 22 | + | |
23 | 23 |
| |
24 | 24 |
| |
25 | 25 |
|
Lines changed: 3 additions & 1 deletion
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
16 | 16 |
| |
17 | 17 |
| |
18 | 18 |
| |
19 |
| - | |
| 19 | + | |
| 20 | + | |
| 21 | + |
Lines changed: 1 addition & 3 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
16 | 16 |
| |
17 | 17 |
| |
18 | 18 |
| |
19 |
| - | |
20 |
| - | |
21 |
| - | |
| 19 | + | |
22 | 20 |
| |
23 | 21 |
| |
24 | 22 |
|
Lines changed: 1 addition & 3 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
16 | 16 |
| |
17 | 17 |
| |
18 | 18 |
| |
19 |
| - | |
20 |
| - | |
21 |
| - | |
| 19 | + | |
22 | 20 |
| |
23 | 21 |
| |
24 | 22 |
|
Lines changed: 2 additions & 2 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
25 | 25 |
| |
26 | 26 |
| |
27 | 27 |
| |
28 |
| - | |
| 28 | + | |
29 | 29 |
| |
30 | 30 |
| |
31 | 31 |
| |
| |||
40 | 40 |
| |
41 | 41 |
| |
42 | 42 |
| |
43 |
| - | |
| 43 | + | |
44 | 44 |
| |
45 | 45 |
| |
46 | 46 |
| |
|
0 commit comments