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

无法通过Spring Boot访问AWS Glue Schema Registry求助

问题:Spring Boot应用无法访问AWS Glue Schema Registry,报错"Failed to get schemaVersionId by schema definition for schema name = hydra-testschema"

我已通过AWS CLI创建了AWS Glue Schema Registry,但Spring Boot应用访问时抛出上述错误,以下是相关信息:

创建Schema的CLI操作

创建Schema的CLI操作

架构文件与测试代码

Customer.avsc

{
  "type": "record",
  "namespace": "ABC_Organization",
  "name": "Employee",
  "fields": [
    {
      "name": "Name",
      "type": "string"
    },
    {
      "name": "Age",
      "type": "int"
    },
    {
      "name": "address",
      "type": {
        "type": "record",
        "name": "addressRecord",
        "fields": [
          {
            "name": "street",
            "type": "string"
          },
          {
            "name": "zipcode",
            "type": "int"
          }
        ]
      }
    }
  ]
}

SampleProducer.java

public class SampleMskProducer{
    private static final Properties properties = new Properties();
    private final static Logger LOGGER = LoggerFactory.getLogger(org.apache.kafka.clients.producer.Producer.class.getName());

    public static void main(String[] args) throws Exception {

        String username = "BrokerUserName";
        String password = "Password";

        String jaasTemplate = "org.apache.kafka.common.security.scram.ScramLoginModule required username=\"%s\" password=\"%s\";";
        String jaasCfg = String.format(jaasTemplate, username, password);

        System.setProperty("software.amazon.awssdk.http.service.impl", "software.amazon.awssdk.http.urlconnection.UrlConnectionSdkHttpService");

        // Setting kafka properties
        properties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "Broker");
        properties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName());
        properties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, GlueSchemaRegistryKafkaSerializer.class.getName());
        properties.put(AWSSchemaRegistryConstants.AWS_REGION, "usa-east-1");
        properties.put(AWSSchemaRegistryConstants.REGISTRY_NAME, "sandbox");
        properties.put(AWSSchemaRegistryConstants.SCHEMA_NAME, "hydra-testschemaa");
        properties.put(AWSSchemaRegistryConstants.DATA_FORMAT, "AVRO");
        properties.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, "SASL_SSL");
        properties.put(SaslConfigs.SASL_MECHANISM, "SCRAM-SHA-512");
        properties.put(SaslConfigs.SASL_JAAS_CONFIG, jaasCfg);

        SystemPropertiesCredentialsProvider systemPropertiesCredentialsProvider=new SystemPropertiesCredentialsProvider();

//Passing the secrets
        System.setProperty("aws.accessKeyId", "A");
        System.setProperty("aws.secretAccessKey", "wC");

// Your AWS SDK or AWS-related code here

        // Declearing and parsing the Schema fields for generic record builder.
        Schema schema_customer = null;
        try {
            schema_customer = new Parser().parse(new File("Customer.avsc"));
        } catch (IOException e) {
            e.printStackTrace();
        }
     
        GenericRecord customer = new GenericData.Record(schema_customer);
        GenericRecord addressRecord = new GenericData.Record(schema_customer.getField("address").schema());
        LOGGER.info("Generic records phase completed...");

        Random rand = new Random();
        int zipcodeValue = 9999;
        int ageMaxValue = 100;

        // Initializing the Producer client, build records and publishing the message to the kafka broker
        try (KafkaProducer<String, GenericRecord> producer = new KafkaProducer<>(properties)) {
            final ProducerRecord<String, GenericRecord> record = new ProducerRecord<String, GenericRecord>("hydra.proxy.updates", customer);
            LOGGER.info("Starting to send records...");
            for (int i = 0; i < 10000; i++) {
                addressRecord.put("street", "city-" + i);
                addressRecord.put("zipcode", rand.nextInt(zipcodeValue));
                customer.put("Name","name-"+ i);
                customer.put("Age",rand.nextInt(ageMaxValue));
                customer.put("address", addressRecord);
                producer.send(record, new ProducerCallback());
            }
            producer.flush();
        } catch (Exception e) {
            e.printStackTrace();
        }
    }

    // Defining method for retrieving secret from Secretmanager

    // Callback class for producer client for logging.
    private static class ProducerCallback implements Callback {
        @Override
        public void onCompletion(RecordMetadata recordMetaData, Exception e) {
            if (e == null) {
                LOGGER.info("Received new metadata. \t" +
                "Topic:" + recordMetaData.topic() + "\t" +
                "Partition: " + recordMetaData.partition() + "\t" +
                "Offset: " + recordMetaData.offset() + "\t" +
                "Timestamp: " + recordMetaData.timestamp());
            } else {
                LOGGER.info("There's been an error from the Producer side");
                e.printStackTrace();
            }
        }
    }
}

执行后报错

执行后报错

排查建议

  • Schema名称拼写校验:代码中配置的Schema名称是hydra-testschemaa(末尾多了一个a),但报错信息中显示的是hydra-testschema,先确认CLI创建的Schema实际名称,保证代码配置与实际Schema名称完全一致。
  • 区域配置修正:代码中配置的AWS区域是usa-east-1,但AWS标准区域标识为us-east-1,需确认CLI创建Schema时使用的区域,保证代码与CLI操作的区域一致。
  • Schema定义一致性检查:确保本地Customer.avsc的Schema定义,与Glue Registry中创建的Schema完全匹配(包括namespace、字段名、类型,甚至格式细节如空格、换行),Glue会通过哈希值严格校验Schema一致性。
  • IAM权限验证:确认代码使用的AWS Access Key拥有glue:GetSchema、glue:GetSchemaVersion等访问Glue Schema Registry的权限,避免因权限不足导致查询失败。
  • 凭证有效性验证:检查代码中设置的AWS Access Key和Secret Access Key是否有效,且属于有权访问该Registry的IAM实体。
  • 网络连通性确认:确保应用所在环境可以正常访问AWS Glue服务端点,无防火墙或VPC配置拦截请求。

内容的提问来源于stack exchange,提问作者Aayush

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.10 20:54:54