File tree
5 files changed
+163
-20
lines changed- flink-connector-kafka/src
- main/java/org/apache/flink/connector/kafka/source/reader
- test/java/org/apache/flink
- connector/kafka/source/reader
- streaming/connectors/kafka/table
5 files changed
+163
-20
lines changedLines changed: 0 additions & 1 deletion
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
41 | 41 |
| |
42 | 42 |
| |
43 | 43 |
| |
44 |
| - | |
45 | 44 |
| |
46 | 45 |
| |
47 | 46 |
| |
|
Lines changed: 10 additions & 10 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
540 | 540 |
| |
541 | 541 |
| |
542 | 542 |
| |
543 |
| - | |
544 |
| - | |
545 |
| - | |
546 |
| - | |
547 |
| - | |
548 |
| - | |
549 |
| - | |
550 |
| - | |
551 |
| - | |
| 543 | + | |
| 544 | + | |
| 545 | + | |
| 546 | + | |
| 547 | + | |
| 548 | + | |
| 549 | + | |
| 550 | + | |
552 | 551 |
| |
553 | 552 |
| |
554 | 553 |
| |
| |||
728 | 727 |
| |
729 | 728 |
| |
730 | 729 |
| |
731 |
| - | |
| 730 | + | |
| 731 | + | |
732 | 732 |
| |
733 | 733 |
| |
734 | 734 |
| |
|
Lines changed: 60 additions & 8 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
26 | 26 |
| |
27 | 27 |
| |
28 | 28 |
| |
| 29 | + | |
29 | 30 |
| |
30 | 31 |
| |
31 | 32 |
| |
| |||
47 | 48 |
| |
48 | 49 |
| |
49 | 50 |
| |
| 51 | + | |
50 | 52 |
| |
51 | 53 |
| |
52 | 54 |
| |
| |||
91 | 93 |
| |
92 | 94 |
| |
93 | 95 |
| |
| 96 | + | |
94 | 97 |
| |
95 | 98 |
| |
96 | 99 |
| |
| |||
212 | 215 |
| |
213 | 216 |
| |
214 | 217 |
| |
215 |
| - | |
| 218 | + | |
| 219 | + | |
216 | 220 |
| |
217 | 221 |
| |
218 | 222 |
| |
| |||
254 | 258 |
| |
255 | 259 |
| |
256 | 260 |
| |
257 |
| - | |
| 261 | + | |
| 262 | + | |
258 | 263 |
| |
259 | 264 |
| |
260 | 265 |
| |
| |||
295 | 300 |
| |
296 | 301 |
| |
297 | 302 |
| |
298 |
| - | |
| 303 | + | |
| 304 | + | |
299 | 305 |
| |
300 | 306 |
| |
301 | 307 |
| |
| |||
346 | 352 |
| |
347 | 353 |
| |
348 | 354 |
| |
349 |
| - | |
| 355 | + | |
| 356 | + | |
350 | 357 |
| |
351 | 358 |
| |
352 | 359 |
| |
| |||
397 | 404 |
| |
398 | 405 |
| |
399 | 406 |
| |
| 407 | + | |
| 408 | + | |
| 409 | + | |
| 410 | + | |
| 411 | + | |
| 412 | + | |
| 413 | + | |
| 414 | + | |
| 415 | + | |
| 416 | + | |
| 417 | + | |
| 418 | + | |
| 419 | + | |
| 420 | + | |
| 421 | + | |
| 422 | + | |
| 423 | + | |
| 424 | + | |
| 425 | + | |
| 426 | + | |
| 427 | + | |
| 428 | + | |
| 429 | + | |
| 430 | + | |
| 431 | + | |
| 432 | + | |
| 433 | + | |
| 434 | + | |
| 435 | + | |
| 436 | + | |
| 437 | + | |
| 438 | + | |
| 439 | + | |
| 440 | + | |
| 441 | + | |
| 442 | + | |
| 443 | + | |
| 444 | + | |
| 445 | + | |
| 446 | + | |
| 447 | + | |
400 | 448 |
| |
401 | 449 |
| |
402 | 450 |
| |
| |||
1128 | 1176 |
| |
1129 | 1177 |
| |
1130 | 1178 |
| |
1131 |
| - | |
| 1179 | + | |
| 1180 | + | |
1132 | 1181 |
| |
1133 | 1182 |
| |
1134 | 1183 |
| |
| |||
1166 | 1215 |
| |
1167 | 1216 |
| |
1168 | 1217 |
| |
1169 |
| - | |
| 1218 | + | |
| 1219 | + | |
1170 | 1220 |
| |
1171 | 1221 |
| |
1172 | 1222 |
| |
| |||
1189 | 1239 |
| |
1190 | 1240 |
| |
1191 | 1241 |
| |
1192 |
| - | |
| 1242 | + | |
| 1243 | + | |
1193 | 1244 |
| |
1194 | 1245 |
| |
1195 | 1246 |
| |
| |||
1207 | 1258 |
| |
1208 | 1259 |
| |
1209 | 1260 |
| |
1210 |
| - | |
| 1261 | + | |
| 1262 | + | |
1211 | 1263 |
| |
1212 | 1264 |
| |
1213 | 1265 |
| |
|
Lines changed: 91 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
29 | 29 |
| |
30 | 30 |
| |
31 | 31 |
| |
| 32 | + | |
32 | 33 |
| |
33 | 34 |
| |
34 | 35 |
| |
| |||
45 | 46 |
| |
46 | 47 |
| |
47 | 48 |
| |
| 49 | + | |
| 50 | + | |
48 | 51 |
| |
49 | 52 |
| |
50 | 53 |
| |
| |||
78 | 81 |
| |
79 | 82 |
| |
80 | 83 |
| |
| 84 | + | |
81 | 85 |
| |
82 | 86 |
| |
83 | 87 |
| |
| |||
1099 | 1103 |
| |
1100 | 1104 |
| |
1101 | 1105 |
| |
| 1106 | + | |
| 1107 | + | |
| 1108 | + | |
| 1109 | + | |
| 1110 | + | |
| 1111 | + | |
| 1112 | + | |
| 1113 | + | |
| 1114 | + | |
| 1115 | + | |
| 1116 | + | |
| 1117 | + | |
| 1118 | + | |
| 1119 | + | |
| 1120 | + | |
| 1121 | + | |
| 1122 | + | |
| 1123 | + | |
| 1124 | + | |
| 1125 | + | |
| 1126 | + | |
| 1127 | + | |
| 1128 | + | |
| 1129 | + | |
| 1130 | + | |
| 1131 | + | |
| 1132 | + | |
| 1133 | + | |
| 1134 | + | |
| 1135 | + | |
| 1136 | + | |
| 1137 | + | |
| 1138 | + | |
| 1139 | + | |
| 1140 | + | |
| 1141 | + | |
| 1142 | + | |
| 1143 | + | |
| 1144 | + | |
| 1145 | + | |
| 1146 | + | |
| 1147 | + | |
| 1148 | + | |
| 1149 | + | |
| 1150 | + | |
| 1151 | + | |
| 1152 | + | |
| 1153 | + | |
| 1154 | + | |
| 1155 | + | |
| 1156 | + | |
| 1157 | + | |
| 1158 | + | |
| 1159 | + | |
| 1160 | + | |
| 1161 | + | |
| 1162 | + | |
| 1163 | + | |
| 1164 | + | |
| 1165 | + | |
| 1166 | + | |
| 1167 | + | |
| 1168 | + | |
| 1169 | + | |
| 1170 | + | |
| 1171 | + | |
| 1172 | + | |
| 1173 | + | |
| 1174 | + | |
| 1175 | + | |
| 1176 | + | |
| 1177 | + | |
| 1178 | + | |
| 1179 | + | |
| 1180 | + | |
| 1181 | + | |
| 1182 | + | |
| 1183 | + | |
| 1184 | + | |
| 1185 | + | |
| 1186 | + | |
| 1187 | + | |
| 1188 | + | |
| 1189 | + | |
| 1190 | + | |
| 1191 | + | |
| 1192 | + | |
1102 | 1193 |
| |
1103 | 1194 |
| |
1104 | 1195 |
| |
|
Lines changed: 2 additions & 1 deletion
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
791 | 791 |
| |
792 | 792 |
| |
793 | 793 |
| |
794 |
| - | |
| 794 | + | |
| 795 | + | |
795 | 796 |
| |
796 | 797 |
| |
797 | 798 |
| |
|
0 commit comments