You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何启用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 12:55:54