This is an automated email from the ASF dual-hosted git repository. fpaul pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/flink-connector-kafka.git
commit e75c00e67bf80bd3dd4ec84a17b8cc935db2074b Author: Aleksandr Savonin <[email protected]> AuthorDate: Sun Jan 18 15:37:40 2026 +0100 [FLINK-38937] Rename KAFKA constant to CP_KAFKA for clarity --- .../java/org/apache/flink/tests/util/kafka/KafkaSinkE2ECase.java | 2 +- .../java/org/apache/flink/tests/util/kafka/KafkaSourceE2ECase.java | 4 ++-- .../java/org/apache/flink/connector/kafka/sink/KafkaSinkITCase.java | 2 +- .../org/apache/flink/connector/kafka/source/KafkaSourceITCase.java | 2 +- .../apache/flink/connector/kafka/testutils/DockerImageVersions.java | 2 +- .../kafka/testutils/FlinkKafkaIntegrationCompatibilityTest.java | 6 +++--- .../java/org/apache/flink/connector/kafka/testutils/KafkaUtil.java | 2 +- .../connector/kafka/testutils/TestKafkaContainerValidationTest.java | 2 +- .../apache/flink/connector/kafka/testutils/TwoKafkaContainers.java | 2 +- .../flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java | 2 +- 10 files changed, 13 insertions(+), 13 deletions(-) diff --git a/flink-connector-kafka-e2e-tests/flink-end-to-end-tests-common-kafka/src/test/java/org/apache/flink/tests/util/kafka/KafkaSinkE2ECase.java b/flink-connector-kafka-e2e-tests/flink-end-to-end-tests-common-kafka/src/test/java/org/apache/flink/tests/util/kafka/KafkaSinkE2ECase.java index 479e3b3b..c34b179c 100644 --- a/flink-connector-kafka-e2e-tests/flink-end-to-end-tests-common-kafka/src/test/java/org/apache/flink/tests/util/kafka/KafkaSinkE2ECase.java +++ b/flink-connector-kafka-e2e-tests/flink-end-to-end-tests-common-kafka/src/test/java/org/apache/flink/tests/util/kafka/KafkaSinkE2ECase.java @@ -51,7 +51,7 @@ public class KafkaSinkE2ECase extends SinkTestSuiteBase<String> { @TestEnv FlinkContainerTestEnvironment flink = new FlinkContainerTestEnvironment(1, 6); private final TestKafkaContainer kafkaContainer = - new TestKafkaContainer(DockerImageVersions.KAFKA).withNetworkAliases(KAFKA_HOSTNAME); + new TestKafkaContainer(DockerImageVersions.CP_KAFKA).withNetworkAliases(KAFKA_HOSTNAME); // Defines ConnectorExternalSystem @SuppressWarnings({"rawtypes", "unchecked"}) diff --git a/flink-connector-kafka-e2e-tests/flink-end-to-end-tests-common-kafka/src/test/java/org/apache/flink/tests/util/kafka/KafkaSourceE2ECase.java b/flink-connector-kafka-e2e-tests/flink-end-to-end-tests-common-kafka/src/test/java/org/apache/flink/tests/util/kafka/KafkaSourceE2ECase.java index 67f7dfa3..37ea90dc 100644 --- a/flink-connector-kafka-e2e-tests/flink-end-to-end-tests-common-kafka/src/test/java/org/apache/flink/tests/util/kafka/KafkaSourceE2ECase.java +++ b/flink-connector-kafka-e2e-tests/flink-end-to-end-tests-common-kafka/src/test/java/org/apache/flink/tests/util/kafka/KafkaSourceE2ECase.java @@ -34,7 +34,7 @@ import org.testcontainers.containers.GenericContainer; import java.util.Arrays; -import static org.apache.flink.connector.kafka.testutils.DockerImageVersions.KAFKA; +import static org.apache.flink.connector.kafka.testutils.DockerImageVersions.CP_KAFKA; import static org.apache.flink.connector.kafka.testutils.KafkaSourceExternalContext.SplitMappingMode.PARTITION; import static org.apache.flink.connector.kafka.testutils.KafkaSourceExternalContext.SplitMappingMode.TOPIC; @@ -49,7 +49,7 @@ public class KafkaSourceE2ECase extends SourceTestSuiteBase<String> { @TestEnv FlinkContainerTestEnvironment flink = new FlinkContainerTestEnvironment(1, 6); TestKafkaContainer kafkaContainer = - new TestKafkaContainer(KAFKA).withNetworkAliases(KAFKA_HOSTNAME); + new TestKafkaContainer(CP_KAFKA).withNetworkAliases(KAFKA_HOSTNAME); // Defines ConnectorExternalSystem @SuppressWarnings({"rawtypes", "unchecked"}) diff --git a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaSinkITCase.java b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaSinkITCase.java index 773b13b8..2504a4cd 100644 --- a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaSinkITCase.java +++ b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/sink/KafkaSinkITCase.java @@ -194,7 +194,7 @@ public class KafkaSinkITCase extends TestLogger { @TestEnv MiniClusterTestEnvironment flink = new MiniClusterTestEnvironment(); private final TestKafkaContainer kafkaContainer = - new TestKafkaContainer(DockerImageName.parse(DockerImageVersions.KAFKA)); + new TestKafkaContainer(DockerImageName.parse(DockerImageVersions.CP_KAFKA)); // Defines external system @SuppressWarnings({"rawtypes", "unchecked"}) diff --git a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceITCase.java b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceITCase.java index c5fe60ee..73028290 100644 --- a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceITCase.java +++ b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/source/KafkaSourceITCase.java @@ -421,7 +421,7 @@ public class KafkaSourceITCase { MiniClusterTestEnvironment flink = new MiniClusterTestEnvironment(); private final TestKafkaContainer kafkaContainer = - new TestKafkaContainer(DockerImageName.parse(DockerImageVersions.KAFKA)); + new TestKafkaContainer(DockerImageName.parse(DockerImageVersions.CP_KAFKA)); // Defines external system @SuppressWarnings({"rawtypes", "unchecked"}) diff --git a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/DockerImageVersions.java b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/DockerImageVersions.java index ec1e1535..f33e9f8a 100644 --- a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/DockerImageVersions.java +++ b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/DockerImageVersions.java @@ -24,7 +24,7 @@ package org.apache.flink.connector.kafka.testutils; */ public class DockerImageVersions { - public static final String KAFKA = "confluentinc/cp-kafka:7.9.2"; + public static final String CP_KAFKA = "confluentinc/cp-kafka:7.9.2"; public static final String APACHE_KAFKA = "apache/kafka:4.1.1"; public static final String SCHEMA_REGISTRY = "confluentinc/cp-schema-registry:7.9.2"; diff --git a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/FlinkKafkaIntegrationCompatibilityTest.java b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/FlinkKafkaIntegrationCompatibilityTest.java index 925dd6b0..938fb7e7 100644 --- a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/FlinkKafkaIntegrationCompatibilityTest.java +++ b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/FlinkKafkaIntegrationCompatibilityTest.java @@ -54,7 +54,7 @@ import java.util.Map; import java.util.UUID; import static org.apache.flink.connector.kafka.testutils.DockerImageVersions.APACHE_KAFKA; -import static org.apache.flink.connector.kafka.testutils.DockerImageVersions.KAFKA; +import static org.apache.flink.connector.kafka.testutils.DockerImageVersions.CP_KAFKA; import static org.assertj.core.api.Assertions.assertThat; /** @@ -86,7 +86,7 @@ class FlinkKafkaIntegrationCompatibilityTest { * <p>This is adapted from {@code KafkaSourceITCase.testValueOnlyDeserializer}. */ @ParameterizedTest - @ValueSource(strings = {KAFKA, APACHE_KAFKA}) + @ValueSource(strings = {CP_KAFKA, APACHE_KAFKA}) void testFlinkKafkaSourceIntegration(String dockerImage) throws Exception { // Start Kafka container kafkaContainer = new TestKafkaContainer(dockerImage); @@ -181,7 +181,7 @@ class FlinkKafkaIntegrationCompatibilityTest { * <p>This is adapted from {@code KafkaSinkITCase.testWriteRecordsToKafkaWithNoneGuarantee}. */ @ParameterizedTest - @ValueSource(strings = {KAFKA, APACHE_KAFKA}) + @ValueSource(strings = {CP_KAFKA, APACHE_KAFKA}) void testFlinkKafkaSinkIntegration(String dockerImage) throws Exception { // Start Kafka container kafkaContainer = new TestKafkaContainer(dockerImage); diff --git a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/KafkaUtil.java b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/KafkaUtil.java index 934c1c92..f0e747e1 100644 --- a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/KafkaUtil.java +++ b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/KafkaUtil.java @@ -62,7 +62,7 @@ public class KafkaUtil { String logLevel = inferLogLevel(logger); Slf4jLogConsumer logConsumer = new Slf4jLogConsumer(logger, true); - return new TestKafkaContainer(DockerImageName.parse(DockerImageVersions.KAFKA)) + return new TestKafkaContainer(DockerImageName.parse(DockerImageVersions.CP_KAFKA)) .withEnv("KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR", "1") .withEnv("KAFKA_TRANSACTION_STATE_LOG_MIN_ISR", "1") .withEnv("KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR", "1") diff --git a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/TestKafkaContainerValidationTest.java b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/TestKafkaContainerValidationTest.java index 38d8a81a..8ddfbbe6 100644 --- a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/TestKafkaContainerValidationTest.java +++ b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/TestKafkaContainerValidationTest.java @@ -27,7 +27,7 @@ class TestKafkaContainerValidationTest { @Test void testConfluentImageIsAccepted() { - assertThatCode(() -> new TestKafkaContainer(DockerImageVersions.KAFKA)) + assertThatCode(() -> new TestKafkaContainer(DockerImageVersions.CP_KAFKA)) .doesNotThrowAnyException(); } diff --git a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/TwoKafkaContainers.java b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/TwoKafkaContainers.java index 247357e4..b15a95c3 100644 --- a/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/TwoKafkaContainers.java +++ b/flink-connector-kafka/src/test/java/org/apache/flink/connector/kafka/testutils/TwoKafkaContainers.java @@ -27,7 +27,7 @@ public class TwoKafkaContainers extends GenericContainer<TwoKafkaContainers> { private final TestKafkaContainer kafka1; public TwoKafkaContainers() { - DockerImageName dockerImageName = DockerImageName.parse(DockerImageVersions.KAFKA); + DockerImageName dockerImageName = DockerImageName.parse(DockerImageVersions.CP_KAFKA); this.kafka0 = new TestKafkaContainer(dockerImageName); this.kafka1 = new TestKafkaContainer(dockerImageName); } diff --git a/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java b/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java index a4d1a8b8..ee55cb01 100644 --- a/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java +++ b/flink-connector-kafka/src/test/java/org/apache/flink/streaming/connectors/kafka/KafkaTestEnvironmentImpl.java @@ -198,7 +198,7 @@ public class KafkaTestEnvironmentImpl extends KafkaTestEnvironment { @Override public String getVersion() { - return DockerImageVersions.KAFKA; + return DockerImageVersions.CP_KAFKA; } @Override
