File tree
5 files changed
+72
-8
lines changed- flink-connector-kafka/src
- main/java/org/apache/flink/connector/kafka/sink
- test/java/org/apache/flink/connector/kafka/sink
5 files changed
+72
-8
lines changedLines changed: 9 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
54 | 54 |
| |
55 | 55 |
| |
56 | 56 |
| |
| 57 | + | |
57 | 58 |
| |
58 | 59 |
| |
59 | 60 |
| |
| |||
77 | 78 |
| |
78 | 79 |
| |
79 | 80 |
| |
| 81 | + | |
80 | 82 |
| |
81 | 83 |
| |
82 | 84 |
| |
| |||
86 | 88 |
| |
87 | 89 |
| |
88 | 90 |
| |
| 91 | + | |
89 | 92 |
| |
90 | 93 |
| |
91 | 94 |
| |
| |||
120 | 123 |
| |
121 | 124 |
| |
122 | 125 |
| |
| 126 | + | |
| 127 | + | |
| 128 | + | |
| 129 | + | |
123 | 130 |
| |
124 | 131 |
| |
125 | 132 |
| |
| |||
133 | 140 |
| |
134 | 141 |
| |
135 | 142 |
| |
| 143 | + | |
136 | 144 |
| |
137 | 145 |
| |
138 | 146 |
| |
139 | 147 |
| |
140 | 148 |
| |
141 | 149 |
| |
| 150 | + | |
142 | 151 |
| |
143 | 152 |
| |
144 | 153 |
| |
|
Lines changed: 27 additions & 1 deletion
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
32 | 32 |
| |
33 | 33 |
| |
34 | 34 |
| |
| 35 | + | |
35 | 36 |
| |
36 | 37 |
| |
37 | 38 |
| |
38 | 39 |
| |
39 | 40 |
| |
40 | 41 |
| |
| 42 | + | |
41 | 43 |
| |
42 | 44 |
| |
43 | 45 |
| |
44 | 46 |
| |
| 47 | + | |
45 | 48 |
| |
46 | 49 |
| |
47 | 50 |
| |
| 51 | + | |
| 52 | + | |
| 53 | + | |
| 54 | + | |
| 55 | + | |
| 56 | + | |
| 57 | + | |
| 58 | + | |
48 | 59 |
| |
49 | 60 |
| |
50 | 61 |
| |
| 62 | + | |
| 63 | + | |
| 64 | + | |
| 65 | + | |
| 66 | + | |
| 67 | + | |
| 68 | + | |
51 | 69 |
| |
52 | 70 |
| |
53 | 71 |
| |
54 | 72 |
| |
| 73 | + | |
55 | 74 |
| |
56 | 75 |
| |
57 | 76 |
| |
| |||
67 | 86 |
| |
68 | 87 |
| |
69 | 88 |
| |
| 89 | + | |
| 90 | + | |
| 91 | + | |
| 92 | + | |
70 | 93 |
| |
71 | 94 |
| |
72 | 95 |
| |
| |||
76 | 99 |
| |
77 | 100 |
| |
78 | 101 |
| |
| 102 | + | |
| 103 | + | |
79 | 104 |
| |
80 | 105 |
| |
81 | 106 |
| |
| |||
94 | 119 |
| |
95 | 120 |
| |
96 | 121 |
| |
| 122 | + | |
97 | 123 |
| |
98 | 124 |
| |
99 | 125 |
| |
100 | 126 |
| |
101 | 127 |
| |
102 |
| - | |
| 128 | + | |
103 | 129 |
| |
104 | 130 |
|
Lines changed: 3 additions & 1 deletion
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
70 | 70 |
| |
71 | 71 |
| |
72 | 72 |
| |
73 |
| - | |
| 73 | + | |
| 74 | + | |
| 75 | + | |
74 | 76 |
| |
75 | 77 |
| |
76 | 78 |
| |
|
Lines changed: 13 additions & 6 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
209 | 209 |
| |
210 | 210 |
| |
211 | 211 |
| |
212 |
| - | |
| 212 | + | |
213 | 213 |
| |
214 | 214 |
| |
215 | 215 |
| |
216 | 216 |
| |
217 |
| - | |
| 217 | + | |
| 218 | + | |
218 | 219 |
| |
219 | 220 |
| |
220 |
| - | |
| 221 | + | |
| 222 | + | |
| 223 | + | |
| 224 | + | |
221 | 225 |
| |
222 | 226 |
| |
223 | 227 |
| |
224 | 228 |
| |
225 |
| - | |
226 |
| - | |
227 |
| - | |
| 229 | + | |
| 230 | + | |
| 231 | + | |
| 232 | + | |
| 233 | + | |
| 234 | + | |
228 | 235 |
| |
229 | 236 |
| |
230 | 237 |
| |
|
Lines changed: 20 additions & 0 deletions
Original file line number | Diff line number | Diff line change | |
---|---|---|---|
| |||
507 | 507 |
| |
508 | 508 |
| |
509 | 509 |
| |
| 510 | + | |
| 511 | + | |
| 512 | + | |
| 513 | + | |
| 514 | + | |
| 515 | + | |
| 516 | + | |
| 517 | + | |
| 518 | + | |
| 519 | + | |
| 520 | + | |
| 521 | + | |
| 522 | + | |
| 523 | + | |
| 524 | + | |
| 525 | + | |
| 526 | + | |
| 527 | + | |
| 528 | + | |
| 529 | + | |
510 | 530 |
| |
511 | 531 |
| |
512 | 532 |
| |
|
0 commit comments