File tree
6 files changed
+63
-32
lines changed- flink-connector-kafka/src
- main/java/org/apache/flink/connector/kafka/dynamic/source
- enumerator
- reader
- test/java/org/apache/flink
- connector/kafka/dynamic/source/enumerator
- streaming/connectors/kafka
6 files changed
+63
-32
lines changedLines changed: 28 additions & 21 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
96 | 96 |
| |
97 | 97 |
| |
98 | 98 |
| |
| 99 | + | |
99 | 100 |
| |
100 | 101 |
| |
101 | 102 |
| |
| |||
151 | 152 |
| |
152 | 153 |
| |
153 | 154 |
| |
| 155 | + | |
154 | 156 |
| |
155 | 157 |
| |
156 | 158 |
| |
| |||
212 | 214 |
| |
213 | 215 |
| |
214 | 216 |
| |
215 |
| - | |
216 |
| - | |
217 |
| - | |
218 |
| - | |
219 |
| - | |
220 |
| - | |
221 |
| - | |
222 |
| - | |
223 |
| - | |
224 |
| - | |
225 |
| - | |
226 |
| - | |
227 |
| - | |
228 |
| - | |
229 |
| - | |
230 |
| - | |
231 |
| - | |
232 |
| - | |
233 |
| - | |
234 |
| - | |
| 217 | + | |
| 218 | + | |
| 219 | + | |
| 220 | + | |
| 221 | + | |
| 222 | + | |
| 223 | + | |
| 224 | + | |
| 225 | + | |
| 226 | + | |
| 227 | + | |
| 228 | + | |
| 229 | + | |
| 230 | + | |
235 | 231 |
| |
236 | 232 |
| |
237 | 233 |
| |
238 | 234 |
| |
239 | 235 |
| |
240 | 236 |
| |
| 237 | + | |
241 | 238 |
| |
242 | 239 |
| |
243 | 240 |
| |
| |||
370 | 367 |
| |
371 | 368 |
| |
372 | 369 |
| |
| 370 | + | |
| 371 | + | |
| 372 | + | |
| 373 | + | |
| 374 | + | |
| 375 | + | |
| 376 | + | |
373 | 377 |
| |
374 | 378 |
| |
375 |
| - | |
| 379 | + | |
| 380 | + | |
| 381 | + | |
| 382 | + | |
376 | 383 |
| |
377 | 384 |
| |
378 | 385 |
| |
|
Lines changed: 24 additions & 5 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
34 | 34 |
| |
35 | 35 |
| |
36 | 36 |
| |
| 37 | + | |
| 38 | + | |
37 | 39 |
| |
38 | 40 |
| |
39 | 41 |
| |
| |||
69 | 71 |
| |
70 | 72 |
| |
71 | 73 |
| |
| 74 | + | |
72 | 75 |
| |
73 | 76 |
| |
74 | 77 |
| |
| |||
79 | 82 |
| |
80 | 83 |
| |
81 | 84 |
| |
| 85 | + | |
82 | 86 |
| |
83 | 87 |
| |
84 | 88 |
| |
85 | 89 |
| |
86 |
| - | |
| 90 | + | |
| 91 | + | |
87 | 92 |
| |
88 | 93 |
| |
89 | 94 |
| |
90 | 95 |
| |
91 | 96 |
| |
92 | 97 |
| |
| 98 | + | |
93 | 99 |
| |
94 | 100 |
| |
95 | 101 |
| |
| |||
147 | 153 |
| |
148 | 154 |
| |
149 | 155 |
| |
150 |
| - | |
| 156 | + | |
| 157 | + | |
| 158 | + | |
151 | 159 |
| |
| 160 | + | |
| 161 | + | |
| 162 | + | |
| 163 | + | |
152 | 164 |
| |
153 | 165 |
| |
154 | 166 |
| |
| |||
286 | 298 |
| |
287 | 299 |
| |
288 | 300 |
| |
289 |
| - | |
| 301 | + | |
| 302 | + | |
290 | 303 |
| |
291 | 304 |
| |
292 |
| - | |
| 305 | + | |
| 306 | + | |
| 307 | + | |
| 308 | + | |
293 | 309 |
| |
294 |
| - | |
| 310 | + | |
| 311 | + | |
| 312 | + | |
| 313 | + | |
295 | 314 |
| |
296 | 315 |
| |
297 | 316 |
|
Lines changed: 2 additions & 2 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
132 | 132 |
| |
133 | 133 |
| |
134 | 134 |
| |
135 |
| - | |
136 |
| - | |
| 135 | + | |
| 136 | + | |
137 | 137 |
| |
138 | 138 |
| |
139 | 139 |
| |
|
Lines changed: 4 additions & 2 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
919 | 919 |
| |
920 | 920 |
| |
921 | 921 |
| |
| 922 | + | |
922 | 923 |
| |
923 | 924 |
| |
924 | 925 |
| |
925 | 926 |
| |
926 |
| - | |
| 927 | + | |
| 928 | + | |
927 | 929 |
| |
928 | 930 |
| |
929 | 931 |
| |
| |||
939 | 941 |
| |
940 | 942 |
| |
941 | 943 |
| |
942 |
| - | |
| 944 | + | |
943 | 945 |
| |
944 | 946 |
| |
945 | 947 |
| |
|
Lines changed: 2 additions & 1 deletion
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
150 | 150 |
| |
151 | 151 |
| |
152 | 152 |
| |
153 |
| - | |
| 153 | + | |
| 154 | + | |
154 | 155 |
| |
155 | 156 |
| |
156 | 157 |
| |
|
Lines changed: 3 additions & 1 deletion
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
204 | 204 |
| |
205 | 205 |
| |
206 | 206 |
| |
| 207 | + | |
207 | 208 |
| |
208 | 209 |
| |
209 | 210 |
| |
| |||
219 | 220 |
| |
220 | 221 |
| |
221 | 222 |
| |
222 |
| - | |
| 223 | + | |
223 | 224 |
| |
| 225 | + | |
224 | 226 |
| |
225 | 227 |
| |
226 | 228 |
| |
|
0 commit comments