如何启用Kafka Schema验证?Docker Stack配置调整咨询
问题
搭建了包含3个Kafka实例(Kraft模式)、生产者、消费者及Schema Registry的Docker Stack,生产者持续向test_topic发送消息,已在Schema Registry中创建AVRO格式的Person Schema。添加CONFLUENT_SUPPORT_SCHEMA_VALIDATION等环境变量后,Kafka并未校验主题/消息与Schema的匹配性,需调整配置实现该校验功能。
当前Docker Stack配置:
version: '3.8' services: kafka1: image: confluentinc/cp-kafka hostname: kafka1 ports: - "9192:9092" volumes: - kafka1-data:/var/lib/kafka/data environment: #Define the name for Broker and controller and protocols KAFKA_INTER_BROKER_LISTENER_NAME: 'BROKER' KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,BROKER:PLAINTEXT,EXTERNAL:PLAINTEXT #Define the ports used for broker, controller and external KAFKA_LISTENERS: BROKER://kafka1:29092,CONTROLLER://kafka1:29093,EXTERNAL://:9092 # TODO: Adjust advertised IP address from 'localhost:9092' to hostname or ip address of running system KAFKA_ADVERTISED_LISTENERS: BROKER://kafka1:29092,EXTERNAL://localhost:9192 # Configure replication factors for autocreate topics, consumer offsets,log-topic and re-balance delay KAFKA_DEFAULT_REPLICATION_FACTOR: 3 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3 # Following parameters are required to run Kafka in Kraft mode (without Zookeeper) KAFKA_PROCESS_ROLES: 'controller,broker' KAFKA_NODE_ID: 1 KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:29093,2@kafka2:29093,3@kafka3:29093' # TODO: Generate ONE kafka-ID using ./bin/kafka-storage.sh random-uuid and set it to all instances CLUSTER_ID: 'aiWo7IWazngRchmPES6q5A==' # Required to avoid log spamming from MetdataLoader KAFKA_LOG4J_LOGGERS: 'org.apache.kafka.image.loader.MetadataLoader=WARN' CONFLUENT_SUPPORT_SCHEMA_VALIDATION: "true" SCHEMA_REGISTRY_URL: "http://schemaregistry1:8085" deploy: placement: constraints: # TODO: To ensure that kafka instances are not placed on the same VM, you should define different constraints such as manager or worker to ensure that at least two instances run on different VMs - "node.role == manager" kafka2: image: confluentinc/cp-kafka hostname: kafka2 ports: - "9193:9093" volumes: - kafka2-data:/var/lib/kafka/data environment: #Define the name for Broker and controller and protocols KAFKA_INTER_BROKER_LISTENER_NAME: 'BROKER' KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,BROKER:PLAINTEXT,EXTERNAL:PLAINTEXT #Define the ports used for broker, controller and external KAFKA_LISTENERS: BROKER://kafka2:29092,CONTROLLER://kafka2:29093,EXTERNAL://:9093 # TODO: Adjust advertised IP address from 'localhost:9093' to hostname or ip address of running system KAFKA_ADVERTISED_LISTENERS: BROKER://kafka2:29092,EXTERNAL://localhost:9193 # Configure replication factors for autocreate topics, consumer offsets,log-topic and re-balance delay KAFKA_DEFAULT_REPLICATION_FACTOR: 3 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3 # Following parameters are required to run Kafka in Kraft mode (without Zookeeper) KAFKA_PROCESS_ROLES: 'controller,broker' KAFKA_NODE_ID: 2 KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:29093,2@kafka2:29093,3@kafka3:29093' # TODO: Generate ONE kafka-ID using ./bin/kafka-storage.sh random-uuid and set it to all instances CLUSTER_ID: 'aiWo7IWazngRchmPES6q5A==' # Required to avoid log spamming from MetdataLoader KAFKA_LOG4J_LOGGERS: 'org.apache.kafka.image.loader.MetadataLoader=WARN' CONFLUENT_SUPPORT_SCHEMA_VALIDATION: "true" SCHEMA_REGISTRY_URL: "http://schemaregistry1:8085" deploy: placement: constraints: # TODO: To ensure that kafka instances are not placed on the same VM, you should define different constraints such as manager or worker to ensure that at least two instances run on different VMs - "node.role == manager" kafka3: image: confluentinc/cp-kafka hostname: kafka3 ports: - "9194:9094" volumes: - kafka3-data:/var/lib/kafka/data environment: # Define the name for Broker and controller and protocols KAFKA_INTER_BROKER_LISTENER_NAME: 'BROKER' KAFKA_CONTROLLER_LISTENER_NAMES: 'CONTROLLER' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,BROKER:PLAINTEXT,EXTERNAL:PLAINTEXT # Define the ports used for broker, controller and external KAFKA_LISTENERS: BROKER://kafka3:29092,CONTROLLER://kafka3:29093,EXTERNAL://:9094 # TODO: Adjust advertised IP address from 'localhost:9094' to hostname or ip address of running system KAFKA_ADVERTISED_LISTENERS: BROKER://kafka3:29092,EXTERNAL://localhost:9194 # Configure replication factors for autocreate topics, consumer offsets,log-topic and re-balance delay KAFKA_DEFAULT_REPLICATION_FACTOR: 3 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 3 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 3 # Following parameters are required to run Kafka in Kraft mode (without Zookeeper) KAFKA_PROCESS_ROLES: 'controller,broker' KAFKA_NODE_ID: 3 KAFKA_CONTROLLER_QUORUM_VOTERS: '1@kafka1:29093,2@kafka2:29093,3@kafka3:29093' # TODO: Generate ONE kafka-ID using ./bin/kafka-storage.sh random-uuid and set it to all instances CLUSTER_ID: 'aiWo7IWazngRchmPES6q5A==' # Required to avoid log spamming from MetdataLoader KAFKA_LOG4J_LOGGERS: 'org.apache.kafka.image.loader.MetadataLoader=WARN' CONFLUENT_SUPPORT_SCHEMA_VALIDATION: "true" SCHEMA_REGISTRY_URL: "http://schemaregistry1:8085" deploy: placement: constraints: # TODO: To ensure that kafka instances are not placed on the same VM, you should define different constraints such as manager or worker to ensure that at least two instances run on different VMs - "node.role == manager" kafka-ui: container_name: kafka-ui image: provectuslabs/kafka-ui:latest ports: - "9180:8080" environment: DYNAMIC_CONFIG_ENABLED: 'true' KAFKA_CLUSTERS_0_NAME: digi-production KAFKA_CLUSTERS_0_BOOTSTRAPSERVERS: kafka1:29092,kafka2:29092,kafka3:29092 KAFKA_CLUSTERS_0_SCHEMAREGISTRY: http://schemaregistry1:8085 KAFKA_CLUSTERS_0_SCHEMAREGISTRYAUTH_USERNAME: admin KAFKA_CLUSTERS_0_SCHEMAREGISTRYAUTH_PASSWORD: letmein healthcheck: test: ["CMD-SHELL", "wget -nv -t1 --spider 'http://localhost:8080'"] interval: 10s timeout: 10s retries: 3 schemaregistry1: image: confluentinc/cp-schema-registry:7.2.1 hostname: schemaregistry1 ports: - 18085:8085 environment: SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS: PLAINTEXT://kafka1:29092,PLAINTEXT://kafka2:29092,PLAINTEXT://kafka3:29092 SCHEMA_REGISTRY_KAFKASTORE_SECURITY_PROTOCOL: PLAINTEXT SCHEMA_REGISTRY_HOST_NAME: schemaregistry1 SCHEMA_REGISTRY_LISTENERS: http://schemaregistry1:8085 # Default credentials: admin/letmein #SCHEMA_REGISTRY_AUTHENTICATION_METHOD: BASIC #SCHEMA_REGISTRY_AUTHENTICATION_REALM: SchemaRegistryProps #SCHEMA_REGISTRY_AUTHENTICATION_ROLES: admin #SCHEMA_REGISTRY_OPTS: -Djava.security.auth.login.config=/conf/schema_registry.jaas SCHEMA_REGISTRY_SCHEMA_REGISTRY_INTER_INSTANCE_PROTOCOL: "http" SCHEMA_REGISTRY_LOG4J_ROOT_LOGLEVEL: INFO SCHEMA_REGISTRY_KAFKASTORE_TOPIC: _schemas kafka-producer: image: confluentinc/cp-kafka command: > bash -c "while true; do echo '{\"id\": 1, \"name\": \"John\", \"age\": 30}' | kafka-console-producer --broker-list kafka1:29092 --topic test_topic \ --property value.schema='\"Person\"' && sleep 5; done" kafka-consumer: image: confluentinc/cp-kafka command: > kafka-console-consumer --bootstrap-server kafka1:29092 --topic test_topic \ --from-beginning --property print.key=true \ --property key.deserializer=org.apache.kafka.common.serialization.StringDeserializer \ --property value.deserializer=io.confluent.kafka.serializers.KafkaAvroDeserializer \ --property schema.registry.url=http://schemaregistry1:8085 \ --property schema.registry.user=admin \ --property schema.registry.password=letmein volumes: kafka1-data: driver: distributed-docker-volumes kafka2-data: driver: distributed-docker-volumes kafka3-data: driver: distributed-docker-volumes
已尝试添加但未生效的环境变量:
CONFLUENT_SUPPORT_SCHEMA_VALIDATION: "true" SCHEMA_REGISTRY_URL: "http://schemaregistry1:8085"
生产者使用--property value.schema='"Person"'参数发送消息,该值为已创建的Schema名称。
配置调整方案
1. 修正Kafka Broker的Schema验证配置
删除每个Kafka Broker中的CONFLUENT_SUPPORT_SCHEMA_VALIDATION和SCHEMA_REGISTRY_URL,替换为以下环境变量,启用Broker端的Schema验证插件:
KAFKA_PLUGIN_PATH: "/usr/share/java/kafka-schema-registry" KAFKA_CONFLUENT_SCHEMA_VALIDATION_ENABLE: "true"
说明:
CONFLUENT_SUPPORT_SCHEMA_VALIDATION是Confluent Support的诊断开关,并非用于Broker端的Schema验证;Broker需通过加载kafka-schema-registry插件来实现验证能力。
2. 为目标主题绑定Schema验证规则
需为test_topic明确配置Schema验证规则,可通过CLI执行以下命令:
# 修改现有主题配置 kafka-configs --alter --topic test_topic --bootstrap-server kafka1:29092 \ --add-config confluent.value.schema.validation=true \ --add-config confluent.value.schema.name=Person
如果主题尚未创建,可在创建时直接指定规则:
kafka-topics --create --topic test_topic --bootstrap-server kafka1:29092 \ --partitions 3 --replication-factor 3 \ --config confluent.value.schema.validation=true \ --config confluent.value.schema.name=Person
3. 修复生产者的消息发送逻辑
当前使用的kafka-console-producer仅发送纯文本消息,未使用Avro序列化器,无法触发Schema验证。需替换为kafka-avro-console-producer,并指定Schema Registry地址:
kafka-producer: image: confluentinc/cp-kafka command: > bash -c "while true; do echo '{\"id\": 1, \"name\": \"John\", \"age\": 30}' | kafka-avro-console-producer --broker-list kafka1:29092 --topic test_topic \ --property schema.registry.url=http://schemaregistry1:8085; sleep 5; done"
提示:若需强制指定Schema,可添加
--property value.schema.id=<SchemaID>,SchemaID可通过Schema Registry API获取:curl http://schemaregistry1:8085/subjects/Person/versions/latest
4. 验证Schema Registry与Broker的连通性
确认Schema Registry的SCHEMA_REGISTRY_KAFKASTORE_BOOTSTRAP_SERVERS配置正确,且Kafka Broker能访问Schema Registry的http://schemaregistry1:8085地址(可在Broker容器内执行curl http://schemaregistry1:8085测试连通性)。
内容的提问来源于stack exchange,提问作者Nico

