MartijnVisser commented on code in PR #26662: URL: https://github.com/apache/flink/pull/26662#discussion_r3960930917
########## flink-end-to-end-tests/flink-confluent-schema-registry/src/test/java/org/apache/flink/schema/registry/test/AvroConfluentITCase.java: ########## @@ -0,0 +1,444 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.schema.registry.test; + +import org.apache.flink.connector.testframe.container.FlinkContainers; +import org.apache.flink.connector.testframe.container.FlinkContainersSettings; +import org.apache.flink.connector.testframe.container.TestcontainersSettings; +import org.apache.flink.core.testutils.CommonTestUtils; +import org.apache.flink.test.resources.ResourceTestUtils; +import org.apache.flink.test.util.SQLJobSubmission; +import org.apache.flink.util.DockerImageVersions; +import org.apache.flink.util.jackson.JacksonMapperFactory; + +import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper; + +import org.apache.avro.generic.GenericRecord; +import org.apache.kafka.clients.admin.Admin; +import org.apache.kafka.clients.admin.NewTopic; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.testcontainers.containers.GenericContainer; +import org.testcontainers.containers.KafkaContainer; +import org.testcontainers.containers.Network; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.utility.DockerImageName; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.HttpRequest; +import java.net.http.HttpResponse; +import java.nio.file.Path; +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Properties; +import java.util.Set; +import java.util.concurrent.TimeoutException; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** E2E Test for Avro-Confluent integration. */ +@Testcontainers +public class AvroConfluentITCase { + + private static final Logger LOG = LoggerFactory.getLogger(AvroConfluentITCase.class); + + private static final String INTER_CONTAINER_KAFKA_ALIAS = "kafka"; + private static final String INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS = "schema-registry"; + + private static final String TOPIC = "test-avro-input"; + private static final String RESULT_TOPIC = "test-avro-output"; + private static final String MANUAL_TOPIC = "test-avro-input-manual"; + private static final String MANUAL_RESULT_TOPIC = "test-avro-output-manual"; + + private static final String INPUT_SCHEMA = SchemaLoader.loadSchema("avro/input-record.avsc"); + private static final String OUTPUT_SCHEMA = SchemaLoader.loadSchema("avro/output-record.avsc"); + + private static final Path sqlToolBoxJar = ResourceTestUtils.getResource(".*/SqlToolbox\\.jar"); + + private final Path sqlConnectorKafkaJar = ResourceTestUtils.getResource(".*kafka.*\\.jar"); + private final Path sqlConnectorUpsertTestJar = + ResourceTestUtils.getResource(".*flink-test-utils.*\\.jar"); + private final Path sqlAvroConfluentJar = + ResourceTestUtils.getResource(".*avro-confluent.*\\.jar"); + + private static final Network NETWORK = Network.newNetwork(); + + @Container + public static final KafkaContainer KAFKA = + new KafkaContainer(DockerImageName.parse(DockerImageVersions.KAFKA)) + .withNetwork(NETWORK) + .withNetworkAliases(INTER_CONTAINER_KAFKA_ALIAS); + + @Container + private static final GenericContainer<?> SCHEMA_REGISTRY = + new GenericContainer<>(DockerImageName.parse(DockerImageVersions.SCHEMA_REGISTRY)) + .withNetwork(NETWORK) + .withNetworkAliases(INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS) + .withExposedPorts(8081) + .withEnv("SCHEMA_REGISTRY_HOST_NAME", INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS) + .withEnv( + "SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", + INTER_CONTAINER_KAFKA_ALIAS + ":9092") + .dependsOn(KAFKA); + + private static final FlinkContainers FLINK = + FlinkContainers.builder() + .withFlinkContainersSettings( + FlinkContainersSettings.builder().numTaskManagers(1).build()) + .withTestcontainersSettings( + TestcontainersSettings.builder() + .network(NETWORK) + .logger(LOG) + .dependsOn(KAFKA) + .build()) + .build(); + + private final HttpClient client = HttpClient.newHttpClient(); + private final ObjectMapper objectMapper = JacksonMapperFactory.createObjectMapper(); + + @BeforeAll + public static void setup() throws Exception { + KAFKA.start(); + SCHEMA_REGISTRY.start(); + FLINK.start(); + } + + @AfterAll + public static void tearDown() { + FLINK.stop(); + SCHEMA_REGISTRY.stop(); + KAFKA.stop(); + } Review Comment: Removed, the container annotations and the Flink extension handle the lifecycle now. ########## flink-end-to-end-tests/flink-confluent-schema-registry/src/test/java/org/apache/flink/schema/registry/test/AvroConfluentITCase.java: ########## @@ -0,0 +1,444 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.schema.registry.test; + +import org.apache.flink.connector.testframe.container.FlinkContainers; +import org.apache.flink.connector.testframe.container.FlinkContainersSettings; +import org.apache.flink.connector.testframe.container.TestcontainersSettings; +import org.apache.flink.core.testutils.CommonTestUtils; +import org.apache.flink.test.resources.ResourceTestUtils; +import org.apache.flink.test.util.SQLJobSubmission; +import org.apache.flink.util.DockerImageVersions; +import org.apache.flink.util.jackson.JacksonMapperFactory; + +import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper; + +import org.apache.avro.generic.GenericRecord; +import org.apache.kafka.clients.admin.Admin; +import org.apache.kafka.clients.admin.NewTopic; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.testcontainers.containers.GenericContainer; +import org.testcontainers.containers.KafkaContainer; +import org.testcontainers.containers.Network; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.utility.DockerImageName; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.HttpRequest; +import java.net.http.HttpResponse; +import java.nio.file.Path; +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Properties; +import java.util.Set; +import java.util.concurrent.TimeoutException; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** E2E Test for Avro-Confluent integration. */ +@Testcontainers +public class AvroConfluentITCase { + + private static final Logger LOG = LoggerFactory.getLogger(AvroConfluentITCase.class); + + private static final String INTER_CONTAINER_KAFKA_ALIAS = "kafka"; + private static final String INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS = "schema-registry"; + + private static final String TOPIC = "test-avro-input"; + private static final String RESULT_TOPIC = "test-avro-output"; + private static final String MANUAL_TOPIC = "test-avro-input-manual"; + private static final String MANUAL_RESULT_TOPIC = "test-avro-output-manual"; + + private static final String INPUT_SCHEMA = SchemaLoader.loadSchema("avro/input-record.avsc"); + private static final String OUTPUT_SCHEMA = SchemaLoader.loadSchema("avro/output-record.avsc"); + + private static final Path sqlToolBoxJar = ResourceTestUtils.getResource(".*/SqlToolbox\\.jar"); + + private final Path sqlConnectorKafkaJar = ResourceTestUtils.getResource(".*kafka.*\\.jar"); + private final Path sqlConnectorUpsertTestJar = + ResourceTestUtils.getResource(".*flink-test-utils.*\\.jar"); + private final Path sqlAvroConfluentJar = + ResourceTestUtils.getResource(".*avro-confluent.*\\.jar"); + + private static final Network NETWORK = Network.newNetwork(); + + @Container + public static final KafkaContainer KAFKA = + new KafkaContainer(DockerImageName.parse(DockerImageVersions.KAFKA)) + .withNetwork(NETWORK) + .withNetworkAliases(INTER_CONTAINER_KAFKA_ALIAS); + + @Container + private static final GenericContainer<?> SCHEMA_REGISTRY = + new GenericContainer<>(DockerImageName.parse(DockerImageVersions.SCHEMA_REGISTRY)) + .withNetwork(NETWORK) + .withNetworkAliases(INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS) + .withExposedPorts(8081) + .withEnv("SCHEMA_REGISTRY_HOST_NAME", INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS) + .withEnv( + "SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", + INTER_CONTAINER_KAFKA_ALIAS + ":9092") + .dependsOn(KAFKA); + + private static final FlinkContainers FLINK = + FlinkContainers.builder() + .withFlinkContainersSettings( + FlinkContainersSettings.builder().numTaskManagers(1).build()) + .withTestcontainersSettings( + TestcontainersSettings.builder() + .network(NETWORK) + .logger(LOG) + .dependsOn(KAFKA) + .build()) + .build(); + + private final HttpClient client = HttpClient.newHttpClient(); + private final ObjectMapper objectMapper = JacksonMapperFactory.createObjectMapper(); + + @BeforeAll + public static void setup() throws Exception { + KAFKA.start(); + SCHEMA_REGISTRY.start(); + FLINK.start(); + } + + @AfterAll + public static void tearDown() { + FLINK.stop(); + SCHEMA_REGISTRY.stop(); + KAFKA.stop(); + } + + @BeforeEach + public void before() throws Exception { + // Create topics using external bootstrap servers since we're outside the Docker network + Properties props = new Properties(); + props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, KAFKA.getBootstrapServers()); + try (Admin admin = Admin.create(props)) { + admin.createTopics( + Arrays.asList( + new NewTopic(TOPIC, 1, (short) 1), + new NewTopic(RESULT_TOPIC, 1, (short) 1), + new NewTopic(MANUAL_TOPIC, 1, (short) 1), + new NewTopic(MANUAL_RESULT_TOPIC, 1, (short) 1))); + + // Poll for topic creation with timeout + try { + CommonTestUtils.waitUntilIgnoringExceptions( + () -> { + try { + Set<String> topics = admin.listTopics().names().get(); + return topics.contains(TOPIC) + && topics.contains(RESULT_TOPIC) + && topics.contains(MANUAL_TOPIC) + && topics.contains(MANUAL_RESULT_TOPIC); + } catch (Exception e) { + LOG.warn("Exception while checking topic creation", e); + return false; + } + }, + Duration.ofSeconds(30), + Duration.ofMillis(100), + "Topics were not created in time"); + } catch (TimeoutException | InterruptedException e) { + throw new RuntimeException("Failed to wait for topic creation", e); + } + + LOG.info( + "Topics {}, {}, {}, and {} created successfully", + TOPIC, + RESULT_TOPIC, + MANUAL_TOPIC, + MANUAL_RESULT_TOPIC); + } + } + + @Test + public void testAvroConfluentIntegrationWithAutoRegister() throws Exception { + // Combine all SQL statements into a single submission + List<String> allSqlStatements = + Arrays.asList( + "SET 'table.dml-sync' = 'true';", + "", + "CREATE TABLE avro_input (", + " name STRING,", + " favoriteNumber STRING,", + " favoriteColor STRING,", + " eventType STRING", + ") WITH (", + " 'connector' = 'kafka',", + " 'topic' = '" + TOPIC + "',", + " 'properties.bootstrap.servers' = '" + + INTER_CONTAINER_KAFKA_ALIAS + + ":9092',", + " 'properties.group.id' = 'test-group',", + " 'scan.startup.mode' = 'earliest-offset',", + " 'scan.bounded.mode' = 'latest-offset',", + " 'format' = 'avro-confluent',", + " 'avro-confluent.url' = 'http://" + + INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS + + ":8081',", + " 'avro-confluent.auto.register.schemas' = 'true'", + ");", + "", + "CREATE TABLE avro_output (", + " name STRING,", + " favoriteNumber STRING,", + " favoriteColor STRING,", + " eventType STRING", + ") WITH (", + " 'connector' = 'kafka',", + " 'topic' = '" + RESULT_TOPIC + "',", + " 'properties.bootstrap.servers' = '" + + INTER_CONTAINER_KAFKA_ALIAS + + ":9092',", + " 'properties.group.id' = 'test-group',", + " 'scan.startup.mode' = 'earliest-offset',", + " 'scan.bounded.mode' = 'latest-offset',", + " 'format' = 'avro-confluent',", + " 'avro-confluent.url' = 'http://" + + INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS + + ":8081',", + " 'avro-confluent.auto.register.schemas' = 'true'", + ");", + "", + "INSERT INTO avro_input VALUES", + " ('Alice', '42', 'blue', 'INSERT'),", + " ('Bob', '7', 'red', 'INSERT'),", + " ('Charlie', '73', 'green', 'INSERT');", + "", + "INSERT INTO avro_output", + "SELECT * FROM avro_input;"); + + LOG.info("Submitting SQL statements: {}", String.join("\n", allSqlStatements)); + + // Execute all SQL statements in a single submission + executeSql(allSqlStatements); + + // Verify output + verifyNumberOfResultRecords(RESULT_TOPIC, 3); + } + + @Test + public void testAvroConfluentIntegrationWithManualRegister() throws Exception { Review Comment: Added `testWritingFailsWhenAutoRegisterDisabledAndSchemaMissing`. It writes to a topic without a registered subject, asserts the job ends in FAILED and that no subject was created afterwards. The failure is the 40401 from the registry, wrapped in the new error message. ########## flink-end-to-end-tests/flink-confluent-schema-registry/src/test/java/org/apache/flink/schema/registry/test/AvroConfluentITCase.java: ########## @@ -0,0 +1,444 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.schema.registry.test; + +import org.apache.flink.connector.testframe.container.FlinkContainers; +import org.apache.flink.connector.testframe.container.FlinkContainersSettings; +import org.apache.flink.connector.testframe.container.TestcontainersSettings; +import org.apache.flink.core.testutils.CommonTestUtils; +import org.apache.flink.test.resources.ResourceTestUtils; +import org.apache.flink.test.util.SQLJobSubmission; +import org.apache.flink.util.DockerImageVersions; +import org.apache.flink.util.jackson.JacksonMapperFactory; + +import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper; + +import org.apache.avro.generic.GenericRecord; +import org.apache.kafka.clients.admin.Admin; +import org.apache.kafka.clients.admin.NewTopic; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.testcontainers.containers.GenericContainer; +import org.testcontainers.containers.KafkaContainer; +import org.testcontainers.containers.Network; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.utility.DockerImageName; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.HttpRequest; +import java.net.http.HttpResponse; +import java.nio.file.Path; +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Properties; +import java.util.Set; +import java.util.concurrent.TimeoutException; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** E2E Test for Avro-Confluent integration. */ +@Testcontainers +public class AvroConfluentITCase { + + private static final Logger LOG = LoggerFactory.getLogger(AvroConfluentITCase.class); + + private static final String INTER_CONTAINER_KAFKA_ALIAS = "kafka"; + private static final String INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS = "schema-registry"; + + private static final String TOPIC = "test-avro-input"; + private static final String RESULT_TOPIC = "test-avro-output"; + private static final String MANUAL_TOPIC = "test-avro-input-manual"; + private static final String MANUAL_RESULT_TOPIC = "test-avro-output-manual"; + + private static final String INPUT_SCHEMA = SchemaLoader.loadSchema("avro/input-record.avsc"); + private static final String OUTPUT_SCHEMA = SchemaLoader.loadSchema("avro/output-record.avsc"); + + private static final Path sqlToolBoxJar = ResourceTestUtils.getResource(".*/SqlToolbox\\.jar"); + + private final Path sqlConnectorKafkaJar = ResourceTestUtils.getResource(".*kafka.*\\.jar"); + private final Path sqlConnectorUpsertTestJar = + ResourceTestUtils.getResource(".*flink-test-utils.*\\.jar"); + private final Path sqlAvroConfluentJar = + ResourceTestUtils.getResource(".*avro-confluent.*\\.jar"); + + private static final Network NETWORK = Network.newNetwork(); + + @Container + public static final KafkaContainer KAFKA = + new KafkaContainer(DockerImageName.parse(DockerImageVersions.KAFKA)) + .withNetwork(NETWORK) + .withNetworkAliases(INTER_CONTAINER_KAFKA_ALIAS); + + @Container + private static final GenericContainer<?> SCHEMA_REGISTRY = + new GenericContainer<>(DockerImageName.parse(DockerImageVersions.SCHEMA_REGISTRY)) + .withNetwork(NETWORK) + .withNetworkAliases(INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS) + .withExposedPorts(8081) + .withEnv("SCHEMA_REGISTRY_HOST_NAME", INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS) + .withEnv( + "SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", + INTER_CONTAINER_KAFKA_ALIAS + ":9092") + .dependsOn(KAFKA); + + private static final FlinkContainers FLINK = + FlinkContainers.builder() + .withFlinkContainersSettings( + FlinkContainersSettings.builder().numTaskManagers(1).build()) + .withTestcontainersSettings( + TestcontainersSettings.builder() + .network(NETWORK) + .logger(LOG) + .dependsOn(KAFKA) + .build()) + .build(); + + private final HttpClient client = HttpClient.newHttpClient(); + private final ObjectMapper objectMapper = JacksonMapperFactory.createObjectMapper(); + + @BeforeAll + public static void setup() throws Exception { + KAFKA.start(); + SCHEMA_REGISTRY.start(); + FLINK.start(); + } + + @AfterAll + public static void tearDown() { + FLINK.stop(); + SCHEMA_REGISTRY.stop(); + KAFKA.stop(); + } + + @BeforeEach + public void before() throws Exception { + // Create topics using external bootstrap servers since we're outside the Docker network + Properties props = new Properties(); + props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, KAFKA.getBootstrapServers()); + try (Admin admin = Admin.create(props)) { + admin.createTopics( + Arrays.asList( + new NewTopic(TOPIC, 1, (short) 1), + new NewTopic(RESULT_TOPIC, 1, (short) 1), + new NewTopic(MANUAL_TOPIC, 1, (short) 1), + new NewTopic(MANUAL_RESULT_TOPIC, 1, (short) 1))); + + // Poll for topic creation with timeout + try { + CommonTestUtils.waitUntilIgnoringExceptions( + () -> { + try { + Set<String> topics = admin.listTopics().names().get(); + return topics.contains(TOPIC) + && topics.contains(RESULT_TOPIC) + && topics.contains(MANUAL_TOPIC) + && topics.contains(MANUAL_RESULT_TOPIC); + } catch (Exception e) { + LOG.warn("Exception while checking topic creation", e); + return false; + } + }, + Duration.ofSeconds(30), + Duration.ofMillis(100), + "Topics were not created in time"); + } catch (TimeoutException | InterruptedException e) { + throw new RuntimeException("Failed to wait for topic creation", e); + } + + LOG.info( + "Topics {}, {}, {}, and {} created successfully", + TOPIC, + RESULT_TOPIC, + MANUAL_TOPIC, + MANUAL_RESULT_TOPIC); + } + } + + @Test + public void testAvroConfluentIntegrationWithAutoRegister() throws Exception { + // Combine all SQL statements into a single submission + List<String> allSqlStatements = + Arrays.asList( + "SET 'table.dml-sync' = 'true';", + "", + "CREATE TABLE avro_input (", + " name STRING,", + " favoriteNumber STRING,", + " favoriteColor STRING,", + " eventType STRING", + ") WITH (", + " 'connector' = 'kafka',", + " 'topic' = '" + TOPIC + "',", + " 'properties.bootstrap.servers' = '" + + INTER_CONTAINER_KAFKA_ALIAS + + ":9092',", + " 'properties.group.id' = 'test-group',", + " 'scan.startup.mode' = 'earliest-offset',", + " 'scan.bounded.mode' = 'latest-offset',", + " 'format' = 'avro-confluent',", + " 'avro-confluent.url' = 'http://" + + INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS + + ":8081',", + " 'avro-confluent.auto.register.schemas' = 'true'", + ");", + "", + "CREATE TABLE avro_output (", + " name STRING,", + " favoriteNumber STRING,", + " favoriteColor STRING,", + " eventType STRING", + ") WITH (", + " 'connector' = 'kafka',", + " 'topic' = '" + RESULT_TOPIC + "',", + " 'properties.bootstrap.servers' = '" + + INTER_CONTAINER_KAFKA_ALIAS + + ":9092',", + " 'properties.group.id' = 'test-group',", + " 'scan.startup.mode' = 'earliest-offset',", + " 'scan.bounded.mode' = 'latest-offset',", + " 'format' = 'avro-confluent',", + " 'avro-confluent.url' = 'http://" + + INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS + + ":8081',", + " 'avro-confluent.auto.register.schemas' = 'true'", + ");", + "", + "INSERT INTO avro_input VALUES", + " ('Alice', '42', 'blue', 'INSERT'),", + " ('Bob', '7', 'red', 'INSERT'),", + " ('Charlie', '73', 'green', 'INSERT');", + "", + "INSERT INTO avro_output", + "SELECT * FROM avro_input;"); Review Comment: Done, both round-trip tests now use `avro_confluent_roundtrip_e2e.sql` with the flag and topics as variables. ########## flink-end-to-end-tests/flink-confluent-schema-registry/src/test/java/org/apache/flink/schema/registry/test/AvroConfluentITCase.java: ########## @@ -0,0 +1,444 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.schema.registry.test; + +import org.apache.flink.connector.testframe.container.FlinkContainers; +import org.apache.flink.connector.testframe.container.FlinkContainersSettings; +import org.apache.flink.connector.testframe.container.TestcontainersSettings; +import org.apache.flink.core.testutils.CommonTestUtils; +import org.apache.flink.test.resources.ResourceTestUtils; +import org.apache.flink.test.util.SQLJobSubmission; +import org.apache.flink.util.DockerImageVersions; +import org.apache.flink.util.jackson.JacksonMapperFactory; + +import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.ObjectMapper; + +import org.apache.avro.generic.GenericRecord; +import org.apache.kafka.clients.admin.Admin; +import org.apache.kafka.clients.admin.NewTopic; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.ConsumerRecords; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; +import org.testcontainers.containers.GenericContainer; +import org.testcontainers.containers.KafkaContainer; +import org.testcontainers.containers.Network; +import org.testcontainers.junit.jupiter.Container; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.utility.DockerImageName; + +import java.net.URI; +import java.net.http.HttpClient; +import java.net.http.HttpRequest; +import java.net.http.HttpResponse; +import java.nio.file.Path; +import java.time.Duration; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Properties; +import java.util.Set; +import java.util.concurrent.TimeoutException; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** E2E Test for Avro-Confluent integration. */ +@Testcontainers +public class AvroConfluentITCase { + + private static final Logger LOG = LoggerFactory.getLogger(AvroConfluentITCase.class); + + private static final String INTER_CONTAINER_KAFKA_ALIAS = "kafka"; + private static final String INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS = "schema-registry"; + + private static final String TOPIC = "test-avro-input"; + private static final String RESULT_TOPIC = "test-avro-output"; + private static final String MANUAL_TOPIC = "test-avro-input-manual"; + private static final String MANUAL_RESULT_TOPIC = "test-avro-output-manual"; + + private static final String INPUT_SCHEMA = SchemaLoader.loadSchema("avro/input-record.avsc"); + private static final String OUTPUT_SCHEMA = SchemaLoader.loadSchema("avro/output-record.avsc"); + + private static final Path sqlToolBoxJar = ResourceTestUtils.getResource(".*/SqlToolbox\\.jar"); + + private final Path sqlConnectorKafkaJar = ResourceTestUtils.getResource(".*kafka.*\\.jar"); + private final Path sqlConnectorUpsertTestJar = + ResourceTestUtils.getResource(".*flink-test-utils.*\\.jar"); + private final Path sqlAvroConfluentJar = + ResourceTestUtils.getResource(".*avro-confluent.*\\.jar"); + + private static final Network NETWORK = Network.newNetwork(); + + @Container + public static final KafkaContainer KAFKA = + new KafkaContainer(DockerImageName.parse(DockerImageVersions.KAFKA)) + .withNetwork(NETWORK) + .withNetworkAliases(INTER_CONTAINER_KAFKA_ALIAS); + + @Container + private static final GenericContainer<?> SCHEMA_REGISTRY = + new GenericContainer<>(DockerImageName.parse(DockerImageVersions.SCHEMA_REGISTRY)) + .withNetwork(NETWORK) + .withNetworkAliases(INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS) + .withExposedPorts(8081) + .withEnv("SCHEMA_REGISTRY_HOST_NAME", INTER_CONTAINER_SCHEMA_REGISTRY_ALIAS) + .withEnv( + "SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS", + INTER_CONTAINER_KAFKA_ALIAS + ":9092") + .dependsOn(KAFKA); + + private static final FlinkContainers FLINK = + FlinkContainers.builder() + .withFlinkContainersSettings( + FlinkContainersSettings.builder().numTaskManagers(1).build()) + .withTestcontainersSettings( + TestcontainersSettings.builder() + .network(NETWORK) + .logger(LOG) + .dependsOn(KAFKA) + .build()) + .build(); + + private final HttpClient client = HttpClient.newHttpClient(); + private final ObjectMapper objectMapper = JacksonMapperFactory.createObjectMapper(); Review Comment: Done. ########## docs/content.zh/docs/connectors/table/formats/avro-confluent.md: ########## @@ -280,6 +280,13 @@ Format 参数 <td>String</td> <td>The URL of the Confluent Schema Registry to fetch/register schemas.</td> </tr> + <tr> + <td><h5>auto.register.schemas</h5></td> + <td>optional</td> + <td style="word-wrap: break-word;">true</td> + <td>Boolean</td> + <td>Whether to automatically register schemas with the Confluent Schema Registry if they don't exist. When set to <code>false</code>, schemas must be manually registered in the Schema Registry before being used. When set to <code>true</code>, schemas will be automatically registered during serialization if they don't already exist. The default value is <code>true</code>.</td> Review Comment: Added to the format description and to the option text: registration only happens when writing, reading always looks up the writer schema by the id in the record. ########## flink-end-to-end-tests/flink-confluent-schema-registry/src/test/resources/log4j2-test.properties: ########## @@ -0,0 +1,34 @@ +################################################################################ +# Licensed to the Apache Software Foundation (ASF) under one +# or more contributor license agreements. See the NOTICE file +# distributed with this work for additional information +# regarding copyright ownership. The ASF licenses this file +# to you under the Apache License, Version 2.0 (the +# "License"); you may not use this file except in compliance +# with the License. You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. +################################################################################ + +# Set root logger level to OFF to not flood build logs +# set manually to INFO for debugging purposes +rootLogger.level = INFO Review Comment: Done. ########## flink-formats/flink-avro-confluent-registry/src/main/java/org/apache/flink/formats/avro/registry/confluent/ConfluentSchemaRegistryCoder.java: ########## @@ -81,13 +100,29 @@ public Schema readSchema(InputStream in) throws IOException { @Override public void writeSchema(Schema schema, OutputStream out) throws IOException { - try { - int registeredId = schemaRegistryClient.register(subject, schema); - out.write(CONFLUENT_MAGIC_BYTE); - byte[] schemaIdBytes = ByteBuffer.allocate(4).putInt(registeredId).array(); - out.write(schemaIdBytes); - } catch (RestClientException e) { - throw new IOException("Could not register schema in registry", e); + int registeredId; + if (registerSchema()) { + try { + registeredId = schemaRegistryClient.register(subject, schema); + } catch (RestClientException e) { + throw new IOException("Could not register schema in registry", e); + } + } else { + try { + registeredId = schemaRegistryClient.getId(subject, schema); + } catch (RestClientException e) { + throw new IOException("Could not retrieve schema in registry", e); + } } + out.write(CONFLUENT_MAGIC_BYTE); + byte[] schemaIdBytes = ByteBuffer.allocate(4).putInt(registeredId).array(); + out.write(schemaIdBytes); + } + + private boolean registerSchema() { Review Comment: It's parsed once in the constructor now, including a check that rejects anything other than true/false. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
