File tree
9 files changed
+253
-25
lines changed- flink-connector-kafka/src
- main/java/org/apache/flink/streaming/connectors/kafka/table
- test/java/org/apache/flink
- connector/kafka/sink
- streaming/connectors/kafka/table
9 files changed
+253
-25
lines changedLines changed: 51 additions & 1 deletion
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
19 | 19 |
| |
20 | 20 |
| |
21 | 21 |
| |
| 22 | + | |
| 23 | + | |
| 24 | + | |
| 25 | + | |
| 26 | + | |
| 27 | + | |
| 28 | + | |
| 29 | + | |
| 30 | + | |
22 | 31 |
| |
23 | 32 |
| |
24 | 33 |
| |
| |||
27 | 36 |
| |
28 | 37 |
| |
29 | 38 |
| |
| 39 | + | |
30 | 40 |
| |
31 | 41 |
| |
32 | 42 |
| |
33 | 43 |
| |
| 44 | + | |
34 | 45 |
| |
35 | 46 |
| |
36 | 47 |
| |
37 | 48 |
| |
| 49 | + | |
38 | 50 |
| |
39 | 51 |
| |
40 | 52 |
| |
41 | 53 |
| |
42 | 54 |
| |
43 | 55 |
| |
44 | 56 |
| |
45 |
| - | |
| 57 | + | |
| 58 | + | |
| 59 | + | |
| 60 | + | |
46 | 61 |
| |
47 | 62 |
| |
48 | 63 |
| |
| |||
170 | 185 |
| |
171 | 186 |
| |
172 | 187 |
| |
| 188 | + | |
| 189 | + | |
| 190 | + | |
| 191 | + | |
| 192 | + | |
| 193 | + | |
| 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 | + | |
173 | 223 |
| |
174 | 224 |
| |
175 | 225 |
| |
|
Lines changed: 25 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
25 | 25 |
| |
26 | 26 |
| |
27 | 27 |
| |
| 28 | + | |
28 | 29 |
| |
29 | 30 |
| |
30 | 31 |
| |
| |||
272 | 273 |
| |
273 | 274 |
| |
274 | 275 |
| |
| 276 | + | |
| 277 | + | |
| 278 | + | |
| 279 | + | |
| 280 | + | |
| 281 | + | |
| 282 | + | |
| 283 | + | |
| 284 | + | |
| 285 | + | |
| 286 | + | |
| 287 | + | |
| 288 | + | |
| 289 | + | |
| 290 | + | |
| 291 | + | |
| 292 | + | |
| 293 | + | |
| 294 | + | |
| 295 | + | |
| 296 | + | |
| 297 | + | |
| 298 | + | |
| 299 | + | |
275 | 300 |
| |
276 | 301 |
| |
277 | 302 |
| |
|
Lines changed: 9 additions & 2 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
27 | 27 |
| |
28 | 28 |
| |
29 | 29 |
| |
| 30 | + | |
30 | 31 |
| |
31 | 32 |
| |
32 | 33 |
| |
| |||
139 | 140 |
| |
140 | 141 |
| |
141 | 142 |
| |
| 143 | + | |
| 144 | + | |
142 | 145 |
| |
143 | 146 |
| |
144 | 147 |
| |
| |||
155 | 158 |
| |
156 | 159 |
| |
157 | 160 |
| |
158 |
| - | |
| 161 | + | |
| 162 | + | |
159 | 163 |
| |
160 | 164 |
| |
161 | 165 |
| |
| |||
168 | 172 |
| |
169 | 173 |
| |
170 | 174 |
| |
| 175 | + | |
171 | 176 |
| |
172 | 177 |
| |
173 | 178 |
| |
| |||
222 | 227 |
| |
223 | 228 |
| |
224 | 229 |
| |
| 230 | + | |
225 | 231 |
| |
226 | 232 |
| |
227 | 233 |
| |
| |||
292 | 298 |
| |
293 | 299 |
| |
294 | 300 |
| |
295 |
| - | |
| 301 | + | |
| 302 | + | |
296 | 303 |
| |
297 | 304 |
| |
298 | 305 |
| |
|
Lines changed: 11 additions & 4 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
27 | 27 |
| |
28 | 28 |
| |
29 | 29 |
| |
| 30 | + | |
30 | 31 |
| |
31 | 32 |
| |
32 | 33 |
| |
| |||
84 | 85 |
| |
85 | 86 |
| |
86 | 87 |
| |
| 88 | + | |
87 | 89 |
| |
88 | 90 |
| |
89 | 91 |
| |
| |||
154 | 156 |
| |
155 | 157 |
| |
156 | 158 |
| |
| 159 | + | |
157 | 160 |
| |
158 | 161 |
| |
159 | 162 |
| |
| |||
171 | 174 |
| |
172 | 175 |
| |
173 | 176 |
| |
174 |
| - | |
| 177 | + | |
| 178 | + | |
175 | 179 |
| |
176 | 180 |
| |
177 | 181 |
| |
| |||
290 | 294 |
| |
291 | 295 |
| |
292 | 296 |
| |
293 |
| - | |
| 297 | + | |
| 298 | + | |
294 | 299 |
| |
295 | 300 |
| |
296 | 301 |
| |
| |||
438 | 443 |
| |
439 | 444 |
| |
440 | 445 |
| |
441 |
| - | |
| 446 | + | |
| 447 | + | |
442 | 448 |
| |
443 | 449 |
| |
444 | 450 |
| |
| |||
455 | 461 |
| |
456 | 462 |
| |
457 | 463 |
| |
458 |
| - | |
| 464 | + | |
| 465 | + | |
459 | 466 |
| |
460 | 467 |
|
Lines changed: 6 additions & 2 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
69 | 69 |
| |
70 | 70 |
| |
71 | 71 |
| |
| 72 | + | |
72 | 73 |
| |
73 | 74 |
| |
74 | 75 |
| |
| |||
117 | 118 |
| |
118 | 119 |
| |
119 | 120 |
| |
| 121 | + | |
120 | 122 |
| |
121 | 123 |
| |
122 | 124 |
| |
123 | 125 |
| |
124 | 126 |
| |
125 |
| - | |
| 127 | + | |
| 128 | + | |
126 | 129 |
| |
127 | 130 |
| |
128 | 131 |
| |
| |||
227 | 230 |
| |
228 | 231 |
| |
229 | 232 |
| |
230 |
| - | |
| 233 | + | |
| 234 | + | |
231 | 235 |
| |
232 | 236 |
| |
233 | 237 |
| |
|
Lines changed: 2 additions & 1 deletion
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
514 | 514 |
| |
515 | 515 |
| |
516 | 516 |
| |
517 |
| - | |
| 517 | + | |
| 518 | + | |
518 | 519 |
| |
519 | 520 |
| |
520 | 521 |
| |
|
0 commit comments