|
24 | 24 | import org.apache.fluss.record.LogRecord; |
25 | 25 | import org.apache.fluss.row.BinaryString; |
26 | 26 | import org.apache.fluss.row.Decimal; |
| 27 | +import org.apache.fluss.row.GenericArray; |
27 | 28 | import org.apache.fluss.row.GenericRow; |
28 | 29 | import org.apache.fluss.row.TimestampLtz; |
29 | 30 | import org.apache.fluss.row.TimestampNtz; |
30 | 31 |
|
| 32 | +import org.apache.paimon.data.InternalArray; |
31 | 33 | import org.apache.paimon.types.RowKind; |
32 | 34 | import org.apache.paimon.types.RowType; |
33 | 35 | import org.junit.jupiter.api.Test; |
@@ -140,4 +142,171 @@ void testPrimaryKeyTableRecord() { |
140 | 142 | assertThat(new FlussRowAsPaimonRow(logRecord.getRow(), tableRowType).getRowKind()) |
141 | 143 | .isEqualTo(RowKind.INSERT); |
142 | 144 | } |
| 145 | + |
| 146 | + @Test |
| 147 | + void testArrayTypeWithIntElements() { |
| 148 | + RowType tableRowType = |
| 149 | + RowType.of( |
| 150 | + new org.apache.paimon.types.IntType(), |
| 151 | + new org.apache.paimon.types.ArrayType( |
| 152 | + new org.apache.paimon.types.IntType())); |
| 153 | + |
| 154 | + long logOffset = 0; |
| 155 | + long timeStamp = System.currentTimeMillis(); |
| 156 | + GenericRow genericRow = new GenericRow(2); |
| 157 | + genericRow.setField(0, 42); |
| 158 | + genericRow.setField(1, new GenericArray(new int[] {1, 2, 3, 4, 5})); |
| 159 | + |
| 160 | + LogRecord logRecord = new GenericRecord(logOffset, timeStamp, APPEND_ONLY, genericRow); |
| 161 | + FlussRowAsPaimonRow flussRowAsPaimonRow = |
| 162 | + new FlussRowAsPaimonRow(logRecord.getRow(), tableRowType); |
| 163 | + |
| 164 | + assertThat(flussRowAsPaimonRow.getInt(0)).isEqualTo(42); |
| 165 | + InternalArray array = flussRowAsPaimonRow.getArray(1); |
| 166 | + assertThat(array).isNotNull(); |
| 167 | + assertThat(array.size()).isEqualTo(5); |
| 168 | + assertThat(array.getInt(0)).isEqualTo(1); |
| 169 | + assertThat(array.getInt(1)).isEqualTo(2); |
| 170 | + assertThat(array.getInt(2)).isEqualTo(3); |
| 171 | + assertThat(array.getInt(3)).isEqualTo(4); |
| 172 | + assertThat(array.getInt(4)).isEqualTo(5); |
| 173 | + } |
| 174 | + |
| 175 | + @Test |
| 176 | + void testArrayTypeWithStringElements() { |
| 177 | + RowType tableRowType = |
| 178 | + RowType.of( |
| 179 | + new org.apache.paimon.types.VarCharType(), |
| 180 | + new org.apache.paimon.types.ArrayType( |
| 181 | + new org.apache.paimon.types.VarCharType())); |
| 182 | + |
| 183 | + long logOffset = 0; |
| 184 | + long timeStamp = System.currentTimeMillis(); |
| 185 | + GenericRow genericRow = new GenericRow(2); |
| 186 | + genericRow.setField(0, BinaryString.fromString("name")); |
| 187 | + genericRow.setField( |
| 188 | + 1, |
| 189 | + new GenericArray( |
| 190 | + new Object[] { |
| 191 | + BinaryString.fromString("a"), |
| 192 | + BinaryString.fromString("b"), |
| 193 | + BinaryString.fromString("c") |
| 194 | + })); |
| 195 | + |
| 196 | + LogRecord logRecord = new GenericRecord(logOffset, timeStamp, APPEND_ONLY, genericRow); |
| 197 | + FlussRowAsPaimonRow flussRowAsPaimonRow = |
| 198 | + new FlussRowAsPaimonRow(logRecord.getRow(), tableRowType); |
| 199 | + |
| 200 | + assertThat(flussRowAsPaimonRow.getString(0).toString()).isEqualTo("name"); |
| 201 | + InternalArray array = flussRowAsPaimonRow.getArray(1); |
| 202 | + assertThat(array).isNotNull(); |
| 203 | + assertThat(array.size()).isEqualTo(3); |
| 204 | + assertThat(array.getString(0).toString()).isEqualTo("a"); |
| 205 | + assertThat(array.getString(1).toString()).isEqualTo("b"); |
| 206 | + assertThat(array.getString(2).toString()).isEqualTo("c"); |
| 207 | + } |
| 208 | + |
| 209 | + @Test |
| 210 | + void testArrayTypeWithNullableElements() { |
| 211 | + RowType tableRowType = |
| 212 | + RowType.of( |
| 213 | + new org.apache.paimon.types.ArrayType( |
| 214 | + new org.apache.paimon.types.IntType().nullable())); |
| 215 | + |
| 216 | + long logOffset = 0; |
| 217 | + long timeStamp = System.currentTimeMillis(); |
| 218 | + GenericRow genericRow = new GenericRow(1); |
| 219 | + genericRow.setField(0, new GenericArray(new Object[] {1, null, 3})); |
| 220 | + |
| 221 | + LogRecord logRecord = new GenericRecord(logOffset, timeStamp, APPEND_ONLY, genericRow); |
| 222 | + FlussRowAsPaimonRow flussRowAsPaimonRow = |
| 223 | + new FlussRowAsPaimonRow(logRecord.getRow(), tableRowType); |
| 224 | + |
| 225 | + InternalArray array = flussRowAsPaimonRow.getArray(0); |
| 226 | + assertThat(array).isNotNull(); |
| 227 | + assertThat(array.size()).isEqualTo(3); |
| 228 | + assertThat(array.getInt(0)).isEqualTo(1); |
| 229 | + assertThat(array.isNullAt(1)).isTrue(); |
| 230 | + assertThat(array.getInt(2)).isEqualTo(3); |
| 231 | + } |
| 232 | + |
| 233 | + @Test |
| 234 | + void testNullArray() { |
| 235 | + RowType tableRowType = |
| 236 | + RowType.of( |
| 237 | + new org.apache.paimon.types.ArrayType(new org.apache.paimon.types.IntType()) |
| 238 | + .nullable()); |
| 239 | + |
| 240 | + long logOffset = 0; |
| 241 | + long timeStamp = System.currentTimeMillis(); |
| 242 | + GenericRow genericRow = new GenericRow(1); |
| 243 | + genericRow.setField(0, null); |
| 244 | + |
| 245 | + LogRecord logRecord = new GenericRecord(logOffset, timeStamp, APPEND_ONLY, genericRow); |
| 246 | + FlussRowAsPaimonRow flussRowAsPaimonRow = |
| 247 | + new FlussRowAsPaimonRow(logRecord.getRow(), tableRowType); |
| 248 | + |
| 249 | + assertThat(flussRowAsPaimonRow.isNullAt(0)).isTrue(); |
| 250 | + } |
| 251 | + |
| 252 | + @Test |
| 253 | + void testNestedArrayType() { |
| 254 | + // Test ARRAY<ARRAY<INT>> |
| 255 | + RowType tableRowType = |
| 256 | + RowType.of( |
| 257 | + new org.apache.paimon.types.ArrayType( |
| 258 | + new org.apache.paimon.types.ArrayType( |
| 259 | + new org.apache.paimon.types.IntType()))); |
| 260 | + |
| 261 | + long logOffset = 0; |
| 262 | + long timeStamp = System.currentTimeMillis(); |
| 263 | + GenericRow genericRow = new GenericRow(1); |
| 264 | + genericRow.setField( |
| 265 | + 0, |
| 266 | + new GenericArray( |
| 267 | + new Object[] { |
| 268 | + new GenericArray(new int[] {1, 2}), |
| 269 | + new GenericArray(new int[] {3, 4, 5}) |
| 270 | + })); |
| 271 | + |
| 272 | + LogRecord logRecord = new GenericRecord(logOffset, timeStamp, APPEND_ONLY, genericRow); |
| 273 | + FlussRowAsPaimonRow flussRowAsPaimonRow = |
| 274 | + new FlussRowAsPaimonRow(logRecord.getRow(), tableRowType); |
| 275 | + |
| 276 | + InternalArray outerArray = flussRowAsPaimonRow.getArray(0); |
| 277 | + assertThat(outerArray).isNotNull(); |
| 278 | + assertThat(outerArray.size()).isEqualTo(2); |
| 279 | + |
| 280 | + InternalArray innerArray1 = outerArray.getArray(0); |
| 281 | + assertThat(innerArray1.size()).isEqualTo(2); |
| 282 | + assertThat(innerArray1.getInt(0)).isEqualTo(1); |
| 283 | + assertThat(innerArray1.getInt(1)).isEqualTo(2); |
| 284 | + |
| 285 | + InternalArray innerArray2 = outerArray.getArray(1); |
| 286 | + assertThat(innerArray2.size()).isEqualTo(3); |
| 287 | + assertThat(innerArray2.getInt(0)).isEqualTo(3); |
| 288 | + assertThat(innerArray2.getInt(1)).isEqualTo(4); |
| 289 | + assertThat(innerArray2.getInt(2)).isEqualTo(5); |
| 290 | + } |
| 291 | + |
| 292 | + @Test |
| 293 | + void testEmptyArray() { |
| 294 | + RowType tableRowType = |
| 295 | + RowType.of( |
| 296 | + new org.apache.paimon.types.ArrayType( |
| 297 | + new org.apache.paimon.types.IntType())); |
| 298 | + |
| 299 | + long logOffset = 0; |
| 300 | + long timeStamp = System.currentTimeMillis(); |
| 301 | + GenericRow genericRow = new GenericRow(1); |
| 302 | + genericRow.setField(0, new GenericArray(new int[] {})); |
| 303 | + |
| 304 | + LogRecord logRecord = new GenericRecord(logOffset, timeStamp, APPEND_ONLY, genericRow); |
| 305 | + FlussRowAsPaimonRow flussRowAsPaimonRow = |
| 306 | + new FlussRowAsPaimonRow(logRecord.getRow(), tableRowType); |
| 307 | + |
| 308 | + InternalArray array = flussRowAsPaimonRow.getArray(0); |
| 309 | + assertThat(array).isNotNull(); |
| 310 | + assertThat(array.size()).isEqualTo(0); |
| 311 | + } |
143 | 312 | } |
0 commit comments