File tree
11 files changed
+246
-11
lines changed- flink-connector-kafka
- src
- main/java/org/apache/flink
- connector/kafka/source
- reader/deserializer
- util
- streaming/connectors/kafka/table
- test/java/org/apache/flink/connector/kafka/source/reader/deserializer
11 files changed
+246
-11
lines changedLines changed: 8 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
72 | 72 |
| |
73 | 73 |
| |
74 | 74 |
| |
| 75 | + | |
75 | 76 |
| |
76 | 77 |
| |
| 78 | + | |
| 79 | + | |
| 80 | + | |
| 81 | + | |
| 82 | + | |
| 83 | + | |
| 84 | + | |
77 | 85 |
| |
78 | 86 |
| |
79 | 87 |
| |
|
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilder.java
Lines changed: 9 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
345 | 345 |
| |
346 | 346 |
| |
347 | 347 |
| |
| 348 | + | |
| 349 | + | |
| 350 | + | |
| 351 | + | |
| 352 | + | |
| 353 | + | |
| 354 | + | |
| 355 | + | |
| 356 | + | |
348 | 357 |
| |
349 | 358 |
| |
350 | 359 |
| |
|
Lines changed: 7 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
96 | 96 |
| |
97 | 97 |
| |
98 | 98 |
| |
| 99 | + | |
| 100 | + | |
| 101 | + | |
| 102 | + | |
| 103 | + | |
| 104 | + | |
| 105 | + | |
99 | 106 |
| |
100 | 107 |
| |
101 | 108 |
| |
|
Lines changed: 13 additions & 1 deletion
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
20 | 20 |
| |
21 | 21 |
| |
22 | 22 |
| |
| 23 | + | |
23 | 24 |
| |
24 | 25 |
| |
25 | 26 |
| |
| |||
35 | 36 |
| |
36 | 37 |
| |
37 | 38 |
| |
| 39 | + | |
38 | 40 |
| |
39 | 41 |
| |
| 42 | + | |
| 43 | + | |
| 44 | + | |
| 45 | + | |
| 46 | + | |
| 47 | + | |
40 | 48 |
| |
| 49 | + | |
41 | 50 |
| |
42 | 51 |
| |
43 | 52 |
| |
| |||
48 | 57 |
| |
49 | 58 |
| |
50 | 59 |
| |
51 |
| - | |
| 60 | + | |
| 61 | + | |
| 62 | + | |
| 63 | + | |
52 | 64 |
| |
53 | 65 |
| |
54 | 66 |
| |
|
Lines changed: 44 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
| 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 | + |
Lines changed: 21 additions & 4 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
20 | 20 |
| |
21 | 21 |
| |
22 | 22 |
| |
| 23 | + | |
23 | 24 |
| |
24 | 25 |
| |
25 | 26 |
| |
| |||
56 | 57 |
| |
57 | 58 |
| |
58 | 59 |
| |
| 60 | + | |
| 61 | + | |
| 62 | + | |
59 | 63 |
| |
60 | 64 |
| |
61 | 65 |
| |
| |||
65 | 69 |
| |
66 | 70 |
| |
67 | 71 |
| |
68 |
| - | |
| 72 | + | |
| 73 | + | |
| 74 | + | |
69 | 75 |
| |
70 | 76 |
| |
71 | 77 |
| |
| |||
84 | 90 |
| |
85 | 91 |
| |
86 | 92 |
| |
| 93 | + | |
| 94 | + | |
87 | 95 |
| |
88 | 96 |
| |
89 | 97 |
| |
| |||
110 | 118 |
| |
111 | 119 |
| |
112 | 120 |
| |
113 |
| - | |
| 121 | + | |
| 122 | + | |
| 123 | + | |
| 124 | + | |
114 | 125 |
| |
115 | 126 |
| |
116 | 127 |
| |
117 | 128 |
| |
118 | 129 |
| |
119 |
| - | |
| 130 | + | |
| 131 | + | |
| 132 | + | |
| 133 | + | |
120 | 134 |
| |
121 | 135 |
| |
122 | 136 |
| |
| |||
127 | 141 |
| |
128 | 142 |
| |
129 | 143 |
| |
130 |
| - | |
| 144 | + | |
| 145 | + | |
| 146 | + | |
| 147 | + | |
131 | 148 |
| |
132 | 149 |
| |
133 | 150 |
| |
|
Lines changed: 14 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
194 | 194 |
| |
195 | 195 |
| |
196 | 196 |
| |
| 197 | + | |
| 198 | + | |
| 199 | + | |
| 200 | + | |
| 201 | + | |
| 202 | + | |
| 203 | + | |
| 204 | + | |
| 205 | + | |
| 206 | + | |
| 207 | + | |
| 208 | + | |
| 209 | + | |
| 210 | + | |
197 | 211 |
| |
198 | 212 |
| |
199 | 213 |
| |
|
Lines changed: 56 additions & 3 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
171 | 171 |
| |
172 | 172 |
| |
173 | 173 |
| |
| 174 | + | |
| 175 | + | |
| 176 | + | |
174 | 177 |
| |
175 | 178 |
| |
176 | 179 |
| |
| |||
189 | 192 |
| |
190 | 193 |
| |
191 | 194 |
| |
| 195 | + | |
| 196 | + | |
| 197 | + | |
| 198 | + | |
| 199 | + | |
| 200 | + | |
| 201 | + | |
| 202 | + | |
| 203 | + | |
| 204 | + | |
| 205 | + | |
| 206 | + | |
| 207 | + | |
| 208 | + | |
| 209 | + | |
| 210 | + | |
| 211 | + | |
| 212 | + | |
| 213 | + | |
| 214 | + | |
| 215 | + | |
| 216 | + | |
| 217 | + | |
| 218 | + | |
| 219 | + | |
| 220 | + | |
| 221 | + | |
| 222 | + | |
| 223 | + | |
| 224 | + | |
| 225 | + | |
| 226 | + | |
| 227 | + | |
| 228 | + | |
| 229 | + | |
| 230 | + | |
| 231 | + | |
| 232 | + | |
| 233 | + | |
| 234 | + | |
| 235 | + | |
| 236 | + | |
192 | 237 |
| |
193 | 238 |
| |
194 | 239 |
| |
| |||
228 | 273 |
| |
229 | 274 |
| |
230 | 275 |
| |
| 276 | + | |
| 277 | + | |
231 | 278 |
| |
232 | 279 |
| |
233 | 280 |
| |
| |||
344 | 391 |
| |
345 | 392 |
| |
346 | 393 |
| |
347 |
| - | |
| 394 | + | |
| 395 | + | |
| 396 | + | |
348 | 397 |
| |
349 | 398 |
| |
350 | 399 |
| |
| |||
384 | 433 |
| |
385 | 434 |
| |
386 | 435 |
| |
387 |
| - | |
| 436 | + | |
| 437 | + | |
| 438 | + | |
388 | 439 |
| |
389 | 440 |
| |
390 | 441 |
| |
| |||
550 | 601 |
| |
551 | 602 |
| |
552 | 603 |
| |
553 |
| - | |
| 604 | + | |
| 605 | + | |
| 606 | + | |
554 | 607 |
| |
555 | 608 |
| |
556 | 609 |
| |
|
Lines changed: 23 additions & 3 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
71 | 71 |
| |
72 | 72 |
| |
73 | 73 |
| |
| 74 | + | |
| 75 | + | |
74 | 76 |
| |
75 | 77 |
| |
76 | 78 |
| |
| |||
152 | 154 |
| |
153 | 155 |
| |
154 | 156 |
| |
| 157 | + | |
| 158 | + | |
155 | 159 |
| |
156 | 160 |
| |
157 | 161 |
| |
| |||
215 | 219 |
| |
216 | 220 |
| |
217 | 221 |
| |
| 222 | + | |
| 223 | + | |
| 224 | + | |
| 225 | + | |
| 226 | + | |
| 227 | + | |
| 228 | + | |
| 229 | + | |
| 230 | + | |
| 231 | + | |
218 | 232 |
| |
219 | 233 |
| |
220 | 234 |
| |
| |||
231 | 245 |
| |
232 | 246 |
| |
233 | 247 |
| |
234 |
| - | |
| 248 | + | |
| 249 | + | |
| 250 | + | |
235 | 251 |
| |
236 | 252 |
| |
237 | 253 |
| |
| |||
395 | 411 |
| |
396 | 412 |
| |
397 | 413 |
| |
398 |
| - | |
| 414 | + | |
| 415 | + | |
| 416 | + | |
399 | 417 |
| |
400 | 418 |
| |
401 | 419 |
| |
| |||
413 | 431 |
| |
414 | 432 |
| |
415 | 433 |
| |
416 |
| - | |
| 434 | + | |
| 435 | + | |
| 436 | + | |
417 | 437 |
| |
418 | 438 |
| |
419 | 439 |
| |
|
0 commit comments