File tree
4 files changed
+85
-7
lines changed- flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source
- reader
4 files changed
+85
-7
lines changedLines changed: 14 additions & 2 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
31 | 31 |
| |
32 | 32 |
| |
33 | 33 |
| |
| 34 | + | |
34 | 35 |
| |
35 | 36 |
| |
36 | 37 |
| |
| |||
102 | 103 |
| |
103 | 104 |
| |
104 | 105 |
| |
| 106 | + | |
105 | 107 |
| |
106 | 108 |
| |
107 | 109 |
| |
108 | 110 |
| |
109 | 111 |
| |
| 112 | + | |
110 | 113 |
| |
111 | 114 |
| |
112 | 115 |
| |
| |||
118 | 121 |
| |
119 | 122 |
| |
120 | 123 |
| |
| 124 | + | |
121 | 125 |
| |
122 | 126 |
| |
123 | 127 |
| |
| |||
171 | 175 |
| |
172 | 176 |
| |
173 | 177 |
| |
174 |
| - | |
| 178 | + | |
| 179 | + | |
175 | 180 |
| |
176 | 181 |
| |
177 | 182 |
| |
| |||
180 | 185 |
| |
181 | 186 |
| |
182 | 187 |
| |
183 |
| - | |
| 188 | + | |
| 189 | + | |
184 | 190 |
| |
185 | 191 |
| |
186 | 192 |
| |
| |||
251 | 257 |
| |
252 | 258 |
| |
253 | 259 |
| |
| 260 | + | |
| 261 | + | |
| 262 | + | |
| 263 | + | |
| 264 | + | |
| 265 | + | |
254 | 266 |
|
flink-connector-kafka/src/main/java/org/apache/flink/connector/kafka/source/KafkaSourceBuilder.java
Lines changed: 23 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
21 | 21 |
| |
22 | 22 |
| |
23 | 23 |
| |
| 24 | + | |
24 | 25 |
| |
25 | 26 |
| |
26 | 27 |
| |
| |||
104 | 105 |
| |
105 | 106 |
| |
106 | 107 |
| |
| 108 | + | |
107 | 109 |
| |
108 | 110 |
| |
109 | 111 |
| |
| |||
350 | 352 |
| |
351 | 353 |
| |
352 | 354 |
| |
| 355 | + | |
| 356 | + | |
| 357 | + | |
| 358 | + | |
| 359 | + | |
| 360 | + | |
| 361 | + | |
| 362 | + | |
| 363 | + | |
| 364 | + | |
| 365 | + | |
| 366 | + | |
| 367 | + | |
| 368 | + | |
| 369 | + | |
| 370 | + | |
| 371 | + | |
| 372 | + | |
| 373 | + | |
| 374 | + | |
353 | 375 |
| |
354 | 376 |
| |
355 | 377 |
| |
| |||
436 | 458 |
| |
437 | 459 |
| |
438 | 460 |
| |
| 461 | + | |
439 | 462 |
| |
440 | 463 |
| |
441 | 464 |
| |
|
Lines changed: 35 additions & 3 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 |
| |
| 31 | + | |
| 32 | + | |
30 | 33 |
| |
| 34 | + | |
| 35 | + | |
31 | 36 |
| |
32 | 37 |
| |
33 | 38 |
| |
34 | 39 |
| |
35 | 40 |
| |
36 | 41 |
| |
37 | 42 |
| |
38 |
| - | |
| 43 | + | |
| 44 | + | |
| 45 | + | |
39 | 46 |
| |
40 |
| - | |
| 47 | + | |
| 48 | + | |
| 49 | + | |
41 | 50 |
| |
| 51 | + | |
| 52 | + | |
| 53 | + | |
42 | 54 |
| |
43 | 55 |
| |
44 | 56 |
| |
| |||
51 | 63 |
| |
52 | 64 |
| |
53 | 65 |
| |
54 |
| - | |
| 66 | + | |
| 67 | + | |
| 68 | + | |
| 69 | + | |
| 70 | + | |
| 71 | + | |
| 72 | + | |
55 | 73 |
| |
56 | 74 |
| |
57 | 75 |
| |
58 | 76 |
| |
59 | 77 |
| |
60 | 78 |
| |
| 79 | + | |
61 | 80 |
| |
62 | 81 |
| |
63 | 82 |
| |
| 83 | + | |
| 84 | + | |
| 85 | + | |
| 86 | + | |
| 87 | + | |
64 | 88 |
| |
65 | 89 |
| |
66 | 90 |
| |
67 | 91 |
| |
| 92 | + | |
| 93 | + | |
| 94 | + | |
68 | 95 |
| |
69 | 96 |
| |
70 | 97 |
| |
| |||
77 | 104 |
| |
78 | 105 |
| |
79 | 106 |
| |
| 107 | + | |
| 108 | + | |
| 109 | + | |
| 110 | + | |
| 111 | + | |
80 | 112 |
| |
81 | 113 |
|
Lines changed: 13 additions & 2 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
23 | 23 |
| |
24 | 24 |
| |
25 | 25 |
| |
| 26 | + | |
26 | 27 |
| |
27 | 28 |
| |
28 | 29 |
| |
| |||
38 | 39 |
| |
39 | 40 |
| |
40 | 41 |
| |
| 42 | + | |
| 43 | + | |
| 44 | + | |
41 | 45 |
| |
42 | 46 |
| |
43 | 47 |
| |
| |||
68 | 72 |
| |
69 | 73 |
| |
70 | 74 |
| |
71 |
| - | |
72 |
| - | |
| 75 | + | |
| 76 | + | |
| 77 | + | |
| 78 | + | |
| 79 | + | |
| 80 | + | |
| 81 | + | |
| 82 | + | |
| 83 | + | |
73 | 84 |
| |
74 | 85 |
| |
75 | 86 |
| |
|
0 commit comments