|
32 | 32 |
|
33 | 33 | import java.util.ArrayList; |
34 | 34 | import java.util.Arrays; |
| 35 | +import java.util.Collections; |
35 | 36 | import java.util.HashMap; |
36 | 37 | import java.util.List; |
37 | 38 | import java.util.Map; |
| 39 | +import java.util.Optional; |
38 | 40 |
|
39 | 41 | /** A test for the {@link org.apache.flink.cdc.common.utils.SchemaUtils}. */ |
40 | 42 | class SchemaUtilsTest { |
@@ -484,4 +486,231 @@ void testInferWiderSchema() { |
484 | 486 | .build())) |
485 | 487 | .isExactlyInstanceOf(IllegalStateException.class); |
486 | 488 | } |
| 489 | + |
| 490 | + // ========================== Tests for duplicate AddColumnEvent handling |
| 491 | + // ========================== |
| 492 | + |
| 493 | + @Test |
| 494 | + void testFilterRedundantAddColumns_allDuplicates() { |
| 495 | + TableId tableId = TableId.parse("default.default.table1"); |
| 496 | + Schema schema = |
| 497 | + Schema.newBuilder() |
| 498 | + .physicalColumn("id", DataTypes.INT()) |
| 499 | + .physicalColumn("name", DataTypes.STRING()) |
| 500 | + .physicalColumn("age", DataTypes.INT()) |
| 501 | + .build(); |
| 502 | + |
| 503 | + List<AddColumnEvent.ColumnWithPosition> addedColumns = new ArrayList<>(); |
| 504 | + addedColumns.add( |
| 505 | + new AddColumnEvent.ColumnWithPosition( |
| 506 | + Column.physicalColumn("name", DataTypes.STRING()))); |
| 507 | + addedColumns.add( |
| 508 | + new AddColumnEvent.ColumnWithPosition( |
| 509 | + Column.physicalColumn("age", DataTypes.INT()))); |
| 510 | + AddColumnEvent addColumnEvent = new AddColumnEvent(tableId, addedColumns); |
| 511 | + |
| 512 | + Optional<AddColumnEvent> result = |
| 513 | + SchemaUtils.filterRedundantAddColumns(schema, addColumnEvent); |
| 514 | + Assertions.assertThat(result).isEmpty(); |
| 515 | + } |
| 516 | + |
| 517 | + @Test |
| 518 | + void testFilterRedundantAddColumns_noDuplicates() { |
| 519 | + TableId tableId = TableId.parse("default.default.table1"); |
| 520 | + Schema schema = |
| 521 | + Schema.newBuilder() |
| 522 | + .physicalColumn("id", DataTypes.INT()) |
| 523 | + .physicalColumn("name", DataTypes.STRING()) |
| 524 | + .build(); |
| 525 | + |
| 526 | + List<AddColumnEvent.ColumnWithPosition> addedColumns = new ArrayList<>(); |
| 527 | + addedColumns.add( |
| 528 | + new AddColumnEvent.ColumnWithPosition( |
| 529 | + Column.physicalColumn("age", DataTypes.INT()))); |
| 530 | + addedColumns.add( |
| 531 | + new AddColumnEvent.ColumnWithPosition( |
| 532 | + Column.physicalColumn("email", DataTypes.STRING()))); |
| 533 | + AddColumnEvent addColumnEvent = new AddColumnEvent(tableId, addedColumns); |
| 534 | + |
| 535 | + Optional<AddColumnEvent> result = |
| 536 | + SchemaUtils.filterRedundantAddColumns(schema, addColumnEvent); |
| 537 | + Assertions.assertThat(result).isPresent(); |
| 538 | + Assertions.assertThat(result.get().getAddedColumns()).hasSize(2); |
| 539 | + Assertions.assertThat(result.get()).isEqualTo(addColumnEvent); |
| 540 | + } |
| 541 | + |
| 542 | + @Test |
| 543 | + void testFilterRedundantAddColumns_partialDuplicates() { |
| 544 | + TableId tableId = TableId.parse("default.default.table1"); |
| 545 | + Schema schema = |
| 546 | + Schema.newBuilder() |
| 547 | + .physicalColumn("id", DataTypes.INT()) |
| 548 | + .physicalColumn("name", DataTypes.STRING()) |
| 549 | + .build(); |
| 550 | + |
| 551 | + List<AddColumnEvent.ColumnWithPosition> addedColumns = new ArrayList<>(); |
| 552 | + addedColumns.add( |
| 553 | + new AddColumnEvent.ColumnWithPosition( |
| 554 | + Column.physicalColumn("name", DataTypes.STRING()))); |
| 555 | + addedColumns.add( |
| 556 | + new AddColumnEvent.ColumnWithPosition( |
| 557 | + Column.physicalColumn("age", DataTypes.INT()))); |
| 558 | + AddColumnEvent addColumnEvent = new AddColumnEvent(tableId, addedColumns); |
| 559 | + |
| 560 | + Optional<AddColumnEvent> result = |
| 561 | + SchemaUtils.filterRedundantAddColumns(schema, addColumnEvent); |
| 562 | + Assertions.assertThat(result).isPresent(); |
| 563 | + Assertions.assertThat(result.get().getAddedColumns()).hasSize(1); |
| 564 | + Assertions.assertThat(result.get().getAddedColumns().get(0).getAddColumn().getName()) |
| 565 | + .isEqualTo("age"); |
| 566 | + } |
| 567 | + |
| 568 | + @Test |
| 569 | + void testFilterRedundantAddColumns_emptyAddedColumns() { |
| 570 | + TableId tableId = TableId.parse("default.default.table1"); |
| 571 | + Schema schema = |
| 572 | + Schema.newBuilder() |
| 573 | + .physicalColumn("id", DataTypes.INT()) |
| 574 | + .physicalColumn("name", DataTypes.STRING()) |
| 575 | + .build(); |
| 576 | + |
| 577 | + AddColumnEvent addColumnEvent = new AddColumnEvent(tableId, Collections.emptyList()); |
| 578 | + |
| 579 | + Optional<AddColumnEvent> result = |
| 580 | + SchemaUtils.filterRedundantAddColumns(schema, addColumnEvent); |
| 581 | + Assertions.assertThat(result).isEmpty(); |
| 582 | + } |
| 583 | + |
| 584 | + @Test |
| 585 | + void testApplyAddColumnEvent_idempotent() { |
| 586 | + TableId tableId = TableId.parse("default.default.table1"); |
| 587 | + Schema schema = |
| 588 | + Schema.newBuilder() |
| 589 | + .physicalColumn("id", DataTypes.INT()) |
| 590 | + .physicalColumn("name", DataTypes.STRING()) |
| 591 | + .build(); |
| 592 | + |
| 593 | + List<AddColumnEvent.ColumnWithPosition> addedColumns = new ArrayList<>(); |
| 594 | + addedColumns.add( |
| 595 | + new AddColumnEvent.ColumnWithPosition( |
| 596 | + Column.physicalColumn("name", DataTypes.STRING()))); |
| 597 | + addedColumns.add( |
| 598 | + new AddColumnEvent.ColumnWithPosition( |
| 599 | + Column.physicalColumn("age", DataTypes.INT()))); |
| 600 | + AddColumnEvent addColumnEvent = new AddColumnEvent(tableId, addedColumns); |
| 601 | + |
| 602 | + Schema result = SchemaUtils.applySchemaChangeEvent(schema, addColumnEvent); |
| 603 | + Assertions.assertThat(result) |
| 604 | + .isEqualTo( |
| 605 | + Schema.newBuilder() |
| 606 | + .physicalColumn("id", DataTypes.INT()) |
| 607 | + .physicalColumn("name", DataTypes.STRING()) |
| 608 | + .physicalColumn("age", DataTypes.INT()) |
| 609 | + .build()); |
| 610 | + } |
| 611 | + |
| 612 | + @Test |
| 613 | + void testApplyAddColumnEvent_allDuplicates() { |
| 614 | + TableId tableId = TableId.parse("default.default.table1"); |
| 615 | + Schema schema = |
| 616 | + Schema.newBuilder() |
| 617 | + .physicalColumn("id", DataTypes.INT()) |
| 618 | + .physicalColumn("name", DataTypes.STRING()) |
| 619 | + .physicalColumn("age", DataTypes.INT()) |
| 620 | + .build(); |
| 621 | + |
| 622 | + List<AddColumnEvent.ColumnWithPosition> addedColumns = new ArrayList<>(); |
| 623 | + addedColumns.add( |
| 624 | + new AddColumnEvent.ColumnWithPosition( |
| 625 | + Column.physicalColumn("name", DataTypes.STRING()))); |
| 626 | + addedColumns.add( |
| 627 | + new AddColumnEvent.ColumnWithPosition( |
| 628 | + Column.physicalColumn("age", DataTypes.INT()))); |
| 629 | + AddColumnEvent addColumnEvent = new AddColumnEvent(tableId, addedColumns); |
| 630 | + |
| 631 | + Schema result = SchemaUtils.applySchemaChangeEvent(schema, addColumnEvent); |
| 632 | + Assertions.assertThat(result) |
| 633 | + .isEqualTo( |
| 634 | + Schema.newBuilder() |
| 635 | + .physicalColumn("id", DataTypes.INT()) |
| 636 | + .physicalColumn("name", DataTypes.STRING()) |
| 637 | + .physicalColumn("age", DataTypes.INT()) |
| 638 | + .build()); |
| 639 | + } |
| 640 | + |
| 641 | + @Test |
| 642 | + void testFilterRedundantAddColumns_withPositions() { |
| 643 | + TableId tableId = TableId.parse("default.default.table1"); |
| 644 | + Schema schema = |
| 645 | + Schema.newBuilder() |
| 646 | + .physicalColumn("id", DataTypes.INT()) |
| 647 | + .physicalColumn("name", DataTypes.STRING()) |
| 648 | + .build(); |
| 649 | + |
| 650 | + List<AddColumnEvent.ColumnWithPosition> addedColumns = new ArrayList<>(); |
| 651 | + addedColumns.add( |
| 652 | + new AddColumnEvent.ColumnWithPosition( |
| 653 | + Column.physicalColumn("name", DataTypes.STRING()), |
| 654 | + AddColumnEvent.ColumnPosition.AFTER, |
| 655 | + "id")); |
| 656 | + addedColumns.add( |
| 657 | + new AddColumnEvent.ColumnWithPosition( |
| 658 | + Column.physicalColumn("age", DataTypes.INT()), |
| 659 | + AddColumnEvent.ColumnPosition.AFTER, |
| 660 | + "name")); |
| 661 | + AddColumnEvent addColumnEvent = new AddColumnEvent(tableId, addedColumns); |
| 662 | + |
| 663 | + Optional<AddColumnEvent> result = |
| 664 | + SchemaUtils.filterRedundantAddColumns(schema, addColumnEvent); |
| 665 | + Assertions.assertThat(result).isPresent(); |
| 666 | + Assertions.assertThat(result.get().getAddedColumns()).hasSize(1); |
| 667 | + |
| 668 | + AddColumnEvent.ColumnWithPosition remaining = result.get().getAddedColumns().get(0); |
| 669 | + Assertions.assertThat(remaining.getAddColumn().getName()).isEqualTo("age"); |
| 670 | + Assertions.assertThat(remaining.getPosition()) |
| 671 | + .isEqualTo(AddColumnEvent.ColumnPosition.AFTER); |
| 672 | + Assertions.assertThat(remaining.getExistedColumnName()).isEqualTo("name"); |
| 673 | + } |
| 674 | + |
| 675 | + @Test |
| 676 | + void testFilterRedundantAddColumns_intraEventDuplicates() { |
| 677 | + Schema schema = Schema.newBuilder().physicalColumn("id", DataTypes.INT()).build(); |
| 678 | + TableId tableId = TableId.tableId("default", "schema", "table"); |
| 679 | + AddColumnEvent event = |
| 680 | + new AddColumnEvent( |
| 681 | + tableId, |
| 682 | + Arrays.asList( |
| 683 | + new AddColumnEvent.ColumnWithPosition( |
| 684 | + Column.physicalColumn("name", DataTypes.STRING())), |
| 685 | + new AddColumnEvent.ColumnWithPosition( |
| 686 | + Column.physicalColumn("name", DataTypes.STRING())))); |
| 687 | + Optional<AddColumnEvent> result = SchemaUtils.filterRedundantAddColumns(schema, event); |
| 688 | + Assertions.assertThat(result).isPresent(); |
| 689 | + Assertions.assertThat(result.get().getAddedColumns()).hasSize(1); |
| 690 | + Assertions.assertThat(result.get().getAddedColumns().get(0).getAddColumn().getName()) |
| 691 | + .isEqualTo("name"); |
| 692 | + } |
| 693 | + |
| 694 | + @Test |
| 695 | + void testApplyAddColumnEvent_duplicateWithinSameEvent() { |
| 696 | + TableId tableId = TableId.parse("default.default.table1"); |
| 697 | + Schema schema = Schema.newBuilder().physicalColumn("id", DataTypes.INT()).build(); |
| 698 | + |
| 699 | + List<AddColumnEvent.ColumnWithPosition> addedColumns = new ArrayList<>(); |
| 700 | + addedColumns.add( |
| 701 | + new AddColumnEvent.ColumnWithPosition( |
| 702 | + Column.physicalColumn("name", DataTypes.STRING()))); |
| 703 | + addedColumns.add( |
| 704 | + new AddColumnEvent.ColumnWithPosition( |
| 705 | + Column.physicalColumn("name", DataTypes.STRING()))); |
| 706 | + AddColumnEvent addColumnEvent = new AddColumnEvent(tableId, addedColumns); |
| 707 | + |
| 708 | + Schema result = SchemaUtils.applySchemaChangeEvent(schema, addColumnEvent); |
| 709 | + Assertions.assertThat(result) |
| 710 | + .isEqualTo( |
| 711 | + Schema.newBuilder() |
| 712 | + .physicalColumn("id", DataTypes.INT()) |
| 713 | + .physicalColumn("name", DataTypes.STRING()) |
| 714 | + .build()); |
| 715 | + } |
487 | 716 | } |
0 commit comments