From ae3c0ab6a77f82f965df883a1e0e88bbf2d25bfc Mon Sep 17 00:00:00 2001 From: Christian Tzolov Date: Fri, 5 Aug 2022 16:05:28 +0200 Subject: [PATCH] fix regressions. Add DB2 Itests Signed-off-by: Christian Tzolov --- .../source/cdc-debezium-source/README.adoc | 2 +- .../cdc/CdcFlatteningIntegrationTest.java | 4 +- .../app/source/cdc/CdcMySqlTestSupport.java | 2 +- .../CdcSourceDatabasesIntegrationTest.java | 52 ++++++++++++++++++- .../cdc/CdcBootStarterIntegrationTest.java | 2 +- functions/common/cdc-debezium-common/pom.xml | 15 +++++- .../fn/common/cdc/CdcCommonConfiguration.java | 6 +-- .../fn/common/cdc/CdcCommonProperties.java | 2 +- .../supplier/cdc-debezium-supplier/pom.xml | 19 ++++++- 9 files changed, 91 insertions(+), 13 deletions(-) diff --git a/applications/source/cdc-debezium-source/README.adoc b/applications/source/cdc-debezium-source/README.adoc index d635293a..2020b9de 100644 --- a/applications/source/cdc-debezium-source/README.adoc +++ b/applications/source/cdc-debezium-source/README.adoc @@ -23,7 +23,7 @@ Properties grouped by prefix: === cdc $$config$$:: $$Spring pass-trough wrapper for debezium configuration properties. All properties with a 'cdc.config.' prefix are native Debezium properties. The prefix is removed, converting them into Debezium io.debezium.config.Configuration.$$ *($$Map$$, default: `$$$$`)* -$$connector$$:: $$Shortcut for the cdc.config.connector.class property. Either of those can be used as long as they do not contradict with each other.$$ *($$ConnectorType$$, default: `$$$$`, possible values: `mysql`,`postgres`,`mongodb`,`oracle`,`sqlserver`)* +$$connector$$:: $$Shortcut for the cdc.config.connector.class property. Either of those can be used as long as they do not contradict with each other.$$ *($$ConnectorType$$, default: `$$$$`, possible values: `mysql`,`postgres`,`mongodb`,`oracle`,`db2`,`sqlserver`)* $$name$$:: $$Unique name for this sourceConnector instance.$$ *($$String$$, default: `$$$$`)* $$schema$$:: $$Include the schema's as part of the outbound message.$$ *($$Boolean$$, default: `$$false$$`)* diff --git a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteningIntegrationTest.java b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteningIntegrationTest.java index 2e60cd75..2cd3d2aa 100644 --- a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteningIntegrationTest.java +++ b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcFlatteningIntegrationTest.java @@ -212,7 +212,7 @@ public class CdcFlatteningIntegrationTest extends CdcMySqlTestSupport { assertThat(messages.get(1).getHeaders().get("cdc_topic")).isEqualTo("my-app-connector.inventory.customers"); JsonAssert.assertJsonEquals("{\"id\":" + newRecordId + "}", toString(messages.get(1).getHeaders().get("cdc_key"))); - if (!StringUtils.hasText(flatteningProps.getAddHeaders()) && flatteningProps.getAddHeaders().contains("op")) { + if (StringUtils.hasText(flatteningProps.getAddHeaders()) && flatteningProps.getAddHeaders().contains("op")) { assertThat(messages.get(1).getHeaders().get("__op")).isEqualTo("u"); } @@ -220,7 +220,7 @@ public class CdcFlatteningIntegrationTest extends CdcMySqlTestSupport { assertThat(toString(messages.get(2).getPayload())).isEqualTo("null"); assertThat(messages.get(1).getHeaders().get("cdc_topic")).isEqualTo("my-app-connector.inventory.customers"); JsonAssert.assertJsonEquals("{\"id\":" + newRecordId + "}", toString(messages.get(1).getHeaders().get("cdc_key"))); - if (!StringUtils.hasText(flatteningProps.getAddHeaders()) + if (StringUtils.hasText(flatteningProps.getAddHeaders()) && flatteningProps.getAddHeaders().contains("op")) { assertThat(messages.get(2).getHeaders().get("__op")).isEqualTo("d"); } diff --git a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcMySqlTestSupport.java b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcMySqlTestSupport.java index 1cfdc0a4..3665b31e 100644 --- a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcMySqlTestSupport.java +++ b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcMySqlTestSupport.java @@ -1,5 +1,5 @@ /* - * Copyright 2020-2021 the original author or authors. + * Copyright 2020-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. diff --git a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcSourceDatabasesIntegrationTest.java b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcSourceDatabasesIntegrationTest.java index dc8e7754..de39eb15 100644 --- a/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcSourceDatabasesIntegrationTest.java +++ b/applications/source/cdc-debezium-source/src/test/java/org/springframework/cloud/stream/app/source/cdc/CdcSourceDatabasesIntegrationTest.java @@ -1,5 +1,5 @@ /* - * Copyright 2020-2021 the original author or authors. + * Copyright 2020-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -154,6 +154,56 @@ public class CdcSourceDatabasesIntegrationTest { postgres.stop(); } + // From within the src/test/docker/db2 folder run: + // docker build -t db2-cdc2 . + // + // docker run -itd --name mydb2 --privileged=true -p 50000:50000 -e LICENSE=accept -e DB2INST1_PASSWORD=password -e + // DBNAME=testdb db2-cdc2 + // + // docker logs -f mydb2 + // docker stop mydb2 + // docker exec -it mydb2 /bin/bash + @Test + @Disabled + public void db2() { + + // TODO + // GenericContainer db2Container = new GenericContainer( + // new ImageFromDockerfile() + // .withFileFromPath(".", file.toPath())) + // .withPrivilegedMode(true) + // .withExposedPorts(50000) + // .withEnv("LICENSE", "accept") + // .withEnv("DBNAME", "testdb") + // .withEnv("DB2INST1_PASSWORD", "password"); + // db2Container.start(); + + try (ConfigurableApplicationContext context = applicationBuilder + .run("--cdc.connector=db2", + "--cdc.config.database.user=db2inst1", + "--cdc.config.database.password=password", + "--cdc.config.slot.name=debezium", + "--cdc.config.database.dbname=TESTDB", + "--cdc.config.database.hostname=localhost", + // "--cdc.config.table.include.list=inventory.*", + "--cdc.config.database.port=" + "50000")) { + OutputDestination outputDestination = context.getBean(OutputDestination.class); + // Using local region here + + List> allMessages = new ArrayList<>(); + Awaitility.await().atMost(Duration.ofMinutes(5)).until(() -> { + List> messageChunk = CdcTestUtils.receiveAll(outputDestination); + if (!CollectionUtils.isEmpty(messageChunk)) { + System.out.println("Chunk size: " + messageChunk.size()); + allMessages.addAll(messageChunk); + } + return allMessages.size() == 30; // Inventory DB entries + }); + } + + // db2Container.stop(); + } + @Test @Disabled public void mongodb() { diff --git a/functions/common/cdc-debezium-boot-starter/src/test/java/org/springframework/cloud/fn/common/cdc/CdcBootStarterIntegrationTest.java b/functions/common/cdc-debezium-boot-starter/src/test/java/org/springframework/cloud/fn/common/cdc/CdcBootStarterIntegrationTest.java index e7019034..351100f1 100644 --- a/functions/common/cdc-debezium-boot-starter/src/test/java/org/springframework/cloud/fn/common/cdc/CdcBootStarterIntegrationTest.java +++ b/functions/common/cdc-debezium-boot-starter/src/test/java/org/springframework/cloud/fn/common/cdc/CdcBootStarterIntegrationTest.java @@ -1,5 +1,5 @@ /* - * Copyright 2020-2021 the original author or authors. + * Copyright 2020-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. diff --git a/functions/common/cdc-debezium-common/pom.xml b/functions/common/cdc-debezium-common/pom.xml index 9b6ea394..bdca7785 100644 --- a/functions/common/cdc-debezium-common/pom.xml +++ b/functions/common/cdc-debezium-common/pom.xml @@ -1,6 +1,6 @@ - + 4.0.0 @@ -19,6 +19,11 @@ + + io.debezium + debezium-api + ${version.debezium} + io.debezium debezium-embedded @@ -60,6 +65,12 @@ true ${version.debezium} + + io.debezium + debezium-connector-db2 + true + ${version.debezium} + org.springframework.cloud.fn diff --git a/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonConfiguration.java b/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonConfiguration.java index e9837cbe..98f08465 100644 --- a/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonConfiguration.java +++ b/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonConfiguration.java @@ -1,5 +1,5 @@ /* - * Copyright 2020-2020 the original author or authors. + * Copyright 2020-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -64,10 +64,10 @@ public class CdcCommonConfiguration { Map config = extractNewRecordState.config().defaultValues(); config.put("drop.tombstones", properties.getFlattening().isDropTombstones()); config.put("delete.handling.mode", properties.getFlattening().getDeleteHandlingMode().name()); - if (!StringUtils.hasText(properties.getFlattening().getAddHeaders())) { + if (StringUtils.hasText(properties.getFlattening().getAddHeaders())) { config.put("add.headers", properties.getFlattening().getAddHeaders()); } - if (!StringUtils.hasText(properties.getFlattening().getAddFields())) { + if (StringUtils.hasText(properties.getFlattening().getAddFields())) { config.put("add.fields", properties.getFlattening().getAddFields()); } diff --git a/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonProperties.java b/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonProperties.java index d20b187f..2369937b 100644 --- a/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonProperties.java +++ b/functions/common/cdc-debezium-common/src/main/java/org/springframework/cloud/fn/common/cdc/CdcCommonProperties.java @@ -1,5 +1,5 @@ /* - * Copyright 2020-2020 the original author or authors. + * Copyright 2020-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. diff --git a/functions/supplier/cdc-debezium-supplier/pom.xml b/functions/supplier/cdc-debezium-supplier/pom.xml index 35b2e656..bc4b921c 100644 --- a/functions/supplier/cdc-debezium-supplier/pom.xml +++ b/functions/supplier/cdc-debezium-supplier/pom.xml @@ -1,5 +1,6 @@ - + 4.0.0 @@ -78,5 +79,21 @@ ${version.debezium} + + io.debezium + debezium-connector-db2 + + + slf4j-log4j12 + org.slf4j + + + ${version.debezium} + + + com.ibm.db2.jcc + db2jcc + db2jcc4 +